mergeMap, concatMap & exhaustMap

Tres respuestas más a «llegó un evento mientras el anterior sigue en marcha»: ejecutarlos todos, ejecutarlos en orden, o ignorar el nuevo.

FxEvents<R> FxEvents<T>.mergeMap<R>(Stream<R> Function(T) f, {int? concurrent}) FxEvents<R> FxEvents<T>.concatMap<R>(Stream<R> Function(T) f) FxEvents<R> FxEvents<T>.exhaustMap<R>(Stream<R> Function(T) f)

Lección

Mapear un evento a un stream interno —una petición, una subida, una consulta— plantea una pregunta que un pipeline pull nunca tiene que responder: ¿qué pasa cuando el siguiente evento llega antes de que el último stream interno haya terminado? Hay exactamente cuatro políticas sensatas, y elegir la equivocada es donde vive la mayoría de los bugs reactivos. switchMap es la respuesta de «gana el último»; estas tres son las otras tres.

mergeMap(f) ejecuta todos los streams internos a la vez e intercala su salida por orden de llegada. Úsalo cuando todos los resultados importan y ninguno reemplaza a otro: subir tres ficheros, abanicarse hacia tres servicios. Con concurrent: n se ejecutan como mucho n a la vez y el resto espera en cola, que es como evitas que un abanico abra doscientos sockets.

concatMap(f) los ejecuta estrictamente en orden, cada uno hasta terminar antes de que empiece el siguiente. Nada se solapa y nada se descarta, así que un stream interno lento atasca la cadena entera; eso es justo lo que quieres cuando el orden es la condición de corrección, como en «aplica estas ediciones en secuencia».

exhaustMap(f) se queda con el primero e ignora el resto: mientras un stream interno está en marcha, los eventos entrantes se descartan sin más — ni encolados ni cancelados. Esta es la protección contra el doble envío. Un segundo toque en un botón cuya petición sigue en vuelo no hace absolutamente nada, que es exactamente lo que quieres cuando la petición es POST /orders.

Capa de eventos de fxdart, siguiendo a flatMap, flatMap(maxConcurrent: 1) y exhaustMap de Rx. Al primero se le llama aquí mergeMap porque flatMap ya significa aplanar iterables en el lado pull.

Demo 1 · mergeMap — todo a la vez

Demo 2 · exhaustMap — la protección contra el doble envío

Pruébalo tú

Ejercicio: el orden de concatMap y un abanico acotado.

Relacionado: switchMap — la cuarta política: gana el más nuevo, el resto se cancela · mapConcurrent — abanico acotado del lado pull, donde los resultados mantienen el orden · debounce — a menudo la mejor solución: corta los eventos de más antes de que se conviertan en streams internos