wandres.dev
FLOW I · streams fríos

Operadores intermedios

Un operador intermedio no está suspendido y devuelve un flujo, y esas dos marcas de su firma bastan para demostrar que no puede ejecutar nada: solo sabe envolver una receta dentro de otra. Esta lección reconstruye map, filter y onEach sobre el operador primitivo transform, muestra que una cadena es un colector dentro de otro colector y no una tubería con etapas, y explica cómo take detiene al productor mediante una excepción que viaja hacia arriba.

⏱ 21 min

Las cadenas de operadores tienen un aire engañoso de tubería industrial, con etapas que procesan y pasan al siguiente puesto, y esa imagen mental produce dos creencias falsas a la vez: que algo empieza a ocurrir cuando se escribe la cadena, y que entre etapa y etapa hay un sitio donde los valores esperan. Ninguna de las dos es cierta, y la demostración no requiere leer el fuente de la librería, solo mirar la firma de cualquier operador. Un operador intermedio no está marcado como suspendido, lo cual significa que no puede llamar a collect, lo cual significa que no puede producir ni consumir ni un solo valor. Todo lo que le queda por hacer es construir un objeto nuevo que recuerde qué habría que hacer más tarde, y eso es exactamente lo que hace.

🎯 Al terminar esta lección sabrás
  • Clasificar cualquier operador como intermedio o terminal leyendo únicamente su firma.
  • Reconstruir map, filter y onEach sobre el operador primitivo transform.
  • Describir una cadena de operadores como un anidamiento de colectores y no como una secuencia de etapas.
  • Explicar el mecanismo por el que take detiene al productor y qué código de usuario puede romperlo.

La firma delata al operador

Pongamos tres declaraciones juntas y leámoslas como quien busca una sola palabra.

fun <T, R> Flow<T>.map(transform: suspend (T) -> R): Flow<R>
fun <T> Flow<T>.filter(predicate: suspend (T) -> Boolean): Flow<T>
suspend fun <T> Flow<T>.toList(): List<T>

Las dos primeras no llevan suspend y devuelven un flujo. La tercera lleva suspend y devuelve datos. Como el único método de la interfaz es suspendido, una función que no lo es no puede invocarlo, ni directamente ni a través de nadie: el compilador se lo impide. Por tanto map y filter no han visto ningún valor cuando retornan, y no podrían haberlo visto aunque quisieran. Lo único que pueden devolver es una descripción nueva.

val f = flowOf(1, 2, 3)
    .map { println("mapeando $it"); it * 2 }
    .filter { println("filtrando $it"); it > 2 }
// Aqui no se ha impreso absolutamente nada

f.collect { println(it) }   // ahora si: mapeando 1, filtrando 2, ...

Nótese además el orden de las impresiones cuando por fin se recolecta. No aparecen los tres mensajes de mapeo seguidos de los tres de filtrado, como sugeriría la imagen de la tubería con etapas, sino un mapeo y un filtrado alternándose valor a valor. Cada elemento recorre la cadena completa antes de que el siguiente se produzca, y esa es la primera pista de que no hay etapas sino anidamiento.

💡
La pregunta que resuelve cualquier duda

Cuando no sepas si un operador ejecuta algo, mira si su firma lleva suspend. Si no la lleva, es imposible que ejecute nada del flujo, sin excepciones y sin casos raros. Si la lleva y devuelve un flujo, estás ante uno de los poquísimos híbridos de la librería y conviene leer su documentación con calma.

transform es el operador primitivo

De todo el catálogo intermedio, uno solo hace falta. transform recibe un bloque suspendido con receptor de tipo FlowCollector, exactamente como el constructor general, y para cada valor de entrada permite emitir ninguno, uno o muchos. Los demás son casos particulares suyos.

// Conceptualmente, cada uno cabe en una linea
fun <T, R> Flow<T>.map(f: suspend (T) -> R): Flow<R> =
    transform { emit(f(it)) }

fun <T> Flow<T>.filter(p: suspend (T) -> Boolean): Flow<T> =
    transform { if (p(it)) emit(it) }

fun <T> Flow<T>.onEach(a: suspend (T) -> Unit): Flow<T> =
    transform { a(it); emit(it) }

