Combinar flujos: combine, zip, flatMapLatest y merge, y la temporización que los separa
Los cuatro operadores de combinación se confunden constantemente porque sus firmas se parecen y sus resultados coinciden en el caso trivial de dos fuentes igual de rápidas. Sus diferencias reales no están en el tipo sino en la temporización: combine reacciona a cualquier cambio usando el último valor de los demás, zip empareja por posición y avanza al ritmo del más lento, flatMapLatest cancela el trabajo anterior cada vez que llega algo nuevo, y merge intercala sin sincronizar nada. Esta lección expone la mecánica de cada uno, los estados intermedios inconsistentes que produce combine y el criterio para no volver a elegir mal.
Elegir entre combine, zip, flatMapLatest y merge no es una cuestión de estilo: cada uno impone una disciplina temporal distinta y produce, con las mismas entradas, secuencias de salida distintas. La confusión persiste porque en el ejemplo de manual —dos flujos que emiten alternadamente y a la misma velocidad— los cuatro parecen razonables. En cuanto una fuente es más rápida que la otra, o una de ellas nunca emite, o el trabajo derivado tarda, las diferencias se vuelven brutales: combine produce combinaciones intermedias que nunca existieron conceptualmente, zip se bloquea esperando a un flujo que no tiene nada que decir, flatMapLatest cancela cosas a mitad y merge mezcla sin decirte de dónde vino cada valor.
- Enunciar la regla de disparo de cada operador y predecir la secuencia de salida ante entradas asíncronas.
- Reconocer los estados intermedios inconsistentes de
combiney decidir si son admisibles en cada dominio. - Usar
flatMapLatestsabiendo exactamente qué se cancela y qué efectos secundarios quedan a medias. - Escoger entre
mergey las alternativas concurrentes según si el orden y el origen importan.
combine y zip: reaccionar frente a emparejar
combine mantiene el último valor de cada fuente y emite cada vez que cualquiera de ellas produce algo, siempre que todas hayan emitido ya al menos una vez. zip empareja el primero con el primero, el segundo con el segundo, y no emite hasta tener ambos: su ritmo es el del más lento y termina cuando termina el más corto.
val a = flowOf(1, 2, 3).onEach { delay(100) }
val b = flowOf("x", "y").onEach { delay(250) }
combine(a, b) { n, s -> "$n$s" } // 2x, 3x, 3y (aprox., según tiempos)
a.zip(b) { n, s -> "$n$s" } // 1x, 2y (termina con el corto)
De ahí salen las dos consecuencias que importan. La primera: combine no espera y por tanto emite parejas que combinan un valor recién llegado con uno viejo de la otra fuente. Si esas dos cosas están correlacionadas —un identificador y los datos de ese identificador— generas estados imposibles: los datos del usuario anterior junto al nombre del nuevo. La segunda: zip sí espera, y por tanto si una fuente es un flujo caliente que rara vez emite, el emparejamiento se queda parado indefinidamente acumulando presión sobre la otra.
combine
Dispara con cualquier cambio, usa el último de los demás. Correcto cuando las fuentes son independientes entre sí y cada una representa una dimensión distinta del estado.
zip
Empareja por posición y avanza al ritmo del más lento. Correcto cuando el elemento n de una fuente y el n de la otra son dos mitades del mismo hecho.
flatMapLatest
Cada valor nuevo cancela el flujo derivado del anterior. Correcto cuando el trabajo en curso deja de tener sentido en cuanto cambia la entrada.
merge
Intercala varias fuentes del mismo tipo sin sincronizar ni transformar. Correcto cuando lo único que quieres es un canal único de sucesos homogéneos.
Conviene también saber que la lambda de combine es una función normal —no suspendida en la variante básica—, mientras que combineTransform sí permite suspender y emitir un número arbitrario de valores por disparo. Y que combine no emite nada en absoluto hasta que todas las fuentes hayan emitido su primer valor: si una de ellas es un SharedFlow sin repetición que aún no ha dicho nada, el resultado permanece mudo, y ese silencio se diagnostica siempre tarde.
Con dos fuentes correlacionadas, combine emite una combinación por cada cambio de cada una, incluidas las que mezclan generaciones distintas. Si esa mezcla es visible para el usuario, el operador correcto no es combine: es flatMapLatest desde la fuente que gobierna.
flatMapLatest y la cancelación como significado
La familia flatMapLatest, flatMapConcat y flatMapMerge transforma cada valor en un flujo nuevo y aplana el resultado; lo que las separa es qué hacen con el flujo anterior cuando llega otro valor. flatMapLatest lo cancela; flatMapConcat espera a que termine antes de empezar el siguiente; flatMapMerge los deja correr en paralelo hasta un límite de concurrencia.
val resultados: Flow<Resultado> = consulta
.debounce(300)
.distinctUntilChanged()
.flatMapLatest { texto -> repositorio.buscar(texto) }
La cancelación de flatMapLatest es su rasgo definitorio y también su riesgo. Cancela en el punto de suspensión, lo cual significa que el trabajo interrumpido puede dejar efectos a medias: una escritura parcial, un fichero abierto, un contador incrementado sin su decremento. Si el flujo interno tiene efectos secundarios no idempotentes, hace falta protegerlos con NonCancellable o, mejor, rediseñar para que el flujo interno sea puramente de lectura.
flowchart TD
S[Llega un valor de la fuente] --> Q{Que operador}
Q -->|combine| C[Emite con el ultimo de cada fuente]
Q -->|zip| Z[Espera al par de la otra fuente]
Q -->|flatMapLatest| L[Cancela el interno anterior y arranca el nuevo]
Q -->|flatMapConcat| K[Encola y espera a que termine el anterior]
Q -->|merge| M[Reemite en cuanto llega sin coordinar]El orden dentro de la cadena tampoco es indiferente. debounce antes de flatMapLatest reduce el número de flujos internos creados, mientras que ponerlo después filtra resultados ya calculados y desperdicia el trabajo. distinctUntilChanged antes evita relanzar una búsqueda idéntica; después, no evita nada. La regla general es que todo filtro que reduzca disparos debe ir lo más arriba posible, porque cada disparo suprimido es un flujo interno que no se crea.
Un caso concreto que aparece siempre: combinar un parámetro cambiante con una consulta. La forma incorrecta es combine entre el parámetro y el flujo de datos, que produce la mezcla de generaciones descrita antes; la forma correcta es flatMapLatest desde el parámetro hacia la consulta, porque expresa que la llegada de un parámetro nuevo invalida el trabajo en curso en lugar de sumarse a él.
merge y los tres errores clásicos
merge toma varios flujos del mismo tipo y reemite lo que llegue en cuanto llegue. No transforma, no empareja, no ordena: es la unión de sucesos. Su compañero natural es un tipo suma que identifique el origen, porque una vez mezclados los valores no llevan etiqueta.
sealed interface Suceso
data class Pulsado(val id: String) : Suceso
data class Recibido(val carga: ByteArray) : Suceso
val todos: Flow<Suceso> = merge(
pulsaciones.map(::Pulsado),
entrantes.map(::Recibido),
)
Los errores recurrentes son tres y todos consisten en pedirle a un operador algo que otro hace. El primero: usar combine para eventos. Como combine reemite con el último valor de las otras fuentes, un evento se vuelve a entregar cada vez que cambia cualquier otra dimensión, lo que produce acciones duplicadas —una navegación repetida, un mensaje que reaparece— sin que nadie las haya emitido dos veces.
El segundo: usar zip con un flujo caliente. Como zip exige un par completo y un flujo caliente no promete emitir nunca, la salida puede quedarse muda para siempre mientras la otra fuente sigue produciendo. El síntoma es una pantalla que nunca sale del estado de carga y unas trazas donde se ve claramente que una de las fuentes sí emitía.
El tercero: esperar concurrencia de flatMapConcat. Es secuencial por definición y su latencia total es la suma de las latencias internas; quien busca paralelismo necesita flatMapMerge con una concurrencia declarada, y quien busca solo lo último necesita flatMapLatest. Elegir flatMapConcat por defecto es la causa habitual de cadenas que se vuelven lentísimas cuando la fuente se acelera.
Para dos fuentes, dibuja en un papel las marcas de emisión de cada una y anota qué debería salir en cada instante. Ese dibujo elige el operador solo, y además se convierte casi literalmente en el test con tiempo virtual que confirmará que acertaste.
Cardinalidad, concurrencia y el ritmo de salida
Los cuatro operadores tienen variantes que amplían la aridad y controlan la concurrencia, y conviene conocerlas porque el código que las ignora acaba anidando combinaciones de dos en dos. combine admite hasta cinco fuentes con lambdas tipadas y, para más, una sobrecarga que recibe una colección y entrega un arreglo; merge acepta un Iterable completo. Anidar combine dentro de combine funciona, pero multiplica las emisiones intermedias, porque cada nivel dispara por su cuenta.
// Aridad alta sin anidar: una sola etapa de disparo.
val estado: Flow<Estado> = combine(
listOf(filtro, orden, pagina, sesion)
) { valores -> construirEstado(valores) }
// Concurrencia declarada: como máximo cuatro flujos internos a la vez.
val enriquecidos: Flow<Detalle> = ids.flatMapMerge(concurrency = 4) { cargar(it) }
Sobre el ritmo de salida hay dos operadores que suelen faltar en estas cadenas y que resuelven problemas que la gente intenta resolver con el operador de combinación equivocado. sample toma el último valor de cada ventana temporal y es la respuesta correcta cuando la fuente es más rápida que lo que el consumidor puede o quiere mostrar; debounce espera a que haya silencio y es la respuesta correcta cuando lo que importa es que el usuario haya terminado de escribir. Confundirlos produce, respectivamente, una interfaz que no se actualiza mientras haya actividad y una que se actualiza demasiado.
Hay finalmente una restricción estructural que explica por qué flatMapMerge y sus parientes concurrentes llevan la anotación de API experimental para el flujo: Flow garantiza que las emisiones al recolector son secuenciales, así que cualquier operador que ejecute cosas en paralelo debe volver a serializar la salida internamente. Eso no es gratuito y, sobre todo, significa que la concurrencia interna nunca se filtra al consumidor: por muchos flujos internos que corran a la vez, tu lambda de collect seguirá ejecutándose de uno en uno, y si el cuello de botella está ahí, subir la concurrencia no mejorará nada.
El criterio final, cuando la duda persiste, se puede reducir a una tabla mental de tres columnas: qué dispara la salida, qué ocurre con el trabajo anterior y qué pasa si una fuente calla para siempre. combine dispara con cualquiera, conserva lo anterior y tolera el silencio de una fuente solo si ya emitió una vez. zip dispara con el par completo, conserva el pendiente y se bloquea indefinidamente ante el silencio. flatMapLatest dispara con la fuente gobernante, cancela lo anterior y sobrevive sin problema al silencio de lo derivado. merge dispara con cualquiera, no tiene nada anterior que conservar y le da igual que alguna fuente no vuelva a hablar.
// La duda típica, resuelta por la relación y no por la forma:
combine(idioma, tema) { i, t -> Ajustes(i, t) } // dimensiones ortogonales
peticion.zip(indices) { p, i -> Numerada(p, i) } // dos mitades del mismo hecho
seleccion.flatMapLatest { repo.observar(it) } // una gobierna a la otra
merge(desdeRed, desdeUsuario) // solo coinciden en el tiempo
Insertar un buffer o un cambio de contexto entre las fuentes y el operador de combinación desacopla los ritmos y altera qué valores coinciden en el tiempo. La cadena que probaste con tiempo virtual puede comportarse distinto en cuanto añadas cualquiera de los dos.
Al mirarlos como utilidades intercambiables se pierde lo único que los distingue, que es la semántica temporal que cada uno presupone sobre el mundo. zip asume que las dos fuentes son dos relatos paralelos del mismo suceso, de modo que el elemento n de una y el n de la otra son piezas de una misma cosa y no tiene sentido hablar de la una sin la otra; esa suposición es exacta cuando las fuentes están generadas por el mismo proceso —dos proyecciones de una lista, una petición y su índice— y es catastróficamente falsa cuando una de ellas es el mundo exterior, porque el mundo exterior no lleva la cuenta de tus emparejamientos y no tiene ninguna obligación de emitir un elemento por cada uno de los tuyos. combine asume lo contrario: que las fuentes son dimensiones independientes de un estado global y que cualquier valor de una es compatible con cualquier valor de otra, lo que es exacto para cosas genuinamente ortogonales —el idioma y el tema visual, el filtro y el orden— y falso en cuanto hay una relación causal entre ellas, porque entonces el producto cartesiano que combine calcula sin pedir permiso incluye combinaciones que jamás existieron y que tu código de más abajo tratará como si fueran estados legítimos. flatMapLatest asume que hay una jerarquía, que una fuente manda y la otra obedece, y que la llegada de un valor nuevo en la que manda no añade información sino que invalida todo el trabajo derivado del anterior; esa es la suposición correcta para búsquedas, selecciones y navegaciones, y es también la más agresiva de las cuatro, porque cancelar en mitad de un flujo interno con efectos deja el mundo en un estado que solo tú puedes reparar. merge, finalmente, no asume nada: renuncia a toda relación entre las fuentes y se limita a la unión temporal, lo que lo convierte en el único honesto cuando lo que tienes son sucesos heterogéneos y en el único inútil cuando lo que necesitas es correlación. La consecuencia práctica es que la pregunta correcta nunca es cómo junto estos dos flujos, porque esa pregunta admite las cuatro respuestas y ninguna es obviamente mala; la pregunta correcta es qué relación hay entre los dos hechos que representan: si son el mismo hecho contado dos veces, empareja; si son dimensiones distintas del mismo estado, combina; si uno gobierna al otro, aplana con cancelación; y si no tienen ninguna relación salvo el instante en que ocurren, mézclalos y ponles una etiqueta.
- Monta dos flujos con retardos distintos, aplica
combineyzipsobre ellos y anota la secuencia real de salida. Compárala con la que predijiste antes de ejecutar. - Correlaciona deliberadamente las dos fuentes —un identificador y su detalle— con
combiney captura una emisión que mezcle generaciones. Reescríbela conflatMapLatesty comprueba que desaparece. - Aplica
zipentre un flujo frío finito y unSharedFlowque no emite. Documenta el bloqueo y explica en qué punto exacto se queda esperando. - Sustituye
flatMapLatestporflatMapConcaten una búsqueda con teclado y mide la latencia del último resultado a medida que escribes más rápido. Explica la curva. - Reescribe con
mergey un tipo suma dos fuentes de eventos que hoy manejes por separado, y comprueba si algún consumidor dependía silenciosamente del orden relativo entre ambas.