Fork Join bei Arraylist

vimar

Bekanntes Mitglied
Hallo,


ich hab eine Arraylist mit 1000 Objekten und will diese auf 4 Threads/Tasks verteilen (quadcore).


vorher hab ich mit 4 normalen threads gearbeitet, muss aber wegen join() immer auf den langsamsten thread warten (arraylist in 4 gleichgroße teile aufgeteilt und jedem thread 1/4 der arraylist gegeben).
ich brauche join weil ich am ende wissen will ob es eine veränderung gab. problem hierbei: 3 von 4 kernen idlen am schluss der berechnung und warten auf den letzten thread. so wie ich join fork verstanden habe, gibts hier einen work-stealing algo, also wenn ein thread fertig ist mit 1/4 seiner zugewiesenen arraylist, dass er sich dann objekte aus den noch nicht abgearbeiteten objekten der anderen 3 threads nimmt. ich denke also dass ich so effizient fast ohne idle time eines kerns schneller zum ergebnis komme. habe ich das richtig verstanden?

nun kam ich zu fork join,

Java:
static int numberOfProcessors = Runtime.getRuntime().availableProcessors(); // = 4
public static ForkJoinPool fjPool = new ForkJoinPool(numberOfProcessors);


und hier steh ich irgendwie grad auf dem schlauch:

Java:
class MyTask extends RecursiveTask<Boolean> {

    static final int SEQUENTIAL_THRESHOLD = 1000;
    ArrayList <cluster> cl;
    int firstcluster; // range lowest
    int lastcluster; // range highest
    
    MyTask(ArrayList <cluster> cl, int firstcluster, int lastcluster){
        this.cl = cl;
        this.firstcluster = firstcluster;
        this.lastcluster = lastcluster;
    }
    
    @Override
    protected Boolean compute() {
     
        
        
        
        return false;
    }
    
}



ich muss doch garnicht mehr "threads" erstellen (new thread..)? inwiefern splitte ich die arraylist nun auf? macht forkjoin da alles automatisch? was müsste ich denn tendenziell schreiben?

ich möchte einen boolean als rückgabewert der gesamten arraylistabarbeitung, daher -> boolean compute(). ich schreibe in meinen berechnungen weder in die arraylist noch remove ich elemente. ich will lediglich objekte auslesen.

mfg vimar
 
Zuletzt bearbeitet:
Das work-stealing bezieht sich auf Tasks bzw. sub-tasks. Die Threads werden automatisch erstellt. Gibt's eine konkretere Frage dazu?
 
ja gibt einen grund. geht darum dass jeder thread einen gewissen booleanwert verändern kann bei einem objekt von false auf true.

am ende wenn alle threads fertig sind schaue ich nach ob dieser boolean auf true geändert wurde. daher muss ich auch auf den letzten thread warten.
 
Huiui ... klingt gefährlich, wenn "viele" Threads schreibend auf was zugreifen... naja, wenn's nur eine boolean-Variable ist, könnte es OK sein (zumindest volatile sollte sie dann aber sein... ggf. auch irgendwas mti syncrhonized drumrum...)
 
ich weiss auch nicht warum das so schwer fällt hier... einfach ne arraylist statt sequentiell auslesen einfach mit 4 threads.. unbegreiflich dass man da so lange dran werkeln muss
 
Wenn du schon weißt, dass es nur 4 Threads sein sollen, und es keine hierarchisch zerlegbare Aufgabe ist, warum dann mit einem ForkJoinPool? Bei einer List würde sich das eigentlich vielleicht anbieten, da könnte man sicher was geschickt mit subList machen, aber wie man das machen könnte, ... dazu müßte ich mir den ForkJoinPool nochmal genauer ansehen.

Ist das Work-Stealing wichtig? Beschreib' doch mal genauer, worum es eigentlich geht. Wenn's nur darum geht, alle 4 cores eine Zeitlang beschäftigt zu halten: Für sowas verwende ich immer so ein Muster wie
Java:
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;


public class ParallelExecutionSample
{
    public static void main(String[] args)
    {
        List<String> list = new ArrayList<String>();
        for (int i=0; i<100000; i++)
        {
            list.add(String.valueOf(i));
        }

        ParallelExecutionSample sample = new ParallelExecutionSample();
        sample.execute(list);
    }
    
    private final ExecutorService executorService;
    private final int numProcessors;
    
    ParallelExecutionSample()
    {
        numProcessors = Runtime.getRuntime().availableProcessors();
        executorService = Executors.newFixedThreadPool(numProcessors);
    }

    public void execute(final List<?> list)
    {
        int batchSize = (int)Math.ceil((double)list.size() / numProcessors);
        List<Callable<Void>> tasks = new ArrayList<Callable<Void>>();
        for (int i=0; i<numProcessors; i++)
        {
            final int minIndex = i * batchSize;
            final int maxIndex = Math.min(list.size(), minIndex + batchSize);
            Callable<Void> callable = new Callable<Void>()
            {
                @Override
                public Void call() throws Exception
                {
                    execute(list, minIndex, maxIndex);
                    return null;
                }
            };
            tasks.add(callable);
        }
        try
        {
            executorService.invokeAll(tasks);
        }
        catch (InterruptedException e)
        {
            Thread.currentThread().interrupt();
        }
    }
    
    private void execute(List<?> list, int minIndex, int maxIndex)
    {
        System.out.println("Execute something in the range "+minIndex+" to "+maxIndex);
        double dummy = 0;
        for (int i=minIndex; i<maxIndex; i++)
        {
            double v = list.get(i).hashCode();
            for (int j=0; j<100; j++)
            {
                dummy += Math.cos(Math.sin(Math.tan(v)));
                dummy *= Math.cos(Math.sin(Math.tan(v)));
                dummy /= Math.cos(Math.sin(Math.tan(v)));
            }
        }
        System.out.println("Result "+dummy);
    }

}
 

Zurück
Oben