wandres.dev
FLOW I · streams fríos

Constructores y recolección

El constructor flow no es una función mágica sino una lambda con receptor suspendido, y esa única observación explica de dónde sale emit, por qué no se puede llamar desde fuera y por qué el productor y el recolector comparten corrutina. Esta lección recorre los constructores de la librería, distingue las operaciones terminales de todo lo demás por su firma, y establece la relación exacta de uno a uno entre cada recolección y una ejecución completa e independiente del cuerpo productor.

⏱ 20 min

La pregunta que descoloca a quien empieza es de dónde sale emit. No es una función global, no está importada, no es miembro de ninguna clase visible, y sin embargo aparece disponible dentro del bloque del constructor y desaparece en cuanto se sale de él. La respuesta no tiene nada que ver con los flujos y sí con una construcción del lenguaje que ya conoces desde el nivel de los dominios específicos: el bloque es una lambda con receptor, el receptor es un FlowCollector, y emit es sencillamente su único método. Entender eso convierte todo el catálogo de constructores en variaciones aburridas de una misma idea, y prepara la observación central de la lección, que es que recolectar y ejecutar son exactamente la misma cosa contada dos veces.

🎯 Al terminar esta lección sabrás
  • Leer la firma del constructor flow y deducir de ella la disponibilidad de emit dentro del bloque.
  • Distinguir los constructores de la librería y elegir el adecuado según de dónde vengan los valores.
  • Separar las operaciones terminales del resto por su firma y saber cuáles son suspendidas.
  • Enunciar y demostrar la relación de uno a uno entre cada recolección y una ejecución completa del productor.

El constructor y su receptor implícito

La firma del constructor no esconde nada. Recibe un bloque que es, a la vez, suspendido y con receptor, y ese receptor es la interfaz consumidora de la lección anterior.

public fun <T> flow(block: suspend FlowCollector<T>.() -> Unit): Flow<T>

Léase el tipo del parámetro de derecha a izquierda. Es una función sin argumentos que devuelve Unit, marcada como suspendida, y con un receptor de tipo FlowCollector parametrizado. Dentro de un bloque con receptor, los miembros del receptor están disponibles sin cualificar, y el único miembro que tiene un FlowCollector es emit. No hay ninguna magia del compilador: hay una lambda con receptor, igual que en cualquier constructor de dominio específico que hayas escrito.

De ahí se siguen dos hechos que suelen enseñarse como reglas arbitrarias. El primero es que emit solo existe dentro del bloque, porque fuera de él no hay ningún receptor de ese tipo. El segundo es que emit puede suspender, porque el bloque entero está marcado como suspendido y por tanto se compila a una función con su continuación.

val temperaturas: Flow<Int> = flow {
    var t = 20
    repeat(5) {
        delay(100)      // se puede suspender: el bloque es suspendido
        emit(t++)       // emit es miembro del receptor implicito
    }
}

Hay una regla más que sí es específica de los flujos y que conviene enunciar ya aunque su explicación completa espere a la última lección. El bloque debe emitir siempre desde su propia corrutina, y por tanto no se puede lanzar una corrutina hija dentro del bloque y emitir desde ella. Si se intenta, la propia librería lo detecta en tiempo de ejecución y falla con un mensaje sobre la invariante del flujo. El constructor alternativo para esos casos, que permite emisión concurrente a cambio de perder algunas garantías, se llama channelFlow y usa send en lugar de emit.

⚠️
Un bloque productor no es un ambito de corrutinas

Dentro del bloque hay una corrutina, la del recolector, y solo una. No es un CoroutineScope, así que no puedes llamar a launch directamente, y si lo consigues envolviéndolo en un ámbito propio, emitir desde el hijo lanzará una excepción. Esa restricción parece incómoda hasta que se entiende que es justo lo que hace que emit sea una llamada directa sin cola y sin bloqueo.

Los otros constructores

El resto del catálogo son atajos, y todos ellos podrían escribirse con el constructor general en dos líneas. Saberlo evita memorizarlos.

flowOf(1, 2, 3)                     // valores conocidos de antemano
listOf("a", "b").asFlow()           // desde cualquier Iterable
(1..100).asFlow()                   // desde un rango
generateSequence(1) { it * 2 }.asFlow()
emptyFlow<Int>()                    // termina sin emitir nada
flow { emit(cargarDeRed()) }        // desde una sola llamada suspendida

