wandres.dev
CORRUTINAS IV · canales y select

Productores y consumidores: produce, iterar y cerrar sin perder nada

Crear un canal es la parte fácil; lo difícil es que alguien lo cierre exactamente una vez, en el momento correcto, y que el fallo del productor llegue al consumidor como una excepción y no como un final de iteración indistinguible del éxito. Esta lección estudia el constructor produce y el ámbito que expone, la diferencia real entre iterar con un bucle y consumir con consumeEach, la asimetría deliberada entre cerrar y cancelar, y los dos caminos simultáneos por los que un fallo del productor recorre el sistema.

⏱ 22 min

Un canal sin dueño es una fuga esperando a ocurrir. La pregunta que decide si tu código de canales sobrevive al contacto con producción no es cómo se envía ni cómo se recibe, sino quién es responsable de cerrar, qué ocurre con los elementos que ya estaban dentro cuando eso pasa y cómo se entera el consumidor de que el productor no terminó porque hubiera acabado su trabajo, sino porque explotó. Kotlin responde a las tres con una sola idea: hacer que el canal sea propiedad de una corrutina y que su ciclo de vida coincida con el de esa corrutina. El constructor produce es esa idea hecha función, y entenderlo bien es entender por qué cerrar a mano casi siempre es un síntoma.

🎯 Al terminar esta lección sabrás
  • Construir productores con produce y explicar qué garantiza su ámbito respecto al cierre del canal.
  • Distinguir con precisión iterar un canal con un bucle de consumirlo con consumeEach o consume.
  • Separar la operación de cerrar, que pertenece al emisor, de la de cancelar, que pertenece al receptor.
  • Trazar los dos caminos por los que un fallo del productor alcanza al consumidor y al ámbito padre.

El productor como corrutina: produce

El patrón de crear un canal, lanzar una corrutina que lo llene y devolverlo es tan común que la librería lo ofrece como un constructor. produce es una extensión de CoroutineScope que arranca una corrutina, le pasa un ProducerScope —que es a la vez un CoroutineScope y un SendChannel— y devuelve el ReceiveChannel correspondiente.

fun CoroutineScope.numeros(hasta: Int): ReceiveChannel<Int> = produce {
    for (n in 1..hasta) send(n)
}   // al terminar el bloque, el canal se cierra solo

Lo relevante no es la brevedad sino las dos invariantes que regala. La primera: el canal se cierra automáticamente cuando el bloque termina, sea porque acabó, sea porque lanzó una excepción, sea porque lo cancelaron. No hay ninguna ruta de salida en la que el canal quede abierto y sin dueño. La segunda: la corrutina es hija del ámbito receptor, de modo que cancelar el ámbito cancela al productor y el productor que falla afecta al ámbito, que es exactamente lo que la concurrencia estructurada promete y lo que un Channel creado a mano no te da.

El constructor acepta además los mismos parámetros que la factoría de canales y un contexto, lo que permite fijar en un solo sitio el buffer, el despachador y el nombre de la corrutina que produce. Ese es el punto donde declaras la política de contrapresión de esa fuente concreta, y tenerlo junto a la lógica que emite es preferible a esconderlo en la construcción de un canal tres capas más arriba.

fun CoroutineScope.lineas(ruta: Path): ReceiveChannel<String> =
    produce(Dispatchers.IO, capacity = 128) {
        Files.newBufferedReader(ruta).use { lector ->
            lector.lineSequence().forEach { send(it) }
        }
    }   // el use cierra el fichero también si nos cancelan

Ese use dentro del bloque ilustra la propiedad más valiosa del constructor: como la cancelación llega al productor en forma de excepción en su próximo send, todos los bloques de limpieza que hayas escrito se ejecutan. El recurso vive dentro de la corrutina que lo usa y muere con ella, sin registros de cierre ni banderas.

Que el tipo de retorno sea ReceiveChannel y no Channel tampoco es casualidad: el llamante no puede enviar ni cerrar, solo consumir. El único que envía es el bloque, y el único que cierra es el ciclo de vida de la corrutina.

⚠️
Marcado como experimental desde siempre

produce lleva años anotado como API experimental de corrutinas, no porque vaya a desaparecer sino porque su interacción con la concurrencia estructurada y el manejo de excepciones ha ido puliéndose. Necesitarás el OptIn correspondiente, y conviene saber que el compilador te lo está recordando por una razón real y no por burocracia.

Iterar, consumir y terminar

