wandres.dev
FLOW II · StateFlow y flujos calientes

shareIn y stateIn: calentar un flujo frío y compartir lo que cuesta caro

Los operadores shareIn y stateIn son la frontera oficial entre los dos mundos: reciben un flujo frío, lo colectan una sola vez dentro de un ámbito que tú eliges y reparten el resultado a todos los suscriptores. Esta lección examina las tres piezas que hay que decidir en cada llamada —el ámbito que gobierna la vida del productor, la política de arranque y parada, y el buffer de repetición—, explica por qué WhileSubscribed con un retardo pequeño es la respuesta correcta casi siempre, y muestra la trampa de colocar el operador dentro de una función que se invoca en cada recomposición.

⏱ 24 min

Un flujo frío ejecuta su productor una vez por cada recolector. Cuando ese productor es barato —transformar una lista, leer memoria— la duplicación es irrelevante. Cuando es una consulta a red, una suscripción a un servicio del sistema o un observador de base de datos, la duplicación es el problema. shareIn y stateIn existen para resolver exactamente eso: colectan la fuente una sola vez desde una corrutina propia y multiplican el resultado hacia todos los interesados. A cambio te obligan a contestar tres preguntas que en el mundo frío no existían —quién mantiene viva esa corrutina, cuándo debe arrancar y parar, y qué hereda quien llegue tarde—, y la calidad de tu sistema caliente será exactamente la calidad de esas tres respuestas.

🎯 Al terminar esta lección sabrás
  • Describir la mecánica de shareIn y stateIn como una única colección compartida dentro de un ámbito explícito.
  • Comparar las tres políticas de SharingStarted en términos de coste, latencia y pérdida.
  • Justificar el retardo de parada de WhileSubscribed a partir de los cambios de configuración y las reconexiones breves.
  • Detectar el error de crear el flujo compartido en un punto que se ejecuta más de una vez.

Una sola colección repartida entre muchos

La firma de ambos operadores es casi la misma: reciben el CoroutineScope donde vivirá el productor y una política de arranque. shareIn pide además cuántos valores repetir y devuelve un SharedFlow<T>; stateIn pide un valor inicial y devuelve un StateFlow<T>, con toda la conflación por igualdad que eso implica.

private val fuenteCara: Flow<Cotizacion> = callbackFlow { /* un socket */ }

val cotizaciones: SharedFlow<Cotizacion> = fuenteCara.shareIn(
    scope = ambitoDelServicio,
    started = SharingStarted.WhileSubscribed(5_000),
    replay = 1,
)

val ultimaCotizacion: StateFlow<Cotizacion?> = fuenteCara.stateIn(
    scope = ambitoDelServicio,
    started = SharingStarted.WhileSubscribed(5_000),
    initialValue = null,
)

Lo que ocurre internamente es sencillo de enunciar y conviene tenerlo presente: se crea un MutableSharedFlow con la configuración pedida y se lanza en el ámbito una corrutina que colecta la fuente y reemite en él. Nada más. De ahí se deducen las consecuencias importantes: el productor se ejecuta en el contexto de ese ámbito y no en el del recolector, su cancelación depende del ámbito y no del recolector, y una excepción de la fuente cae en el manejador del ámbito, no en el catch de quien colecta. Si quieres que un fallo de red no derribe nada, el catch va antes del shareIn, sobre el flujo frío, donde todavía puedes convertirlo en un valor.

Existe también una variante suspendida de stateIn sin valor inicial: espera al primer valor de la fuente y devuelve un StateFlow ya poblado. Es honesta —no inventa un estado que no existe— y es peligrosa en arranques, porque suspende hasta que la fuente responda.

ℹ️
El operador va al final de la cadena

Todo lo que pongas después de shareIn vuelve a ejecutarse una vez por suscriptor, porque el resultado se sigue colectando. Compartir tarde ahorra menos de lo que crees: coloca el operador justo detrás del trabajo caro y antes de las transformaciones baratas y específicas de cada consumidor.

Las tres políticas de arranque

SharingStarted es una interfaz con tres implementaciones de fábrica que responden a la pregunta de cuándo debe estar viva la corrutina productora. La diferencia entre ellas no es de rendimiento, es de qué estás dispuesto a pagar cuando nadie mira.

🔥

Eagerly

Arranca al crear el flujo compartido y no para nunca hasta que muera el ámbito. Máxima frescura, coste permanente y emisiones perdidas si el buffer de repetición es cero.

🐢

Lazily

Arranca con el primer suscriptor y tampoco para nunca. Evita el coste si nadie llega a mirar, pero una vez encendido se queda encendido aunque se vayan todos.

⏱️

WhileSubscribed

Arranca con el primero y para cuando se va el último, opcionalmente tras un retardo. Es la única que devuelve recursos y la única que sobrevive bien a consumidores efímeros.