La última línea merece un comentario porque es el caso más frecuente y el peor entendido. Envolver una única llamada suspendida en un flujo no aporta nada por sí mismo: si solo va a haber un valor, la función suspendida ya era suficiente. Lo que sí aporta es la frialdad con reintento, es decir, la posibilidad de que alguien aplique los operadores de reintento y de captura del último tema del nivel sobre una operación que se puede repetir entera. El criterio para decidir es si tiene sentido ejecutar esa operación más de una vez.

Existe también un constructor pensado para adaptar mundos que no suspenden, y es el puente natural entre los flujos y las interfaces de devolución de llamada que ya estudiamos.

fun ubicaciones(): Flow<Ubicacion> = callbackFlow {
    val oyente = OyenteDeUbicacion { u -> trySend(u) }
    servicio.registrar(oyente)
    awaitClose { servicio.desregistrar(oyente) }
}

Ese constructor no es un atajo del general: pertenece a la familia de channelFlow porque internamente crea un canal, y por eso ofrece trySend en lugar de emit. El bloque final es obligatorio y su ausencia se detecta en ejecución, porque sin él la suscripción al mundo externo quedaría viva para siempre después de que la recolección termine.

🏗️

flow

El general. Bloque suspendido con receptor, emisión secuencial desde la misma corrutina. Es el que sirve para casi todo.

📃

flowOf y asFlow

Azúcar sobre el anterior para valores ya disponibles en memoria. No suspenden nunca y no aportan asincronía.

📻

channelFlow

Permite emitir desde varias corrutinas a cambio de introducir un canal y por tanto un búfer y una frontera.

🔌

callbackFlow

Especialización del anterior para envolver interfaces de devolución de llamada. Exige cerrar el registro con awaitClose.

Recolectar es ejecutar

La firma vuelve a decidirlo todo. Una operación es terminal si es suspendida y devuelve algo que no es un flujo; es intermedia si no es suspendida y devuelve un flujo. No hay una tercera categoría ni hace falta memorizar listas, porque la firma no miente.

suspend fun <T> Flow<T>.collect(action: suspend (T) -> Unit)   // terminal
suspend fun <T> Flow<T>.toList(): List<T>                      // terminal
suspend fun <T> Flow<T>.first(): T                             // terminal
suspend fun <T> Flow<T>.single(): T                            // terminal
fun <T> Flow<T>.launchIn(scope: CoroutineScope): Job           // terminal, no suspende

La última rompe el patrón por una buena razón: en lugar de recolectar aquí y ahora, lanza una corrutina que recolecta y devuelve su trabajo. No es una excepción a la regla sino un envoltorio, y equivale exactamente a lanzar un launch cuyo cuerpo llama a collect. Que devuelva un Job en lugar de suspender es lo que la hace apta para arrancar recolecciones desde código que no suspende.

Ahora la observación central. Cada llamada a una operación terminal ejecuta el cuerpo productor entero, desde su primera línea, con un colector recién construido. No hay caché, no hay reparto, no hay progreso compartido.

val f = flow { println("ejecuto"); emit(1) }

val a = f.toList()      // ejecuto
val b = f.first()       // ejecuto
launch { f.collect { } } // ejecuto
launch { f.collect { } } // ejecuto

Cuatro operaciones terminales, cuatro ejecuciones completas e independientes. Si el cuerpo hace una petición de red, se harán cuatro peticiones. Esto no es un defecto que haya que remediar con cuidado, es la definición de frío, y las herramientas para compartir una sola ejecución entre varios recolectores son un asunto distinto que corresponde a los flujos calientes del nivel siguiente.

Hay una operación terminal que ilustra bien la relación y que sorprende al medirla. Cuando alguien pide solo el primer valor, el flujo no se ejecuta entero: en cuanto ese valor llega, la recolección se detiene y el productor recibe la señal de abandono que estudiaremos en la lección siguiente. Un flujo infinito, por tanto, admite perfectamente esa operación terminal y termina en un instante, algo que sería imposible si recolectar significase materializar.

val infinito = flow { var n = 0; while (true) emit(n++) }
val primero = infinito.first()   // termina de inmediato, vale 0

