wandres.dev
FLOW Y STATEFLOW · streams en Kotlin

Flow: streams fríos que producen bajo demanda

Una función suspend resuelve devolver un valor que tarda; Flow resuelve devolver muchos valores que llegan a lo largo del tiempo. Esta lección construye el concepto desde su raíz: por qué Flow ocupa la casilla que faltaba en el cuadro de valores únicos y múltiples, síncronos y asíncronos; qué significa con precisión que un Flow sea frío —que su cuerpo es una receta apagada que se ejecuta entera y por separado para cada coleccionista, y que sin coleccionistas no consume nada—; cómo se construyen flujos con el constructor flow, con flowOf y asFlow, y cómo se tiende un puente desde una API de callbacks con callbackFlow y awaitClose; y por qué collect y los demás operadores terminales son el interruptor que enciende la producción. Cierra explicando que la contrapresión, la preservación de contexto y la cancelación de un Flow no son funciones añadidas sino consecuencias directas de que emitir y recolectar sean operaciones suspend dentro de la concurrencia estructurada.

⏱ 17 min

Una función suspend resuelve un problema muy concreto: devolver un valor que tarda sin bloquear el hilo que espera. Pero buena parte de una aplicación no va de valores que tardan, va de valores que siguen llegando: la tabla que cambia bajo tus pies, el sensor que emite cada cien milisegundos, el campo de búsqueda que el usuario teclea letra a letra, el socket que empuja mensajes hasta que se cierra. Para eso existe Flow: una secuencia de valores producidos a lo largo del tiempo, sin bloquear a nadie. Y antes de aprender un solo operador conviene interiorizar su rasgo definitorio, porque explica casi todo lo que viene después: un Flow es frío. No es un chorro que ya corre y al que te asomas; es una receta apagada que nadie ejecuta hasta que alguien decide recolectarla, y que se ejecuta entera, desde el principio y por separado, para cada quien la recolecte.

🎯 Al terminar esta lección sabrás
  • Situar Flow en el cuadro de valores únicos y múltiples, síncronos y asíncronos, y ver qué casilla ocupa.
  • Comprender con precisión qué significa que un flujo sea frío y qué consecuencias tiene para memoria y recursos.
  • Construir flujos con el constructor flow, con flowOf y asFlow, y puentear callbacks con callbackFlow.
  • Reconocer collect y los terminales como el interruptor, y derivar contrapresión y cancelación de la suspensión.

Un valor no basta: la casilla que faltaba

Ordena el terreno con dos preguntas: ¿cuántos valores devuelve esto, uno o muchos? y ¿puede tardar sin bloquear el hilo, o no? El cruce da cuatro casillas y Kotlin las llena todas. Un valor síncrono es un T corriente. Un valor asíncrono es una función suspend que devuelve T. Muchos valores síncronos son una List<T>, que los calcula todos y los guarda, o una Sequence<T>, que los produce bajo demanda pero bloqueando el hilo en cada paso. La cuarta casilla —muchos valores, producidos bajo demanda, sin bloquear— estuvo vacía en el lenguaje hasta que llegó Flow.

Esa vacante importa porque las tres primeras respuestas fallan justo donde más falta hacen. Una List te obliga a esperar a que el último elemento exista antes de ver el primero, y a tenerlos todos en memoria a la vez: inviable para un histórico paginado o para un stream que no termina nunca. Una Sequence produce perezosamente, que es lo que quieres, pero su iterator no es suspend, así que cada valor que tarde bloquea el hilo que lo pide: inviable en Android. Flow es exactamente la Sequence que sí puede suspender.

// Un valor que tarda: una sola respuesta y se acabo
suspend fun cargarUsuario(id: String): Usuario = api.usuario(id)

// Muchos valores que tardan: una secuencia asincrona que no termina
fun mensajesDe(sala: String): Flow<Mensaje> = flow {
    var cursor = 0L
    while (true) {
        val lote = api.mensajesDesde(sala, cursor)   // suspende, no bloquea
        lote.forEach { mensaje -> emit(mensaje) }    // los entrega uno a uno
        cursor = lote.lastOrNull()?.instante ?: cursor
        delay(1_000)                                 // suspende de nuevo
    }
}

Lee el segundo con atención, porque hay más de lo que parece. El cuerpo del constructor flow es una lambda suspend: dentro puedes esperar a la red, dormir con delay, abrir una transacción o llamar a cualquier otra función suspendida. emit también es suspend, y eso es lo que permite que el productor y el consumidor se acompasen sin buffers infinitos. Y el bucle while (true) no es una temeridad: un Flow no tiene obligación de terminar, igual que una conversación no la tiene.

