Общая рассылка в Kotlin
SharedFlow похож на шину событий:
отправитель вызывает emit или
tryEmit, активные подписчики
получают значения. В отличие от
StateFlow, здесь нет обязательной
ячейки «текущее» - поведение для опоздавшего
задает параметр replay.
При replay = 0 новый collect
не видит прошлые элементы - только то, что
придет после подписки:
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.collect
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
fun main() {
runBlocking {
val bus = MutableSharedFlow<String>(
replay = 0,
extraBufferCapacity = 1
)
bus.tryEmit("early")
launch {
bus.collect { msg ->
println("sub $msg")
}
}
delay(20)
bus.emit("late")
delay(20)
}
}
Подписчик не печатает "early" - событие
ушло до старта сбора. "late" уже
попадает в активную подписку.
Несколько подписчиков на одном объекте получают одни и те же новые элементы параллельно:
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.collect
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
fun main() {
runBlocking {
val hub = MutableSharedFlow<Int>(replay = 0)
launch {
hub.collect { n -> println("one $n") }
}
launch {
hub.collect { n -> println("two $n") }
}
delay(20)
hub.emit(7)
delay(20)
}
}
Число 7 дублируется в двух строках -
рассылка общая. Без replay опоздавший
коллектор не восстановит пропущенное, пока
источник снова не отправит данные.
Шина с replay равным нулю принимает
строку "ping" до подписки, затем
запускается сбор в фоне и отправляется
"pong". В консоли должна быть только
вторая строка.
Два фоновых сборщика на одном
MutableSharedFlow без повтора должны
оба напечатать метку "sync" после
одной отправки этого текста.