Variable auf Thread in C # Parallel.ForEach beibehalten

Oct 30 2020

Ich möchte die Verarbeitung einer Aufgabe parallelisieren, die von einem Objekt ( State) abhängig ist , das nicht threadsicher ist und dessen Konstruktion zeitaufwändig ist .

Aus diesem Grund habe ich nach partition-lokalen Variablen gesucht , aber entweder mache ich es falsch oder ich suche nach etwas anderem. Dies entspricht mehr oder weniger meiner aktuellen Implementierung:

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) => { });

Ich habe jedoch ein Console.WriteLinein meinem StateInitialisierer hinzugefügt , und ich sehe, dass für jede Iteration der Schleife der StateKonstruktor aufgerufen wird, was zu einem großen Leistungseinbruch führt. Ich möchte das Beispiel von Statein einem Thread auf dem gleichen Thread zu der nachfolgenden Iteration geleitet wird.

Wie kann ich das erreichen?

Antworten

2 TheodorZoulias Oct 30 2020 at 21:45

Sie haben mehrere Möglichkeiten. Am einfachsten ist es, ein einzelnes StateObjekt zu erstellen und den Zugriff darauf zu synchronisieren, indem Sie Folgendes verwenden 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);
});

Ich gehe davon aus, dass dies nicht praktikabel ist, da die DoSomethingMethode zeitaufwändig ist und die Synchronisierung die Parallelität zunichte macht.

Eine andere Möglichkeit ist die Verwendung von a ThreadLocal. Diese Klasse bietet eine threadlokale Speicherung von Daten, sodass die Anzahl der Stateerstellten Objekte der Anzahl der von der verwendeten Threads entspricht 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);
});

Dies wird wahrscheinlich weniger StateObjekte als die Parallel.ForEach<TSource, TLocal>Überladung erzeugen , aber immer noch nicht gleich der konfigurierten MaxDegreeOfParallelism. Das Parallel.ForEachverwendet Threads aus dem ThreadPool, und es ist durchaus möglich, dass es alle während der Berechnung verwendet, vorausgesetzt, die Liste von foldersist ausreichend lang. Und Sie haben wenig Kontrolle über die Größe der ThreadPool. Dies ist also auch keine besonders verlockende Lösung.

Die dritte und letzte Option, die ich mir vorstellen kann, besteht darin, einen Pool von StateObjekten und Rent/ Returneinen in jeder Schleife zu erstellen :

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);
});

Auf diese Weise entspricht die Anzahl der instanziierten StateObjekte dem maximalen Parallelitätsgrad.

Das einzige Problem ist, dass es ObjectPoolin der .NET-Plattform keine Klasse gibt (es gibt nur eine ArrayPoolKlasse), daher müssen Sie eine finden. Hier ist eine einfache Implementierung basierend auf 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();
}