Rendi lo stream parallelo al risultato di flatMap
Considera il seguente semplice codice:
Stream.of(1)
.flatMap(x -> IntStream.range(0, 1024).boxed())
.parallel() // Moving this before flatMap has the same effect because it's just a property of the entire stream
.forEach(x -> {
System.out.println("Thread: " + Thread.currentThread().getName());
});
Per molto tempo ho pensato che Java avrebbe avuto un'esecuzione parallela per gli elementi anche dopo flatMap. Ma il codice di cui sopra stampa tutto "Thread: main", il che dimostra che il mio pensiero è sbagliato.
Un modo semplice per renderlo parallelo dopo flatMapsarebbe quello di raccogliere e quindi trasmettere di nuovo in streaming:
Stream.of(1)
.flatMap(x -> IntStream.range(0, 1024).boxed())
.parallel() // Moving this before flatMap has the same effect because it's just a property of the entire stream
.collect(Collectors.toList())
.parallelStream()
.forEach(x -> {
System.out.println("Thread: " + Thread.currentThread().getName());
});
Mi chiedevo se esiste un modo migliore, e la scelta del design di flatMapquesto parallelizza solo il flusso prima della chiamata, ma non dopo la chiamata.
========= Ulteriori chiarimenti sulla domanda ========
Da alcune risposte, sembra che la mia domanda non sia completamente trasmessa. Come ha detto @Andreas, se inizio con uno stream di 3 elementi, potrebbero esserci 3 thread in esecuzione.
Ma la mia domanda è davvero: Java Stream utilizza un ForkJoinPool comune che ha una dimensione predefinita uguale a uno in meno rispetto al numero di core, secondo questo post . Supponiamo ora che io abbia 64 core, quindi mi aspetto che il mio codice sopra veda molti thread diversi dopo flatMap, ma in realtà ne vede solo uno (o 3 nel caso di Andreas). A proposito, ho usato isParallelper osservare che il flusso è parallelo.
Ad essere onesti, non stavo facendo questa domanda per puro interesse accademico. Mi sono imbattuto in questo problema in un progetto che presenta una lunga catena di operazioni di flusso per trasformare un set di dati. La catena inizia con un singolo file ed esplode in molti elementi flatMap. Ma a quanto pare, nel mio esperimento, NON sfrutta appieno la mia macchina (che ha 64 core), ma utilizza solo un core (dall'osservazione dell'utilizzo della CPU).
Risposte
Mi chiedevo [...] sulla scelta del design
flatMapche parallelizza il flusso solo prima della chiamata, ma non dopo la chiamata.
Ti sbagli. Tutti i passaggi prima e dopo flatMapvengono eseguiti in parallelo, ma divide solo il flusso originale tra i thread. L' flatMapoperazione viene quindi gestita da uno di questi thread e il relativo flusso non viene suddiviso.
Poiché il tuo stream originale ha solo 1 elemento, non può essere diviso e quindi parallelnon ha alcun effetto.
Provare a passare a Stream.of(1, 2, 3), e vedrete che il forEach, che è dopo la flatMap, è in realtà gestita in 3 diversi thread.
La documentazione perforEach specifica:
Per ogni dato elemento, l'azione può essere eseguita in qualsiasi momento e in qualsiasi thread la libreria scelga.
In particolare, "esegui tutte le operazioni sul thread di richiamo" sembra una buona implementazione ampiamente sicura.
Tieni presente che il tuo tentativo di parallelizzare il flusso non richiede alcun parallelismo specifico, ma è molto più probabile che tu veda un effetto con questo:
IntStream.range(0, 1024).boxed()
.parallel()
.map(i -> "Thread: " + Thread.currentThread().getName())
.forEach(System.out::println);