WhileSubscribed acepta dos parámetros temporales cuya diferencia hay que tener clarísima. stopTimeoutMillis es cuánto espera tras marcharse el último suscriptor antes de cancelar el productor, y existe porque un consumidor de interfaz se destruye y se reconstruye constantemente: sin ese margen, cada rotación de pantalla cancelaría el socket y volvería a abrirlo. replayExpirationMillis es cuánto sobrevive el buffer de repetición una vez parada la producción; pasado ese tiempo la caché se olvida y el siguiente suscriptor no verá un valor rancio.

SharingStarted.WhileSubscribed(
    stopTimeoutMillis = 5_000,        // sobrevive a un cambio de configuración
    replayExpirationMillis = 60_000,  // pasado un minuto, el dato ya no vale
)
flowchart TD
A[Llega el primer suscriptor] --> B[Arranca la corrutina productora]
B --> C[Se reparte a todos los suscriptores]
C --> D{Se va el ultimo}
D -->|No| C
D -->|Si| E[Espera el retardo de parada]
E --> F{Vuelve alguien antes del plazo}
F -->|Si| C
F -->|No| G[Se cancela la produccion y se libera el recurso]

Elegir entre las tres se reduce a dos preguntas encadenadas. Primera: ¿cuesta algo mantener viva la fuente? Si no cuesta nada, Eagerly es aceptable y ahorra latencia. Segunda, si sí cuesta: ¿los consumidores van y vienen? Si van y vienen, WhileSubscribed con un retardo de unos segundos; si hay un consumidor estable durante toda la vida del ámbito, Lazily es equivalente y más simple.

Dónde se crea importa más que cómo

El error más caro de este operador no está en sus parámetros sino en el sitio donde se escribe la llamada. shareIn y stateIn crean un objeto nuevo cada vez que se ejecutan. Si la llamada está dentro de una función que se invoca repetidamente, cada invocación fabrica un flujo compartido distinto, con su propia corrutina productora, y el efecto neto es peor que no haber compartido nada: mismo número de suscripciones a la fuente, más una fuga por cada una que quede huérfana.

// Mal: cada llamada crea un compartido nuevo y una corrutina más.
fun cotizaciones(): SharedFlow<Cotizacion> =
    fuenteCara.shareIn(ambito, SharingStarted.Lazily, 1)

// Bien: una sola instancia, creada una vez, con vida ligada al ámbito.
val cotizaciones: SharedFlow<Cotizacion> =
    fuenteCara.shareIn(ambito, SharingStarted.WhileSubscribed(5_000), 1)

Cuando el flujo depende de un parámetro —una consulta por identificador— la propiedad no basta, y la solución correcta es una caché explícita indexada por ese parámetro, no una función que fabrica. Si el parámetro cambia con el tiempo, la construcción idiomática es partir de un StateFlow del parámetro y usar flatMapLatest, compartiendo el resultado una sola vez al final.

val detalle: StateFlow<Detalle> = idSeleccionado
    .flatMapLatest { repositorio.observar(it) }   // frío hasta aquí
    .stateIn(ambito, SharingStarted.WhileSubscribed(5_000), Detalle.Vacio)

Queda un matiz de coste que casi nadie mide: compartir no es gratis. Añade una corrutina, un buffer, una indirección por emisión y una clase de errores nueva. Para un flujo cuyo productor es una transformación en memoria, shareIn gasta más de lo que ahorra. El operador se justifica cuando el productor tiene un coste externo —conexión, consulta, sensor— o cuando la semántica exige que todos los observadores vean exactamente la misma secuencia, que es un requisito distinto y a veces más importante que el ahorro.

💡
El ámbito es la decisión de fondo

Antes de discutir políticas, mira qué CoroutineScope estás pasando. Si su vida es más larga que la utilidad del dato, ninguna política te salvará: WhileSubscribed te devolverá el recurso, pero el objeto compartido y su buffer seguirán ahí mientras el ámbito exista.

Qué se expone en el borde del módulo

Compartir tiene además una dimensión de diseño de API que suele decidirse por descuido. Un repositorio que devuelve Flow<T> está diciendo pídeme esto cuando lo necesites y pagaré por consulta; uno que devuelve SharedFlow<T> o StateFlow<T> está diciendo yo mantengo esto vivo, tú asómate. La segunda promesa es mucho más fuerte: obliga al módulo a poseer un ámbito, a gestionar su ciclo de vida y a responder por lo que ocurra cuando nadie mire.

class RepositorioCotizaciones(private val ambito: CoroutineScope) {
    // Frío hacia fuera: cada consumidor decide cuándo y cuánto.
    fun historico(dia: Fecha): Flow<Cotizacion> = fuente.historico(dia)

