parallel

Solapa trabajo de CPU entre isolates, en el orden de la fuente. No es concurrent con otro nombre.

int get parallelWorkers FxAsyncIterable<R> parallel<A, R>(int workers, FutureOr<R> Function(A input) worker, Iterable<A> iterable, {int chunk = 1, bool chunked = false}) FxAsyncIterable<R> parallelAsync<A, R>(int workers, FutureOr<R> Function(A input) worker, FxAsyncIterable<A> iterable, {int chunk = 1, bool chunked = false}) FxAsyncIterable<R> mapParallel<A, R>(int workers, FutureOr<R> Function(A input) worker, Iterable<A> iterable, {int chunk = 1, bool chunked = false}) FxAsyncIterable<R> mapParallelAsync<A, R>(int workers, FutureOr<R> Function(A input) worker, FxAsyncIterable<A> iterable, {int chunk = 1, bool chunked = false}) fxPipe2(decodePng, thumbnail) // one hop, see fxPipe IsolatePool.spawn(int workers) / IsolatePool.using(int workers, body) FxAsyncIterable<R> parallelOn<A, R>(IsolatePool pool, FutureOr<R> Function(A input) worker, Iterable<A> iterable, {int chunk = 1, bool chunked = false}) FxAsync<R> Fx.parallel(int workers, FutureOr<R> Function(T input) worker, {int chunk = 1, bool chunked = false}) FxAsync<R> Fx.parallelOn(IsolatePool pool, FutureOr<R> Function(T input) worker, {int chunk = 1, bool chunked = false}) FxAsync<R> FxAsync.parallel(int workers, FutureOr<R> Function(T input) worker, {int chunk = 1, bool chunked = false}) FxAsync<R> FxAsync.parallelOn(IsolatePool pool, FutureOr<R> Function(T input) worker, {int chunk = 1, bool chunked = false})

Lección

concurrent(n) solapa Futures en el mismo isolate — I/O. La historia de CPU de Dart son los isolates. parallel(n, worker) es el gemelo: un pool reutilizado de n isolates, resultados en el orden de la fuente, como mapConcurrent es la forma combinada de map-más-concurrent. No son el mismo operador — la comparación está en concurrent or parallel.

Prefiere una función top-level o static. Un closure que capture algo no enviable (un ReceivePort, un socket abierto) lanza ArgumentError al spawn — el contrato de isolate, no un invento de fxdart. Un input o un resultado no enviable falla ese pull del mismo modo, en lugar de colgarse. En la web el operador lanza UnsupportedError — usa concurrent(n) ahí. Este listado es solo VM y no es un playground en vivo. El worker puede devolver un Future (FutureOr, la misma forma que mapConcurrent) — un callback síncrono sigue siendo el camino rápido. Un parallel anidado dentro de un worker async está permitido: ese isolate crea su propio pool, y el cancel de la cadena exterior apaga el pool interno. Un nivel de anidación es el contrato — un tercer parallel anidado muere con su padre, así que no puede apagar sus hijos.

¿No quieres elegir n? parallelWorkers es el número de procesadores de la VM — pásalo como primer argumento. Una List más corta que n ajusta el pool a la lista, así que parallel(8, w) sobre dos elementos arranca dos isolates, no ocho. Quien venga de mapConcurrent puede escribir mapParallel; es el mismo operador.

int timesTen(int x) => x * 10;

Future<void> main() async {
  print(await fx([1, 2, 3, 4]).parallel(2, timesTen).toList());
  // [10, 20, 30, 40]
}

chunk — cuántos elementos viajan en un mensaje

Por defecto cada elemento cruza a un worker por su cuenta. Ese viaje de ida y vuelta cuesta unos 5µs, más que la mayoría de callbacks, y es toda la razón por la que un worker barato es más lento bajo parallel que en un bucle simple. chunk: k lo paga una vez cada k elementos:

// 20,000 elements, ~0.4µs of work each, 4 workers:
await fx(rows).parallel(4, parseRow).toList();             // ~142ms
await fx(rows).parallel(4, parseRow, chunk: 512).toList(); //   ~3ms

// the same work in a plain loop, no isolates:            //   ~8ms

47× en esa forma, y la forma por lotes es la primera que de verdad gana al bucle que sustituye. Elige k para que k × callback supere holgadamente 5µs, y aun así queden varios lotes por worker para equilibrar — length ~/ (workers * 4) es un buen punto de partida.

Un lote no cambia lo que observas: el orden es el mismo, la contrapresión es la misma, y un worker que lanza sigue emitiendo los resultados de los elementos anteriores y luego raisea en el elemento que de verdad falló. Cambian dos cosas. El primer elemento ahora espera a todo su lote, así que un take(1) quiere un chunk pequeño o ninguno. Y un input o resultado no enviable falla todo el lote en lugar de solo su pull — averiguar qué elemento tuvo la culpa significaría enviarlos por separado, que es el coste que el lote existe para evitar.

¿No quieres escribir chunk: n ~/ (workers * 4) y repetir el número de workers? chunked: true lo hace a partir de la longitud:

await fx(rows).parallel(4, parseRow, chunked: true);
// k = rows.length ~/ 16 — un 4, no dos

La fuente tiene que ser una List. Un generator o una fuente async no tiene longitud — pasa chunk: k. chunk: y chunked: juntos lanzan; la llamada tiene una sola política.

Dos etapas de CPU, un hop

Dos llamadas a .parallel copian cada resultado de vuelta a este isolate y otra vez hacia fuera. Compón los workers con fxPipe2 para que ambas etapas corran en el worker:

await fx(blobs)
    .parallel(4, fxPipe2(decodePng, thumbnail), chunk: 64)
    .toList();

decodePng y thumbnail tienen que ser enviables, igual que cualquier worker de parallel. La función devuelta captura ambos. Añade .then para más etapas — no hay tope de aridad. El último .then es el worker.

Reutilizar el pool

parallel hace spawn en el primer pull y mata los isolates cuando esa cadena termina. Dos trabajos pagan el arranque dos veces. IsolatePool es el corchete de spawn-una-vez. IsolatePool.using mata en el finally, incluso si el body lanza. Cancelar una cadena parallelOn no mata el pool — la siguiente cadena puede usarlo.

await IsolatePool.using(4, (pool) async {
  final a = await fx(batchA).parallelOn(pool, parseRow, chunk: 256).toList();
  final b = await fx(batchB).parallelOn(pool, parseRow, chunk: 256).toList();
  return (a, b);
});
Relacionado: concurrent — I/O, cualquier closure · mapConcurrent — la forma combinada de I/O · concurrent or parallel — I/O vs CPU · mapParallel — el mismo operador que parallel · fxPipe — compón workers para que dos etapas paguen un hop · ¿merece la pena parallel? — el mismo trabajo de cinco maneras, medido