🎯

Un valor, sincrono

Un T corriente. Lo tienes en el acto o no lo tienes; si tarda, el hilo se queda esperando y en Android eso se traduce en una interfaz congelada.

Un valor, asincrono

Una función suspend que devuelve T. Puede tardar sin bloquear, pero responde una sola vez: sirve para cargar, no para observar.

📚

Muchos valores, sincronos

List los calcula todos y los guarda; Sequence los produce perezosamente pero bloqueando en cada paso. Ninguna de las dos sabe esperar sin ocupar el hilo.

🌊

Muchos valores, asincronos

Flow. Produce bajo demanda, uno a uno, suspendiendo entre valor y valor. Es la Sequence que aprendió a esperar sin bloquear a nadie.

La tabla no es un adorno pedagógico: es una herramienta de diseño que se usa a diario. Cuando dudes qué debe devolver una función de tu capa de datos, responde primero a las dos preguntas —cuántos valores, y si puede tardar— y el tipo aparece solo. Confundir la casilla es el origen de la mayoría de las capas de datos incómodas: un suspend que devuelve List allí donde el dato cambia constantemente obliga a la interfaz a preguntar en bucle, y un Flow de un solo valor allí donde bastaba un suspend añade ceremonia sin comprar nada.

Frío: la receta y el plato

Aquí está la idea que hay que dejar bien asentada. Cuando escribes val f = mensajesDe("general") no ha ocurrido nada. No se ha abierto una conexión, no se ha llamado a la red, no hay ninguna corrutina viva, no se ha reservado un solo byte para valores. Lo único que existe es un objeto que guarda la lambda que sabe producir. Es una receta escrita en una ficha, no un plato en el fuego. El fuego lo enciende collect, y lo enciende para quien colecciona, no para todos.

De ahí se sigue la consecuencia más contraintuitiva para quien viene del mundo de los eventos: si dos coleccionistas recogen el mismo Flow, el cuerpo se ejecuta dos veces, entero, desde el principio, en dos corrutinas distintas y sin compartir absolutamente nada. Dos llamadas a la red, dos cursores independientes, dos secuencias de emisiones. No hay difusión: hay duplicación. Quien espera que el segundo coleccionista se enganche a lo que ya estaba pasando está pensando en un flujo caliente, y eso es otra cosa —StateFlow, SharedFlow— que estudiaremos en las lecciones siguientes.

flowchart LR
DEF[Definicion del flow apagada sin ejecutar] --> C1[Coleccionista uno]
DEF --> C2[Coleccionista dos]
C1 --> E1[Ejecucion propia desde el principio]
C2 --> E2[Otra ejecucion independiente desde el principio]
style DEF fill:#89b4fa,color:#11111b
style E1 fill:#a6e3a1,color:#11111b
style E2 fill:#a6e3a1,color:#11111b

Conviene ver el mecanismo por dentro, porque desmitifica el asunto. Un Flow no es más que un objeto con un único método que recibe un coleccionista y ejecuta la lambda contra él; emit es literalmente una llamada al coleccionista que le pasa el valor. No hay planificador oculto, no hay hilo propio, no hay cola: hay una función que llama a otra función, dentro de la corrutina de quien colecciona. Por eso la ejecución ocurre “en” el coleccionista y por eso dos coleccionistas son dos ejecuciones. Toda la magia aparente del stream se reduce a una inversión de control que cabe en una interfaz de un solo método.

Antes de leer la duplicación como un defecto, mira lo que compra. Un Flow frío no tiene estado que corromper, no acumula valores que nadie lee, no mantiene suscripciones vivas ni recursos abiertos, no puede filtrar memoria mientras nadie lo colecciona y no necesita que lo apagues, porque nunca estuvo encendido. Es declarativo: describe qué valores existirían si alguien los pidiera. Por eso un repositorio puede devolver flujos alegremente, guardarlos en propiedades y pasarlos por parámetros sin coste alguno; la factura solo llega en el punto exacto donde alguien decide consumirlos.

Conviene, eso sí, no confundir frío con inofensivo. Que la ejecución se repita por coleccionista significa que también se repiten sus efectos: si el cuerpo del flujo escribe en una caché, incrementa un contador o dispara una petición de pago, cada colección lo hará otra vez. La frialdad garantiza que nada ocurre antes del terminal, no que lo que ocurra después sea idempotente. La disciplina que evita el susto es la de siempre: dentro de un constructor de flujo, produce y no decidas; deja los efectos observables para el consumidor, que sabe cuántas veces está consumiendo.