La segunda mitad de la observación es igual de importante y menos evidente: el productor y el recolector viven en la misma corrutina. Cuando el bloque llama a emit, esa llamada entra directamente en el cuerpo del recolector, ejecuta lo que haya allí, y solo cuando termina vuelve al productor. La pila de llamadas es una sola y se puede leer entera en un depurador.

sequenceDiagram
participant C as Llamante
participant P as Bloque productor
participant R as Cuerpo del recolector
C->>P: collect inicia la ejecucion
P->>R: emit del primer valor
R-->>P: el cuerpo termina y emit retorna
P->>R: emit del segundo valor
R-->>P: el cuerpo termina y emit retorna
P-->>C: el bloque acaba y collect retorna

Hay una consecuencia práctica que se toca con la mano al medir. Si el productor tarda cien milisegundos por valor y el recolector otros cien, el total de diez valores son dos segundos y no uno, porque nada se solapa. Esa suma es el precio exacto de no tener búfer, y la lección cuarta enseña a pagarlo o a evitarlo de forma consciente. Antes de llegar allí conviene interiorizar que el estado por defecto es el secuencial estricto, porque casi todos los errores de rendimiento con flujos vienen de haber supuesto lo contrario.

Que el productor y el recolector compartan corrutina convierte una cadena asincrona en algo que se puede leer como una funcion ordinaria y esa legibilidad es el rendimiento real del diseno

Lo más difícil de comunicar sobre los flujos es que su principal virtud no aparece en ningún banco de pruebas. Que emit sea una llamada directa al cuerpo del recolector, dentro de la misma corrutina y sobre la misma pila lógica, tiene un efecto medible y modesto: evita colas, evita objetos intermedios y evita saltos entre hilos que no hacen falta. Pero tiene otro efecto que no se mide y que decide cuánto cuesta mantener un programa durante años. Significa que la cadena entera, con sus constructores y sus operadores y su recolector final, es una única unidad de ejecución que se puede razonar como se razona una función corriente. Un try colocado alrededor de un collect captura de verdad lo que falla arriba, porque no hay ninguna frontera entre medias por la que la excepción tuviera que ser reempaquetada y transportada. Un finally se ejecuta de verdad cuando la recolección acaba o se cancela, porque la cancelación es la de la corrutina que lo contiene y no un mensaje enviado a otra parte. Una traza de pila muestra de verdad la cadena entera, con el operador y el productor y el consumidor, y no un fragmento decapitado que empieza en el planificador. Una variable local declarada antes del bloque sigue siendo accesible dentro sin trucos, porque es una lambda ordinaria capturada como cualquier otra. Cada una de esas propiedades parece pequeña y ninguna aparecería en una comparativa de operaciones por segundo, pero juntas son la diferencia entre un fallo que se depura en veinte minutos y uno que se depura en tres días. La historia del código asíncrono es la historia de haber ido perdiendo esas propiedades una a una, a cambio de capacidades reales, y de haberse acostumbrado a la pérdida hasta considerarla el precio inevitable de la asincronía. Los flujos demuestran que no lo era: la mayor parte de aquella pérdida venía de que cada librería introducía sus propias fronteras de ejecución, y basta con no introducirlas para recuperarlo casi todo. Cuando en la cuarta lección aprendas a insertar una frontera deliberadamente con un búfer, recuerda que estás gastando exactamente esto, y que la pregunta correcta no es cuánta latencia ganas sino cuánta legibilidad estás dispuesto a vender por ella.

⚔️ Demuestra la relacion uno a uno
  1. Escribe un flujo cuyo cuerpo incremente un contador global e imprima su valor. Aplícale cuatro operaciones terminales distintas y predice el número final antes de ejecutarlo.
  2. Construye la misma secuencia con el constructor general, con flowOf y con asFlow. Escribe cuál elegirías en cada caso y por qué.
  3. Mide el tiempo total de un flujo con productor y consumidor lentos, y comprueba que es la suma y no el máximo. Guarda el número para la lección cuarta.
  4. Intenta emitir desde un launch dentro del bloque productor. Lee el mensaje de error completo y reescribe la solución con channelFlow.
  5. Envuelve una interfaz de devolución de llamada real con callbackFlow, omite el bloque de cierre a propósito y observa qué ocurre. Después arréglalo con awaitClose.