SharedFlow: repetición, capacidad extra y las dos formas de emitir
Un SharedFlow es un flujo caliente de multidifusión cuyo comportamiento entero queda determinado por tres parámetros del constructor que casi nadie lee con atención: cuántos valores guarda para quien llegue tarde, cuánto colchón adicional concede al emisor por delante del suscriptor más lento, y qué sacrifica cuando ese colchón se agota. Esta lección reconstruye el buffer real que resulta de sumar los tres, explica por qué emit suspende exactamente cuando suspende, y establece la regla dura que gobierna tryEmit: solo puede tener éxito si has comprado de antemano el derecho a no suspender.
SharedFlow<T> suele presentarse como el flujo para eventos, y esa etiqueta esconde lo único que hay que entender de él: es una estructura con un buffer compartido por todos los suscriptores, en el que cada suscriptor mantiene su propio índice de lectura, y en el que el emisor solo puede avanzar mientras quepa por delante del lector más rezagado. De esa mecánica se derivan, sin excepción, todos sus comportamientos: por qué un suscriptor tardío ve historia, por qué un suscriptor lento frena a todos los demás, por qué tryEmit devuelve false en unos montajes y nunca en otros, y por qué el valor por defecto —cero repetición, cero colchón, suspender— es el más estricto de todos los posibles y también el más honesto.
- Reconstruir el tamaño y la semántica del buffer real a partir de
replayyextraBufferCapacity. - Explicar la condición exacta bajo la cual
emitsuspende y cómo influye el suscriptor más lento. - Aplicar las tres políticas de
BufferOverflowsabiendo qué se pierde con cada una y en qué extremo del buffer. - Decidir entre
emitytryEmita partir de la garantía que puedes ofrecer, no de la comodidad del punto de llamada.
El buffer real es la suma de dos números
El constructor tiene tres parámetros y ninguno es decorativo. replay es cuántos de los valores ya emitidos se le entregan de golpe a un suscriptor nuevo en el momento de suscribirse. extraBufferCapacity es cuántos valores adicionales puede aceptar el flujo por delante del suscriptor más lento antes de tener que suspender al emisor. Y onBufferOverflow decide qué ocurre cuando ese espacio se agota.
val eventos = MutableSharedFlow<Evento>(
replay = 0,
extraBufferCapacity = 0,
onBufferOverflow = BufferOverflow.SUSPEND, // valores por defecto
)
La primera consecuencia, que es la que sorprende, es que la capacidad total del buffer es la suma de los dos primeros. Un replay de tres no solo da historia: da también tres huecos de holgura al emisor, porque esos valores están físicamente en el mismo buffer. Por eso un MutableSharedFlow con repetición tres y sin capacidad extra tiene un colchón efectivo de tres, y por eso mucha gente cree haber configurado solo historia cuando en realidad ha configurado también tolerancia a ráfagas.
La segunda consecuencia es que con los valores por defecto —suma cero— el flujo se comporta como una cita: emit suspende hasta que todos los suscriptores actuales han recibido el valor. Es un comportamiento de contrapresión total, correcto para eventos que no se pueden perder, y también la explicación de por qué un solo suscriptor lento arrastra al sistema entero: no hay múltiples colas independientes, hay un buffer con un puntero por suscriptor, y el emisor mira siempre al más rezagado.
Si no hay ningún suscriptor y replay es cero, emit no suspende: el valor se descarta inmediatamente porque no queda nadie a quien deba nada. La contrapresión de un SharedFlow es siempre relativa a los suscriptores vivos, nunca absoluta.
Cuándo suspende emit y qué se pierde si no
La regla es corta: emit suspende cuando la política es SUSPEND y el nuevo valor no cabe sin desalojar un valor que el suscriptor más lento todavía no ha leído. Con las otras dos políticas no suspende jamás, y lo que cambia es cuál de los dos extremos del buffer se sacrifica.
SUSPEND
El emisor espera al más lento. Cero pérdida y contrapresión real que sube hasta la fuente. Es el único modo en que un evento emitido está garantizado como entregado.
DROP_OLDEST
Se desaloja el valor más antiguo aún no leído. El emisor nunca frena y el suscriptor lento se salta huecos sin enterarse. Correcto para estado, letal para secuencias con significado acumulativo.
DROP_LATEST
Se descarta el valor recién emitido. Conserva la historia ya aceptada y sacrifica lo nuevo. Es lo que quieres en muestreo y telemetría, donde una serie veraz vale más que una completa.
Hay un detalle que rompe expectativas y conviene fijar: las políticas de descarte solo pueden usarse si hay algún buffer, es decir, si la suma de replay y extraBufferCapacity es mayor que cero. Pedir DROP_OLDEST con buffer cero es una contradicción —no hay nada viejo que tirar— y la biblioteca lo rechaza.
// Estado de conexión: al último le importa el presente, no la historia.
val conexion = MutableSharedFlow<Estado>(
replay = 1,
onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
// Comandos del usuario: ninguno puede perderse, el emisor espera.
val comandos = MutableSharedFlow<Comando>(extraBufferCapacity = 8)
Un segundo detalle, que se descubre siempre depurando: replay no es un caché consultable. No hay una propiedad valor actual en un SharedFlow genérico, aunque expone replayCache como lista de solo lectura para inspección. Si lo que necesitas es el último valor, siempre disponible, sincrónicamente, no estás describiendo un SharedFlow con repetición uno: estás describiendo un StateFlow, y la diferencia entre ambos no es de forma sino de conflación, como veremos en la lección siguiente.
flowchart TD
E[emit de un valor nuevo] --> Q{Cabe por delante del suscriptor mas lento}
Q -->|Si| OK[Se guarda y se entrega]
Q -->|No| P{Politica de desbordamiento}
P -->|SUSPEND| W[El emisor espera al mas lento]
P -->|DROP_OLDEST| D1[Se desaloja el mas antiguo no leido]
P -->|DROP_LATEST| D2[Se descarta el valor recien emitido]emit frente a tryEmit, y quién los observa
tryEmit es la variante no suspendida y devuelve un Boolean que indica si el valor entró. Su regla es tan estricta que resuelve casi todas las dudas: tryEmit devuelve true si y solo si el valor pudo aceptarse sin suspender. Con la política SUSPEND y buffer cero, eso significa que devolverá false siempre que haya un suscriptor que no esté esperando activamente. Con cualquiera de las políticas de descarte, en cambio, devolverá true siempre, porque descartar es una forma legítima de aceptar.
// Un callback del sistema no puede suspender: hay que comprar el derecho.
private val _pulsaciones = MutableSharedFlow<Pulsacion>(
extraBufferCapacity = 64,
onBufferOverflow = BufferOverflow.DROP_OLDEST,
)
fun alPulsar(p: Pulsacion) {
_pulsaciones.tryEmit(p) // aquí siempre entra: se pagó con buffer
}
De ahí sale la única heurística sólida sobre el asunto: si escribes tryEmit sobre un flujo con los parámetros por defecto y no compruebas el resultado, has escrito una pérdida silenciosa con apariencia de emisión. Y si lo escribes ignorando el booleano porque siempre da true, entonces el diseño real está en el constructor y no en el punto de llamada, así que el comentario que explica la política pertenece al constructor.
La otra pieza que completa el cuadro es subscriptionCount, un StateFlow<Int> que publica cuántos suscriptores hay en cada momento. Sirve para dos cosas legítimas y una discutible. Legítimas: arrancar la fuente cara solo cuando aparece el primer suscriptor y pararla cuando se va el último —que es justamente lo que automatiza WhileSubscribed—, y esperar activamente a que haya alguien antes de emitir, cerrando la ventana de carrera del arranque.
// Cerrar la ventana: no emitas al vacío durante el arranque.
suspend fun emitirCuandoHayaPublico(e: Evento) {
eventos.subscriptionCount.first { it > 0 }
eventos.emit(e)
}
Lo discutible es usar subscriptionCount para tomar decisiones de negocio, porque es una cantidad inherentemente de carrera: entre que la lees y actúas, puede haber cambiado. Como señal de ciclo de vida es excelente; como condición de corrección, no.
Queda una propiedad estructural que conviene grabar: MutableSharedFlow no completa nunca y no propaga excepciones del productor. No existe cerrar un SharedFlow. Si tu dominio necesita señalar el fin o el fallo, tienes que modelarlo dentro del tipo emitido —un tipo suma con un caso terminal— porque el canal fuera de banda que ofrecía Channel con su close aquí sencillamente no está.
Declara el mutable como privado y publica asSharedFlow(). No es cosmética: MutableSharedFlow es también un FlowCollector, así que quien tenga la referencia puede emitir. El tipo es la única barrera real entre publicar y observar.
Difundir no es repartir
Hay una confusión que sobrevive a todo lo anterior y que conviene cortar de raíz, porque lleva a usar SharedFlow donde hacía falta un Channel. Un SharedFlow difunde: cada suscriptor recibe todos los valores, y dos suscriptores ven la misma secuencia. Un Channel reparte: cada elemento lo recibe exactamente un receptor, y dos receptores se dividen el trabajo. Son semánticas opuestas y ninguna emula bien a la otra.
flowchart LR P[Productor] --> SF[SharedFlow difunde] SF --> A[Suscriptor A ve todo] SF --> B[Suscriptor B ve todo] P2[Productor] --> CH[Channel reparte] CH --> C[Receptor C ve la mitad] CH --> D[Receptor D ve la otra mitad]
La consecuencia práctica es que un evento que debe procesarse una sola vez —cerrar sesión, navegar, enviar una petición— no pertenece a un SharedFlow si existe la posibilidad de que haya dos suscriptores, porque entonces se procesará dos veces. El síntoma clásico es una navegación duplicada tras una reconstrucción de pantalla en la que el suscriptor viejo aún no se había ido.
// Difusión: varios observadores del mismo hecho, todos lo ven.
val sesionExpirada: SharedFlow<Unit> = _sesionExpirada.asSharedFlow()
// Reparto: exactamente un consumidor debe atender cada trabajo.
private val trabajos = Channel<Trabajo>(capacity = 64)
La segunda diferencia es la terminación, ya mencionada pero decisiva aquí: un Channel se cierra y ese cierre viaja al receptor como información; un SharedFlow no tiene cierre. Si tu dominio necesita decir se acabó o falló, con SharedFlow tienes que codificarlo dentro del valor. Elegir entre ambos tipos, por tanto, se decide con dos preguntas encadenadas: cuántos deben ver cada elemento, y si existe un final que alguien tenga que observar.
Hay una asimetría profunda entre cómo se leen esos tres parámetros y lo que realmente significan. Se leen como afinado —un poco más de buffer aquí, un poco de historia allá— y significan una declaración jurídica sobre qué le debe tu emisor a tus suscriptores. Con los valores por defecto, la deuda es total y se salda en el acto: cuando emit retorna, todos los suscriptores presentes tienen el valor, sin excepción y sin demora; el precio es que un suscriptor lento congela al emisor y, a través de él, a la fuente entera. Esa propagación no es un defecto sino la única forma de que la lentitud sea observable en el sitio donde alguien puede hacer algo con ella. En el momento en que añades capacidad extra, la deuda deja de saldarse en el acto y pasa a diferirse: emit retorna diciendo aceptado, que es una palabra distinta de entregado, y la diferencia entre ambas —medida en número de valores en vuelo— es exactamente el número que escribiste en el constructor. En el momento en que además eliges una política de descarte, la deuda deja de existir: has declarado por adelantado que hay circunstancias en las que un evento emitido no llegará a nadie, y has elegido cuáles. Eso es perfectamente respetable, y de hecho es lo correcto para todo lo que sea estado observable, telemetría, posiciones, progresos o cualquier magnitud en la que el presente sustituya al pasado en lugar de sumarse a él. Lo que no es respetable es no saber cuál de los tres regímenes has elegido, y es lo habitual, porque el código que resulta es idéntico en el punto de emisión: la misma línea emit significa espero a que todos lo tengan, lo dejo en cola y sigo o lo tiro si no cabe dependiendo de tres números escritos en otro fichero. Por eso la disciplina útil no es aprenderse las políticas sino invertir la pregunta: antes de tocar el constructor, escribe la frase que describe qué se pierde y bajo qué condición; si no sabes escribirla, deja los valores por defecto, porque la contrapresión total es el único contrato que no requiere explicación y el único cuyo incumplimiento se manifiesta como lentitud visible en vez de como datos que faltan.
- Crea un
MutableSharedFlowcon repetición cero y capacidad extra cero, suscribe un recolector que tarde cien milisegundos por elemento y mide cuánto tarda un bucle de diezemit. Explica el número. - Añade
extraBufferCapacity = 5sin tocar nada más y repite la medición. Deduce de la diferencia el tamaño efectivo del colchón. - Monta el mismo flujo con
replay = 3y capacidad extra cero, y comprueba con el experimento anterior que la holgura del emisor también creció. Confirma tu hipótesis inspeccionandoreplayCache. - Escribe un
tryEmitsobre el flujo por defecto con un suscriptor lento y cuenta cuántos devuelvenfalse. Repite conDROP_OLDESTy explica por qué la cuenta cae a cero. - Suscribe dos recolectores de velocidades muy distintas y registra el retraso de cada uno frente al emisor. Identifica en la gráfica el instante exacto en que el lento empieza a gobernar al rápido.