💡
Devolver un Flow es devolver una promesa de trabajo, no trabajo hecho

Una función que devuelve Flow no necesita ser suspend: no hace nada, solo entrega la receta. Por eso verás fun observarTareas(): Flow<List<Tarea>> sin suspend, mientras que suspend fun cargarTareas(): List<Tarea> sí lo lleva. Esa diferencia de firma no es cosmética: te dice, de un vistazo, si la llamada cuesta algo ahora o solo cuando la colecciones. Cuando diseñes tus propias capas de datos, deja esa señal limpia.

Construir, recolectar y apagar

Los constructores son pocos y cubren casi todo. flowOf(a, b, c) envuelve valores ya conocidos. El método asFlow convierte una colección, un rango o una secuencia que ya tienes. El constructor flow es el general: dentro emites tantas veces como quieras y suspendes cuando haga falta. Y callbackFlow es el puente hacia el pasado imperativo, el que convierte una API de escuchadores —sensores, geolocalización, un cliente de terceros— en un stream que participa de la concurrencia estructurada.

fun ubicaciones(cliente: LocationClient): Flow<Ubicacion> = callbackFlow {
    val escucha = object : LocationListener {
        override fun onLocation(u: Ubicacion) { trySend(u) }   // no suspende
        override fun onError(e: Throwable) { close(e) }
    }
    cliente.registrar(escucha)
    awaitClose { cliente.desregistrar(escucha) }   // se ejecuta al cancelar
}

// Terminales: cada uno enciende la produccion de una forma distinta
suspend fun ejemplos(f: Flow<Ubicacion>) {
    f.collect { u -> println(u) }        // consume todo hasta que termine
    val primera = f.first()              // consume una y cancela el resto
    val todas = f.toList()               // solo si el flujo termina
}

Fíjate en trySend dentro del escuchador: los callbacks de una API imperativa no son suspend, así que no pueden esperar a que haya sitio en el búfer; trySend intenta entregar sin suspender y te devuelve si lo consiguió. Esa asimetría es la razón de que exista callbackFlow como constructor aparte y no baste con el flow general.

awaitClose es la pieza que casi todo el mundo olvida y la que hace que el puente sea seguro: mantiene viva la corrutina productora hasta que el coleccionista se va, y entonces —y solo entonces— desregistra el escuchador. Sin ella tendrías la fuga clásica: un listener registrado para siempre contra una pantalla que ya no existe. Con ella, cancelar la corrutina que colecciona limpia el recurso automáticamente, porque la cancelación viaja hacia arriba por el mismo árbol que estudiaste en el nivel anterior.

El otro protagonista es collect, y conviene nombrarlo con precisión: es un operador terminal, una función suspend que ejecuta el flujo y no vuelve hasta que este se agota o alguien lo cancela. Los demás terminales —first, toList, fold, single— son variantes de lo mismo con una condición de parada distinta. Todo lo que no es terminal es intermedio, y los intermedios no ejecutan nada: solo devuelven un Flow nuevo que envuelve al anterior, tan apagado como él. Un flujo con diez operadores encadenados y ningún terminal sigue siendo una ficha de receta.

Hay un terminal con un papel especial, launchIn, que merece nombrarse aparte porque resuelve una incomodidad muy frecuente. Como collect suspende hasta que el flujo se agota, coleccionar dos flujos seguidos en la misma corrutina nunca funciona: el segundo collect no se alcanza jamás. launchIn lanza la colección en un scope y devuelve el control inmediatamente, de modo que puedes arrancar varias suscripciones independientes sin anidar bloques ni escribir un launch por cada una.

fun observar(scope: CoroutineScope) {
    repo.tareas()
        .onEach { lista -> pintar(lista) }        // lo que harias en el collect
        .flowOn(Dispatchers.IO)                   // cambia el contexto aguas arriba
        .catch { e -> log("fallo el stream", e) } // captura lo de arriba, no lo de abajo
        .onCompletion { causa -> log("fin: $causa") }
        .launchIn(scope)                          // terminal que no suspende al llamante
}