Las tres líneas cuentan la historia entera. map emite siempre uno, transformado. filter emite cero o uno, sin tocarlo. onEach emite siempre el mismo, después de hacer algo con él; es el operador para efectos secundarios y su valor es que no altera el flujo, lo cual lo convierte en el instrumento natural de registro y depuración dentro de una cadena.

El bloque de transform puede además emitir varias veces por cada entrada, y ahí es donde el operador deja de ser una curiosidad académica. Aplanar, intercalar valores de control o expandir un elemento en varios son casos que ningún operador con nombre cubre y que aquí salen solos.

val conMarcas = numeros.transform { n ->
    emit(Evento.Antes(n))
    emit(Evento.Valor(n * 2))
    if (n % 10 == 0) emit(Evento.Hito(n))
}

Conviene mencionar dos parientes cercanos que aparecen a diario y que se suelen confundir. onEach no emite un valor nuevo, mientras que map sí; usar map cuando lo único que se quería era registrar produce cadenas donde no se sabe qué tipo lleva cada tramo. Y filterNotNull con filterIsInstance no son solo filtros: además estrechan el tipo del flujo resultante, lo cual es una razón excelente para preferirlos a un filter escrito a mano que dejaría el tipo intacto y obligaría a una conversión más abajo.

🔁

transform

El primitivo. Cero, uno o muchos valores de salida por cada valor de entrada, con el colector como receptor.

🧭

map y filter

Casos particulares con nombre. Uno transformado y cero o uno respectivamente. Nada más ocurre dentro.

🪵

onEach

Efecto secundario sin alterar el flujo. El instrumento correcto para registrar sin cambiar tipos.

✂️

take y drop

Recortan por posición. El primero es especial porque necesita apagar al productor cuando ya tiene bastante.

La cadena es una matrioska

Escribamos qué ocurre cuando alguien recolecta una cadena de tres tramos. La llamada a collect va al último flujo construido, que es el del filtro. Su cuerpo recolecta el flujo del mapeo pasándole un colector propio. El cuerpo de aquel recolecta el flujo original pasándole otro colector propio. Y solo entonces el productor empieza a emitir, hacia dentro de esa pila de colectores anidados.

flowchart TD
A[collect sobre el flujo del filtro] --> B[El filtro recolecta el flujo del mapeo]
B --> C[El mapeo recolecta el flujo original]
C --> D[El productor emite un valor]
D --> E[Colector del mapeo transforma y emite]
E --> F[Colector del filtro decide y emite]
F --> G[Cuerpo del usuario recibe el valor]
G --> D

De ese dibujo salen tres conclusiones que conviene tener presentes al medir. La primera es que no existe ningún lugar donde un valor pueda esperar: la única memoria del sistema es la pila de llamadas, y la pila solo puede contener un valor en tránsito. La segunda es que el número de objetos creados es proporcional al número de operadores, no al número de elementos, porque los colectores se construyen una vez al empezar la recolección y se reutilizan para todos los valores. La tercera es que si una etapa suspende, la cadena entera queda detenida, incluido el productor, que es la contrapresión de la lección siguiente vista desde otro ángulo.

Hay un detalle de implementación que conviene conocer porque explica por qué esto no es tan caro como parece. La mayoría de los operadores intermedios están declarados como funciones en línea con sus lambdas marcadas para no escapar, de modo que el cuerpo que escribes no genera un objeto función por llamada, sino que se incrusta en el sitio. Combinado con lo anterior, una cadena de cinco operadores sobre un millón de elementos crea del orden de cinco o diez objetos en total, no cinco millones.

take y la excepción que apaga al productor

Queda el caso que no encaja en el molde. Un operador como take necesita, después de recibir los elementos que le pedían, detener al productor que sigue arriba y que no tiene manera de enterarse. Como la única comunicación va hacia abajo, la única forma de mandar información hacia arriba es lanzar.

// Version conceptual del mecanismo
fun <T> Flow<T>.take(n: Int): Flow<T> = flow {
    var restantes = n
    try {
        collect { valor ->
            emit(valor)
            if (--restantes == 0) throw AbortFlowException(this)
        }
    } catch (e: AbortFlowException) {
        // Solo se traga la suya propia: la comprueba por identidad
    }
}

