Fan-out, fan-in y pipelines: topologías sobre canales
Con una sola primitiva y tres formas de conectarla aparecen las tres topologías que sostienen casi todo el procesamiento concurrente: repartir un canal entre varios trabajadores, fundir varios productores en un único destino y encadenar etapas donde la salida de una es la entrada de la siguiente. Esta lección estudia el reparto equitativo y por qué no es difusión, la regla de cierre que solo un ámbito puede aplicar cuando hay varios emisores, la propagación de la cancelación aguas arriba en una tubería y el criterio para saber cuándo un canal es la herramienta correcta y cuándo deberías estar usando Flow.
Una vez que existe una primitiva que conecta a un emisor con un receptor y que suspende a quien vaya demasiado rápido, todo lo demás es topología. Poner varios receptores sobre un mismo canal produce un repartidor de trabajo con equilibrado automático que nadie tuvo que programar. Poner varios emisores sobre un mismo canal produce un embudo que serializa flujos concurrentes sin un solo cerrojo. Encadenar canales produce una tubería donde cada etapa avanza al ritmo que le permite la siguiente. Las tres emergen de la contrapresión, no de código de coordinación, y las tres comparten el mismo punto delicado: quién declara el final y cómo viaja la cancelación por la topología cuando alguien se rinde.
- Construir un reparto de trabajo entre varios consumidores y explicar por qué reparte en lugar de difundir.
- Fundir varios productores en un canal y aplicar la única regla de cierre que funciona con varios emisores.
- Encadenar etapas en una tubería y garantizar que la cancelación de la última alcance a todas las anteriores.
- Decidir con criterio cuándo la topología pide canales y cuándo pide
Flow.
Fan-out: un canal, muchos consumidores
Nada impide que varias corrutinas reciban del mismo canal. Cada elemento lo entrega el canal a exactamente uno de los receptores en espera, y lo hace en orden de llegada, de modo que el reparto es justo. El resultado es un repartidor de tareas donde el trabajador que termina antes pide antes, y por tanto el equilibrado de carga es una consecuencia gratuita de la contrapresión.
fun CoroutineScope.trabajador(id: Int, tareas: ReceiveChannel<Tarea>) = launch {
for (t in tareas) { // for, nunca consumeEach
procesar(t)
}
}
coroutineScope {
val tareas = produce { repeat(1_000) { send(Tarea(it)) } }
repeat(8) { trabajador(it, tareas) }
}
El detalle que separa el código correcto del que falla en producción está en el comentario: aquí hay que usar el bucle for y no consumeEach, porque este último cancela el canal al salir y el primer trabajador que terminase dejaría a los otros siete sin trabajo y con una excepción de cancelación. Con un bucle, cada trabajador simplemente sale cuando el canal se cierra, que es cuando ya no queda nada para nadie.
El reparto se completa casi siempre con su simétrico: los trabajadores no solo consumen, también producen resultados hacia un canal común, y entonces la topología entera cabe en un ámbito.
suspend fun procesarTodo(entradas: List<Tarea>): List<Resultado> = coroutineScope {
val tareas = produce { entradas.forEach { send(it) } }
val salidas = Channel<Resultado>(capacity = entradas.size)
val obreros = List(8) { launch { for (t in tareas) salidas.send(calcular(t)) } }
obreros.joinAll()
salidas.close()
salidas.toList()
}
Ese pequeño programa contiene ya las tres decisiones difíciles del nivel: quién cierra, con qué capacidad y quién espera a quién. Cámbiale cualquiera de las tres y deja de funcionar por motivos distintos.
La segunda propiedad que hay que interiorizar es que esto reparte, no difunde. Con ocho trabajadores y mil tareas, cada tarea se procesa una vez. Si lo que querías era que los ocho vieran las mil, ningún ajuste de canal te lo dará: necesitas un SharedFlow, que es la primitiva pensada para varios observadores del mismo evento.
Con trabajo ligado a la CPU, tantos trabajadores como paralelismo real tenga el despachador. Con trabajo ligado a entrada y salida, muchos más, porque están suspendidos casi todo el tiempo. La única forma de acertar es medir el rendimiento total variando el número, no razonarlo desde el hardware.
Fan-in: muchos productores, un canal
La dirección contraria es igual de simple de escribir y bastante más fácil de estropear. Varias corrutinas pueden enviar al mismo canal y el canal serializa los elementos sin que nadie tome un cerrojo; lo que no puede hacerse es que cada productor cierre cuando acaba, porque el primero en terminar cerraría el canal para todos y los demás verían fallar sus envíos.
suspend fun fundir(fuentes: List<ReceiveChannel<Evento>>): ReceiveChannel<Evento> =
coroutineScope {
val salida = Channel<Evento>()
launch {
coroutineScope { // espera a todos los hijos
fuentes.forEach { fuente ->
launch { for (e in fuente) salida.send(e) }
}
}
salida.close() // exactamente un cierre
}
salida
}
El patrón es siempre el mismo: un ámbito interno que contiene a todos los emisores, y el cierre justo después de que ese ámbito termine. El coroutineScope anidado no está de adorno, es lo que convierte todos han terminado en una línea de código en lugar de en un contador atómico y una condición de carrera. Cualquier variante con banderas compartidas o con un contador de productores vivos es una reimplementación peor de esa misma espera.
El orden de salida del fan-in no está especificado más allá de que cada fuente conserva su orden interno. Si necesitas una fusión ordenada por marca de tiempo, ningún canal te la va a dar: tendrás que mirar el primer elemento de cada fuente y elegir, que es un caso natural para lo que estudiarás en la lección siguiente.
Pipelines y la cancelación aguas arriba
Encadenar etapas es la topología más vistosa y la que más fugas produce. Cada etapa es un produce que consume el canal de la anterior y emite el suyo.
fun CoroutineScope.enteros() = produce { var n = 1; while (true) send(n++) }
fun CoroutineScope.filtrar(entrada: ReceiveChannel<Int>, p: Int) = produce {
for (n in entrada) if (n % p != 0) send(n)
}
coroutineScope {
var canal = enteros()
repeat(20) {
val primo = canal.receive()
println(primo)
canal = filtrar(canal, primo)
}
coroutineContext.cancelChildren() // sin esto, veinte corrutinas eternas
}
La capacidad de los canales intermedios es, en una tubería, un parámetro con significado propio: desacopla el ritmo de dos etapas contiguas y permite que una siga trabajando mientras la otra atiende un pico. Es exactamente el mismo papel que juega el operador de amortiguación en un flujo, y el mismo peligro: cada hueco entre dos etapas es trabajo ya hecho que todavía no ha servido para nada, y una tubería de seis etapas con buffers generosos puede tener más elementos en tránsito que en proceso.
Esa última línea es la lección entera. Las etapas de una tubería infinita no terminan solas: cuando el consumidor final decide que ya tiene bastante, alguien debe cancelar hacia arriba. Si todas las etapas son hijas del mismo ámbito, cancelar ese ámbito basta y el orden de apagado lo resuelve la propia jerarquía, porque cada etapa cancelada cierra su canal y la siguiente termina su bucle. Si en cambio has creado los canales a mano o has lanzado alguna etapa en un ámbito ajeno, acabas de construir un sistema donde la mitad de las corrutinas no tienen padre y no hay forma de pararlas.
flowchart LR P[Productor] --> C1[Canal 1] C1 --> T1[Etapa filtro] T1 --> C2[Canal 2] C2 --> W1[Trabajador 1] C2 --> W2[Trabajador 2] C2 --> W3[Trabajador 3] W1 --> C3[Canal de resultados] W2 --> C3 W3 --> C3 C3 --> S[Consumidor final]
Reparto justo
Los receptores en espera se atienden en orden de llegada. El equilibrado de carga aparece sin código porque el trabajador libre es el que está esperando.
Serialización sin cerrojos
Un fan-in convierte N emisores concurrentes en una secuencia ordenada. Es exclusión mutua obtenida por topología y no por bloqueo.
Cierre único
Con varios emisores, el cierre pertenece al ámbito que los contiene a todos. Cerrar desde un emisor es una carrera contra sus compañeros.
Cancelación descendente
La jerarquía de trabajos es el único mecanismo de apagado fiable de una tubería. Cancelar el ámbito propaga el final por todas las etapas.
La tubería de números primos es un ejemplo pedagógico soberbio y un antipatrón de producción, y conviene entender por qué las dos cosas son ciertas a la vez. Es soberbia porque enseña en quince líneas que la contrapresión compone: cada etapa avanza exactamente al ritmo que le permite la siguiente y nadie tuvo que programar esa coordinación. Es un antipatrón porque cada etapa es un objeto vivo con su propia corrutina, su propio buffer y su propia necesidad de ser cancelada, y basta con que una sola quede huérfana para que el proceso arrastre memoria y trabajo para siempre. Flow existe precisamente para eliminar esa categoría de error en el caso lineal: un flujo frío no es un objeto vivo sino una descripción de cómo producir elementos, la ejecución entera ocurre dentro de la corrutina que colecta, no hay canales intermedios que cerrar ni etapas que cancelar, y cuando la colección termina —normal o excepcionalmente— todo lo que existía deja de existir sin que nadie tenga que acordarse. Por eso la regla práctica es tan asimétrica: si tu topología es una cadena, usa Flow y no lo pienses más; los operadores intermedios cubren el noventa por ciento de lo que ibas a escribir a mano y el diez restante lo cubre transform. Los canales ganan su sitio en tres situaciones muy concretas, y solo en esas tres. Primera, cuando de verdad ramifica: repartir un flujo entre trabajadores que compiten por los elementos no es una operación de flujo, es una topología, y ahí el canal es la primitiva correcta —de hecho es lo que flatMapMerge y buffer usan por dentro—. Segunda, cuando la fuente es caliente y ajena a tu ciclo de vida: un callbackFlow es un canal con un envoltorio, porque un callback que llega cuando le apetece necesita un sitio donde dejar el elemento aunque nadie esté mirando. Y tercera, cuando el elemento es un comando y no un dato, es decir, cuando lo que atraviesa la tubería son peticiones a un estado con protocolo y orden, que es justo el caso de los actores. Fuera de esas tres, un canal en tu código suele ser un Flow que alguien escribió antes de aprender que existía, con la diferencia de que este sí puede quedarse encendido cuando te vayas a casa.
- Monta un fan-out de ocho trabajadores sobre mil tareas de duración aleatoria y registra cuántas atendió cada uno. Comprueba que el reparto sigue la velocidad de cada trabajador y no un turno fijo.
- Cambia el bucle
forde los trabajadores porconsumeEachy observa el fallo exacto que provoca. Escribe el mensaje de error que aparece y por qué. - Implementa el fan-in de tres fuentes con el patrón del ámbito anidado y luego rómpelo cerrando desde dentro de cada emisor. Identifica la excepción que produce la carrera.
- Ejecuta la tubería de primos sin la cancelación final dentro de un proceso que siga vivo y observa el número de hilos y la memoria retenida. Añade la cancelación y repite la medida.
- Reescribe esa misma tubería con
Flowy operadores intermedios. Compara el número de líneas, el número de objetos vivos y el número de formas distintas de dejar algo encendido.