Un ReceiveChannel se puede recorrer con un bucle for porque implementa un operador de iteración suspendida. El bucle termina limpiamente cuando el canal se cierra sin causa, y relanza la excepción si se cerró con una.

for (n in canal) {          // termina al cerrarse el canal
    procesar(n)
}

Ese bucle tiene un agujero que la firma no revela: si el cuerpo lanza una excepción o si sales con break, el bucle se abandona y el canal queda abierto. El productor seguirá vivo, suspendido para siempre en un send que nadie va a atender, salvo que el ámbito que lo contiene se cancele por otro motivo. Para el caso habitual —quiero este canal para mí y cuando yo termine ya no le sirve a nadie— existe una familia de extensiones que cancelan el canal al salir, pase lo que pase.

canal.consumeEach { procesar(it) }        // cancela el canal al salir

val primero = canal.consume { receive() } // consume uno y cancela el resto

La diferencia entre las dos familias se entiende mejor si se lee lo que cada una afirma sobre la propiedad del canal. El bucle for dice estoy leyendo de algo que no es mío, y por eso lo deja como estaba al salir. consumeEach dice este canal es mío y su vida termina cuando yo termine, y por eso lo cancela. Ninguna de las dos es más segura que la otra en abstracto: cada una es correcta bajo una afirmación de propiedad distinta, y el error consiste en escribirlas sin haber decidido cuál de las dos afirmaciones es cierta en tu caso.

De ahí sale una regla que parece un detalle y no lo es: consumeEach es correcto cuando hay un único consumidor y venenoso cuando hay varios, porque el primero que termine cancelará el canal debajo de los pies de sus compañeros. En cuanto varias corrutinas comparten un canal, el bucle for es la forma correcta y el cierre pasa a ser responsabilidad del ámbito común.

Queda una tentación que hay que desactivar de raíz: preguntar al canal si está cerrado antes de actuar. Las propiedades que informan de ello existen, pero consultarlas es una carrera contra el otro extremo, porque entre la pregunta y la acción cualquier cosa puede haber pasado. La forma correcta de saber si un canal terminó es intentar la operación y mirar su resultado, no consultar un estado que ya está caduco cuando lo lees.

🚪

close pertenece al emisor

Marca el final de la secuencia. Los elementos ya presentes en el buffer se entregan igualmente antes de que el receptor vea el final. Es idempotente y devuelve si fue el primero en cerrar.

✂️

cancel pertenece al receptor

Declara que ya no se quiere nada más. Descarta el buffer de inmediato y hace fallar los envíos pendientes. No es un cierre educado, es una interrupción.

🧾

close con causa

Adjunta una excepción al final de la secuencia. Todo receptor que llegue a ese punto la relanzará, lo que convierte el canal en un transporte de errores además de datos.

🧹

onUndeliveredElement

La única red bajo los elementos que ya entraron y nunca saldrán, porque un cancel llegó antes que su receptor. Sin él, cancelar es perder en silencio.

Cuando el productor falla

Aquí está la parte que más código real rompe. Si el bloque de produce lanza una excepción, ocurren dos cosas a la vez y hay que tener las dos en la cabeza.

flowchart TD
A[Excepcion dentro de produce] --> B[El canal se cierra con esa causa]
A --> C[La corrutina hija falla]
B --> D[El consumidor la relanza en receive]
C --> E[El ambito padre se cancela]
E --> F[Los demas hijos se cancelan tambien]
D --> G[Se maneja como una excepcion normal]

La primera: el canal se cierra con esa excepción como causa, de modo que el consumidor la recibirá en su próximo receive o al final de su bucle. La segunda: la corrutina del productor es hija del ámbito, así que su fallo cancela al padre y con él al resto de hermanos, siguiendo las reglas normales de la concurrencia estructurada. El error de diseño clásico consiste en creer que el canal sustituye a la propagación estructural y envolver el bucle consumidor en un try esperando que eso contenga el fallo. No lo contiene: lo captura en un sitio mientras el ámbito ya se está desmontando por el otro.

suspend fun leerLineas(ruta: Path) = coroutineScope {
    val lineas = produce {
        Files.newBufferedReader(ruta).use { lector ->
            lector.lineSequence().forEach { send(it) }
        }
    }
    lineas.consumeEach(::procesar)   // la IOException llega aquí
}                                    // y además cancela este ámbito

Para el caso opuesto —quiero que el fallo viaje solo por el canal, como un dato— la respuesta no es luchar contra la jerarquía sino cambiar lo que se envía: un canal de Result o de un tipo sellado propio convierte el error en un elemento más de la secuencia, y entonces el productor nunca falla y el consumidor decide qué hacer con cada caso.