El mecanismo es una excepción interna de la librería que se lanza desde dentro del colector, sube por toda la cadena atravesando al productor, y la captura el propio operador que la lanzó, que comprueba por identidad que es la suya y no la de otro take anidado. El productor no llega a saber por qué se le interrumpió, pero sus bloques finally se ejecutan con normalidad y sus recursos se cierran.

Esto tiene una consecuencia práctica desagradable y muy frecuente. Si el cuerpo productor captura excepciones de forma indiscriminada, se tragará también esa señal y romperá el operador que la envió.

// Roto: este flujo no se puede cortar con take
val malo = flow {
    try {
        while (true) emit(leerSensor())
    } catch (e: Exception) {
        registrar(e)          // se traga la senal de aborto y de cancelacion
    }
}

La regla que se deduce es la misma que ya regía para la cancelación de corrutinas y ahora vale también para los flujos: nunca captures el tipo raíz de las excepciones dentro de un cuerpo productor. Si necesitas limpiar, usa finally, que se ejecuta igual sin interceptar nada. Si necesitas capturar de verdad, captura el tipo concreto que esperas y deja pasar el resto.

Un operador que no puede ejecutar nada es mucho mas poderoso que uno que si puede porque la impotencia es lo que garantiza que la composicion sea gratis

Vale la pena invertir el juicio natural sobre esta parte del diseño. Que un operador intermedio no pueda suspender parece una limitación y suena a que el diseñador se quedó corto, cuando en realidad es la restricción que hace que todo lo demás encaje. Considérese qué se perdería si map estuviese marcado como suspendido y pudiera, por tanto, empezar a consumir de su flujo de origen en el momento de construirse. Se perdería la propiedad de que escribir una cadena no cuesta nada, y con ella la posibilidad de construir cadenas en sitios donde no hay corrutina, de guardarlas en propiedades, de devolverlas desde funciones ordinarias y de pasarlas por parámetro sin haber decidido todavía quién ni cuándo las va a consumir. Se perdería la reutilización, porque un flujo que ya empezó no se puede volver a recolectar desde el principio, y con ella desaparecerían el reintento y la sustitución del flujo entero por otro tras un fallo. Se perdería la fusión de operadores, porque un operador que ya está ejecutando no puede negociar con su vecino para colapsarse en una sola etapa, y esa negociación es lo que permite que varias llamadas consecutivas de contexto o de búfer se resuelvan en un solo canal. Y se perdería, sobre todo, la propiedad de que el consumidor manda: hoy quien decide el ritmo, el contexto de ejecución, el momento de empezar y el momento de parar es el que recolecta, y el que escribió la cadena no le impuso ninguna de esas cuatro cosas. Esa asimetría entre quien describe y quien ejecuta es la misma que separa una consulta declarativa de un bucle imperativo, y produce el mismo beneficio: mientras la descripción siga siendo solo una descripción, el sistema puede reordenarla, fusionarla, repetirla o descartarla entera. En cuanto una sola etapa empieza a actuar por su cuenta, todas esas libertades se evaporan a la vez, porque ya hay efectos observados que no se pueden deshacer. La lección general es que el poder de una abstracción compositiva no viene de lo que sus piezas saben hacer, sino de lo que se les ha prohibido hacer antes de tiempo.

⚔️ Desmonta la cadena
  1. Escribe una cadena con map, filter y onEach, con un println en cada uno, y predice el orden exacto de las impresiones para tres elementos antes de ejecutarla.
  2. Reescribe map, filter y onEach con transform y comprueba que tu cadena da el mismo resultado que la original.
  3. Implementa con transform un operador que intercale un separador entre elementos consecutivos y no lo emita ni antes del primero ni después del último.
  4. Aplica take a un flujo infinito y coloca un finally en el cuerpo productor. Comprueba que se ejecuta y explica por qué medio se enteró el productor.
  5. Rompe ese mismo caso capturando el tipo raíz de las excepciones dentro del productor. Observa el fallo y arréglalo sin dejar de registrar los errores reales.