SharedFlow acaba de ser introducido en coroutines 1.4.0-M1, y está destinado a reemplazar todas las implementaciones de BroadcastChannel (como se indica en la descripción del problema de diseño ).
Tengo un caso de uso en el que uso un BroadcastChannel para representar marcos de socket web entrantes, de modo que varios oyentes puedan "suscribirse" a los marcos. El problema que tengo cuando me muevo a SharedFlow es que no puedo "finalizar" el flujo cuando recibo un cuadro cerrado o un error ascendente (que me gustaría hacer para informar a todos los suscriptores que el flujo ha terminado).
¿Cómo puedo hacer que todas las suscripciones finalicen cuando quiero "cerrar" SharedFlow de manera efectiva? ¿Hay alguna manera de saber la diferencia entre el cierre normal y el cierre con excepción? (como canales)
Si MutableSharedFlow no permite transmitir el final del flujo a los suscriptores, ¿cuál es la alternativa si BroadcastChannel queda obsoleto/eliminado?
La documentación SharedFlow describe lo que necesita:
Tenga en cuenta que la mayoría de los operadores de terminal, como Flow.toList, tampoco se completarían cuando se aplicaran a un flujo compartido, pero los operadores de truncamiento de flujo como Flow.take y Flow.takeWhile se pueden usar en un flujo compartido para convertirlo en uno de finalización.
SharedFlow no se puede cerrar como BroadcastChannel y nunca puede representar una falla. Todos los errores y señales de finalización deben materializarse explícitamente si es necesario.
Básicamente, deberá introducir un objeto especial que pueda emitir desde el flujo compartido para indicar que el flujo ha finalizado, utilizando takeWhile en el extremo del consumidor puede hacer que emitan hasta que se reciba ese objeto especial.
Creo que una posible solución es crear un indicador booleano isValid y exponer públicamente solo los flujos con .takeWhile { isValid } . Luego simplemente llame a isValid = false y sFlow.emit() cuando desee cerrar todos los suscriptores.
Posible implementación:
private var isValid = true // In real scenario use atomic boolean private val _sharedFlow = MutableSharedFlow<Unit>() val sharedFlow: Flow<Unit> get() = _sharedFlow.takeWhile { isValid } suspend fun cancelSharedFlow() { isValid = false _sharedFlow.emit(Unit) } EDITAR: En mi caso .emit() siempre se suspendía, así que tuve que usar BufferOverflow.DROP_LATEST (que no es adecuado para muchos casos de uso). No estoy seguro si el problema está en este ejemplo o en otra parte de mi aplicación. Si ves un problema, por favor comenta :)