Dos matices sobre ese fragmento valen tanto como el resto de la lección. flowOn afecta solo a lo que está encima de él en la cadena, nunca a lo de debajo, y por eso puede aparecer varias veces con dispatchers distintos: cada tramo produce donde le conviene y el coleccionista sigue recibiendo en su propio contexto. Y catch obedece a la misma direccionalidad: intercepta las excepciones que vienen de aguas arriba y deja intactas las que lance tu propio consumidor, lo cual evita el error clásico de tragarse un fallo de la interfaz creyendo que se estaba protegiendo la red. La transparencia de excepciones —que un operador solo responda por lo que tiene encima— es lo que permite razonar sobre una cadena leyéndola por tramos en vez de entera.

⚠️
No emitas desde otra corrutina dentro del constructor flow

El constructor flow exige preservación de contexto: debes emitir desde la misma corrutina que ejecuta el cuerpo, y si intentas hacerlo desde un withContext interno o desde una corrutina lanzada aparte, la biblioteca lanza IllegalStateException. No es una arbitrariedad, es lo que garantiza que el contexto del coleccionista mande. Si necesitas producir en otro dispatcher, el operador es flowOn, que cambia el contexto aguas arriba sin romper la regla; si necesitas emitir desde varias corrutinas a la vez, el constructor correcto es channelFlow, que para eso te da send y trySend.

ℹ️
La contrapresion te sale gratis

Como emit es suspend y collect procesa cada valor antes de pedir el siguiente, un productor rápido frente a un consumidor lento simplemente se detiene: la emisión número tres no ocurre hasta que la dos se ha consumido. Eso es contrapresión, y no la implementa ningún componente aparte —ni colas, ni estrategias, ni tokens—, la implementa la propia suspensión. Cuando esa sincronía estricta te estorbe, tienes buffer para desacoplar productor y consumidor, conflate para quedarte solo con el último valor y collectLatest para cancelar el procesamiento anterior si llega uno nuevo. Elegir entre ellos es elegir qué prefieres perder cuando no puedes seguir el ritmo.

La frialdad es lo que convierte un stream en una expresion

Detente en lo que realmente significa que un Flow sea frío, porque es la decisión de diseño de la que cuelga todo lo demás. Un flujo caliente es un proceso: ya está ocurriendo, tiene identidad, tiene estado interno, tiene un antes y un después, y por lo tanto conectarse a él es un acto con historia —llegas tarde o pronto, te pierdes cosas, alteras algo—. Un flujo frío no es un proceso: es una expresión, una descripción de qué valores existirían bajo qué condiciones. Y las expresiones tienen una propiedad maravillosa que los procesos no tienen: se pueden componer sin miedo, porque componerlas no las ejecuta. Puedes encadenarle diez operadores, guardarla en una variable, devolverla desde tres capas distintas, pasarla como argumento y no ha pasado nada todavía; el mundo sigue exactamente igual que antes. Esa es la razón profunda de que la programación reactiva moderna sea legible mientras que la de los años del bus de eventos era un laberinto: cuando componer no ejecuta, puedes razonar sobre el qué antes de preocuparte por el cuándo, y separar esas dos preguntas es la mitad del trabajo intelectual de un programa concurrente. La contrapartida —que dos coleccionistas signifiquen dos ejecuciones— no es un accidente que haya que parchear, es el precio exacto y honesto de esa pureza; y cuando de verdad necesites compartir una sola ejecución entre varios observadores, el remedio no será renunciar a la frialdad, sino aplicar deliberadamente un operador que la convierta en calor bajo un scope que tú eliges. Frío por defecto, caliente por decisión explícita: ese orden, y no el contrario, es lo que hace que el sistema entero se pueda razonar.

⚔️ Enciende y apaga tus propios flujos
  1. Escribe una función que devuelva Flow<Int> emitiendo los cinco primeros cuadrados con un delay entre cada uno, y demuestra con impresiones que el cuerpo no se ejecuta hasta el collect.
  2. Colecciona ese mismo flujo dos veces desde dos corrutinas distintas y explica, con la salida en la mano, por qué el cuerpo se ejecutó dos veces enteras.
  3. Convierte una API de callbacks que conozcas —un sensor, un escuchador de red, un cliente de terceros— en un callbackFlow con awaitClose, y describe qué fuga aparecería si omitieras esa última línea.
  4. Reemplaza collect por first y razona qué le ocurre a la corrutina productora justo después de la primera emisión.
  5. Justifica por qué una función que devuelve Flow no necesita ser suspend, y qué información le da esa firma a quien lee tu capa de datos sin abrirla.