sealed interface Registro
data class Valida(val linea: String) : Registro
data class Corrupta(val numero: Int, val causa: String) : Registro

fun CoroutineScope.registros(ruta: Path) = produce<Registro> {
    lecturaDe(ruta).forEachIndexed { i, texto ->
        send(runCatching { parsear(texto) }
            .fold({ Valida(it) }, { Corrupta(i, it.message ?: "desconocida") }))
    }
}

Es más código y bastante más honesto, porque hace visible en el tipo que esta secuencia puede contener fallos parciales sin que eso termine el trabajo. La distinción que estás codificando es la que separa un fallo del elemento de un fallo de la fuente: el primero es un dato y debe viajar como dato; el segundo es el final de la secuencia y debe viajar como cierre con causa. Mezclarlos —capturar la excepción de disco y enviarla como elemento, o dejar que un carácter mal codificado tumbe la lectura entera— es el origen de la mayoría de los productores que se comportan de forma incomprensible bajo datos reales.

💡
Si escribes close a mano, pregúntate quién más podría cerrarlo

Un close explícito solo es correcto cuando hay exactamente un emisor y ese emisor conoce el final de la secuencia. En cuanto hay dos, cerrar es una carrera: el primero que termina apaga la luz para todos. La solución no es coordinar los cierres con banderas, es dejar que el ámbito que los contiene a todos sea quien decida el final.

El cierre de un canal no es una operación de limpieza sino el último elemento de la secuencia, y tratarlo como tal resuelve de golpe la mitad de los problemas de ciclo de vida

La intuición que arruina el código de canales es la que hereda de los ficheros y los sockets: abrir, usar, cerrar en un bloque final. Bajo esa intuición, cerrar es higiene, algo que se hace después de que lo importante haya ocurrido, y por eso se escribe en el sitio donde uno recuerda escribirlo y se olvida en las tres rutas de salida que uno no había previsto. La intuición correcta es la contraria: en un canal, el cierre es información y viaja por el mismo carril que los datos. Cuando un productor cierra, no está liberando un recurso, está enviando el último mensaje de la conversación, un mensaje que dice no habrá más y que, si lleva causa adjunta, dice además no habrá más porque esto se rompió. Esa distinción tiene consecuencias inmediatas y muy concretas. Explica por qué el cierre no descarta el buffer: los elementos que ya estaban dentro fueron enviados antes que el mensaje de final, y entregarlos primero es simplemente respetar el orden de la secuencia. Explica por qué cancelar sí descarta: cancelar no es enviar un último mensaje sino declarar que la conversación entera dejó de importar, y ahí ya no hay orden que respetar. Explica por qué el cierre debe ser idempotente y por qué solo el emisor puede hacerlo: un mensaje de final que pudiera enviarse dos veces, o que pudiera enviar quien no está hablando, sería un mensaje sin significado. Y explica, sobre todo, por qué produce es superior a construir el canal a mano: al ligar el final de la secuencia al final de la corrutina, hace que ese último mensaje se envíe siempre y exactamente una vez, en cada ruta de salida imaginable, incluida la que tú no imaginaste. La consecuencia práctica es una regla que puedes aplicar sin pensar: si en tu código aparece un close fuera de un bloque cuya terminación es el final de la secuencia, no estás cerrando un canal, estás intentando sincronizar dos ciclos de vida a mano. Y sincronizar ciclos de vida a mano es precisamente el trabajo que la concurrencia estructurada existe para quitarte.

⚔️ Rompe y repara un productor
  1. Escribe un produce que lea un fichero línea a línea y haz que falle a mitad. Comprueba dónde aparece la excepción en el consumidor y qué le ocurre al ámbito que los contiene.
  2. Sustituye el consumeEach por un bucle for con un break en el elemento diez. Verifica con un finally en el productor si llegó a ejecutarse y explica el resultado.
  3. Repite el punto anterior con consume y compara. Escribe en una frase la regla que has descubierto.
  4. Lanza dos consumidores sobre el mismo canal usando consumeEach en ambos y observa el desastre. Arréglalo cambiando a bucles for y cerrando desde el ámbito común.
  5. Reescribe el productor para que emita un tipo sellado con casos de éxito y de fallo en lugar de fallar. Compara qué garantiza cada versión y decide cuál querrías en un servicio que no puede caerse por una línea mal formada.