Pertahankan variabel pada thread di C # Parallel.ForEach

Oct 30 2020

Saya ingin memparalelkan pemrosesan tugas yang bergantung pada objek ( State), yang tidak aman untuk thread, dan yang konstruksinya mahal waktu .

Untuk alasan ini, saya sedang mencari variabel partisi-lokal , tetapi saya salah melakukannya atau mencari yang lain. Ini kurang lebih mewakili implementasi saya saat ini:

Parallel.ForEach<string, State>(folders, config, () => new State(), (source, loopState, index, threadState) => {
    var content = File.ReadAllText(source);        // read file
    var result = threadState.doSomething(content); // do something
    File.WriteAllText(outputFile, result);         // write output
    return threadState;
}, (threadState) => { });

Namun, saya telah menambahkan a Console.WriteLinedi Statepenginisialisasi saya , dan saya melihat bahwa untuk setiap iterasi loop, Statekonstruktor dipanggil, menghasilkan kinerja yang besar. Saya ingin contoh dari Statedi satu thread yang dilewatkan ke iterasi berikutnya pada thread yang sama.

Bagaimana saya bisa mencapai itu?

Jawaban

2 TheodorZoulias Oct 30 2020 at 21:45

Anda memiliki dua pilihan. Yang paling sederhana adalah membuat satu Stateobjek, dan menyinkronkan aksesnya dengan menggunakan lock:

var state = new State();

Parallel.ForEach(folders, config, source =>
{
    var content = File.ReadAllText(source);
    string result;
    lock (state) { result = state.DoSomething(content); }
    File.WriteAllText(outputFile, result);
});

Saya berasumsi bahwa ini tidak dapat dijalankan karena DoSomethingmetode ini memakan waktu, dan menyinkronkannya akan mematikan paralelisme.

Pilihan lainnya adalah menggunakan file ThreadLocal. Kelas ini menyediakan penyimpanan data lokal utas, sehingga jumlah Stateobjek yang dibuat akan sama dengan jumlah utas yang digunakan oleh Parallel.ForEach.

var threadLocalState = new ThreadLocal<State>(() => new State());

Parallel.ForEach(folders, config, source =>
{
    var content = File.ReadAllText(source);
    var result = threadLocalState.Value.DoSomething(content);
    File.WriteAllText(outputFile, result);
});

Ini mungkin akan membuat Stateobjek lebih sedikit daripada Parallel.ForEach<TSource, TLocal>kelebihan beban, tetapi masih tidak sama dengan yang dikonfigurasi MaxDegreeOfParallelism. The Parallel.ForEachpenggunaan benang dari ThreadPool, dan sangat mungkin bahwa ia akan menggunakan semua dari mereka selama perhitungan, asalkan daftar folderscukup panjang. Dan Anda memiliki sedikit kendali atas ukuran file ThreadPool. Jadi ini juga bukan solusi yang menarik.

Opsi ketiga dan terakhir yang dapat saya pikirkan adalah membuat kumpulan Stateobjek, dan Rent/ Returnsatu di setiap loop:

var statePool = new ObjectPool<State>(() => new State());

Parallel.ForEach(folders, config, source =>
{
    var state = statePool.Rent();
    var content = File.ReadAllText(source);
    var result = state.DoSomething(content);
    File.WriteAllText(outputFile, result);
    statePool.Return(state);
});

Dengan cara ini jumlah Stateobjek yang dipakai akan sama dengan derajat maksimum paralelisme.

Satu-satunya masalah adalah tidak ada ObjectPoolkelas di platform .NET (hanya ada satu ArrayPoolkelas), jadi Anda harus menemukannya. Berikut adalah implementasi sederhana berdasarkan ConcurrentBag:

public class ObjectPool<T> : IEnumerable<T> where T : new()
{
    private readonly ConcurrentBag<T> _bag = new ConcurrentBag<T>();
    private readonly Func<T> _factory;

    public ObjectPool(Func<T> factory = null) => _factory = factory;

    public T Rent()
    {
        if (_bag.TryTake(out var obj)) return obj;
        return _factory != null ? _factory() : new T();
    }

    public void Return(T obj) => _bag.Add(obj);

    public IEnumerator<T> GetEnumerator() => _bag.GetEnumerator();
    IEnumerator IEnumerable.GetEnumerator() => this.GetEnumerator();
}