    // Caliente hacia fuera: el repositorio posee la conexión y la comparte.
    val enVivo: StateFlow<Cotizacion?> = fuente.enVivo()
        .stateIn(ambito, SharingStarted.WhileSubscribed(5_000), null)
}

La regla que se sigue de ahí es que el calentamiento debe ocurrir donde vive el recurso, no donde resulta cómodo. Compartir en la capa de presentación un flujo cuya fuente cara está en la capa de datos duplica el coste en cuanto haya dos pantallas, porque cada una crea su propio compartido sobre el mismo flujo frío. Compartir en la capa de datos lo resuelve para todos y de paso concentra la decisión de política en un solo sitio auditable.

Queda una consideración sobre pruebas que se olvida siempre: un flujo compartido con Eagerly empieza a producir en el momento en que se construye el objeto que lo contiene, lo cual significa que instanciar la clase en un test ya arranca trabajo. Con WhileSubscribed no ocurre, y por eso esa política no solo es mejor para el recurso: también hace que los tests sean deterministas, porque nada sucede hasta que alguien colecta.

Por último, conviene recordar que estos operadores no son magia irreversible: siempre puedes volver al mundo frío conservando el ahorro, colocando después del compartido los operadores específicos de cada consumidor. Esa es la forma idiomática de tener una sola conexión y, sin embargo, diez vistas distintas de ella.

val soloErrores: Flow<Cotizacion> = cotizaciones.filter { it.esAnomala }
val redondeadas: Flow<Int> = cotizaciones.map { it.valor.toInt() }
Compartir un flujo no es una optimización local: es declarar que existe una sola versión de la verdad y decidir quién paga por mantenerla viva cuando nadie la mira

El razonamiento habitual para introducir shareIn es económico y va así: la fuente es cara, hay varios consumidores, luego conviene colectarla una vez. Ese razonamiento es correcto y sin embargo omite lo más importante, porque el efecto principal del operador no es ahorrar sino unificar. En cuanto compartes, dejas de tener varias ejecuciones independientes de una descripción y pasas a tener una única instancia viva de la fuente, con una historia concreta, un instante de arranque concreto y una secuencia de valores que ya no depende de quién mire ni de cuándo. Eso cambia la naturaleza de las garantías que puedes ofrecer: en el mundo frío, dos pantallas que consultan lo mismo pueden ver valores distintos porque cada una hizo su propia consulta en su propio instante, y esa divergencia es invisible hasta el día en que dos partes de la interfaz muestran cifras que no cuadran; en el mundo compartido, esa divergencia es imposible por construcción, y a cambio aparecen otras dos que no existían. La primera es la de la historia: quien llega tarde no ve lo que pasó antes salvo que hayas comprado repetición, así que la pregunta qué sabe un consumidor recién nacido pasa a tener una respuesta configurable, y si la contestas mal —repetición cero para un estado, repetición uno para un evento— el fallo no aparecerá en el arranque sino la primera vez que alguien se suscriba en un momento inconveniente. La segunda es la del vacío: hay un intervalo en que nadie mira, y hay que decidir qué ocurre entonces, que es una decisión de negocio disfrazada de constante numérica, porque mantener viva la fuente cuando nadie escucha significa gastar batería, cuota o conexiones para que el primer valor tras el regreso sea instantáneo, y apagarla significa ahorrar todo eso a cambio de una latencia y, muchas veces, de un hueco en la serie. WhileSubscribed con un retardo de unos segundos se ha convertido en la respuesta por defecto no porque sea óptima sino porque codifica una observación empírica muy sólida: los consumidores de interfaz desaparecen y reaparecen continuamente por motivos que no tienen nada que ver con el interés del usuario en el dato, y confundir esa desaparición técnica con una pérdida de interés real es lo que produce esos sistemas que reconectan un socket cada vez que giras el teléfono. El número que escribas ahí no es un plazo: es tu definición operativa de cuándo un usuario ha dejado de estar interesado.

⚔️ Mide lo que ahorras y lo que pagas
  1. Instrumenta un flujo frío caro con un contador de arranques y coléctalo desde tres corrutinas. Anota el contador; añade shareIn y repite. La diferencia es el ahorro real.
  2. Prueba las tres políticas sobre esa misma fuente registrando el instante de arranque y de parada. Dibuja la línea temporal de cada una con dos suscriptores que entran y salen.
  3. Ajusta stopTimeoutMillis a cero y simula una reconstrucción del consumidor. Cuenta las reconexiones. Súbelo a cinco segundos y vuelve a contarlas.
  4. Escribe deliberadamente la versión con función que crea el compartido en cada llamada, invócala cinco veces y demuestra con el contador de arranques que has empeorado el sistema.
  5. Coloca un catch después de shareIn y provoca un fallo en la fuente. Comprueba que no lo captura, muévelo antes del operador y explica en dos frases por qué ahora sí.