JavaRush /Курсы /Kotlin SELF /SharedFlow как шина событий: collect, replay и жизненный ...

SharedFlow как шина событий: collect, replay и жизненный цикл подписки

Kotlin SELF
57 уровень , 4 лекция
Открыта

1. Публикация и подписка в SharedFlow

Если вы только что написали аккуратный Event<T> со snapshot и токеном отписки, то логичный вопрос звучит так: «А зачем мне ещё какой-то SharedFlow? Я же только что победил хаос». И это правильный вопрос: свой Event<T> действительно годится для синхронной рассылки в рамках обычного кода.

Но как только мы оказываемся в мире корутин, появляется новый стиль жизни: подписчики становятся корутинами, «жизненный цикл» подписки удобнее контролировать через Job, а сама рассылка событий естественно ложится в модель потоков данных (Flow). В Kotlin есть несколько способов обмена данными между корутинами, и SharedFlow как раз описывается как механизм, который «делится каждым значением со всеми активными коллекторами».

Смысл лекции: научиться использовать SharedFlow как «шину событий» в корутинном коде, понимать, как устроены подписка и «отписка», и почему параметр replay делает учебные примеры стабильными и не зависящими от «успел/не успел».

MutableSharedFlow и SharedFlow: «пульт» и «экран»

Когда вы впервые видите MutableSharedFlow<T>, может возникнуть желание выдать его всем подряд: «пусть кто хочет — тот и emit-ит». Это примерно как оставить кнопку «пуск ракеты» на кухонном столе рядом с печеньем. Формально удобно, но потом вы полдня ищете, кто и почему «выстрелил» лишнее событие.

В корутинном мире принято разделять роли: MutableSharedFlow<T> — это точка публикации (туда «пушат» события), а SharedFlow<T> — безопасный вид «только для чтения» (оттуда события только читают через collect). Этот подход хорошо согласуется с общей идеей про разделение обязанностей и уменьшение связанности: наружу мы отдаём подписку, но не отдаём возможность публиковать откуда попало.

Сделаем минимальный «event bus» для нашего консольного приложения (пусть это будет небольшой CLI‑трекер задач, где есть команды добавления/удаления, а события нужны для аудита и уведомлений).

import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.asSharedFlow

sealed interface AppEvent
data class TaskAdded(val text: String) : AppEvent

class EventBus {
    private val _events = MutableSharedFlow<AppEvent>(replay = 1)
    val events: SharedFlow<AppEvent> = _events.asSharedFlow()

    suspend fun emit(event: AppEvent) = _events.emit(event)
}

Обратите внимание на две вещи. Во-первых, _events — приватный, «мутабельный» (его можно менять, в него можно emit). Во-вторых, наружу торчит events: SharedFlow<AppEvent> — подписывайтесь сколько хотите, но «публиковать» напрямую нельзя.

Подписка через collect: «слушатель» живёт в корутине

В нашей самописной системе событий слушатель был обычной функцией (T) -> Unit, которую мы складывали в список. В SharedFlow слушатель превращается в корутину-коллектор: вы запускаете launch { flow.collect { ... } }, и пока эта корутина жива — она получает события.

Это очень удобно для управления жизненным циклом: «отписка» — это просто отмена Job. Никаких поисков «той самой лямбды», никаких remove(listener). Подписка живёт ровно столько, сколько живёт корутина, и это отлично совпадает с идеей structured concurrency: дочерние задачи живут внутри родительского scope.

Минимальный пример «два подписчика получают все события»:

import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    val bus = EventBus()

    val a = launch {
        bus.events.collect { println("A got: $it") }
    }
    val b = launch {
        bus.events.collect { println("B got: $it") }
    }

    bus.emit(TaskAdded("Buy milk")) // A got: TaskAdded(text=Buy milk), B got: TaskAdded(text=Buy milk)

    a.cancel() // отписались
    b.cancel() // отписались
}

Здесь важна мысль: collect — это «подписка», которая обычно не заканчивается сама по себе. Если вы не отмените job, коллектор будет ждать новые события. Для «event bus» это нормально: события потенциально бесконечны.

«Отписка» в SharedFlow: отменяем Job, а не ищем лямбду

В лекции про токены отписки мы подчёркивали «ссылочную природу лямбд»: отписаться можно только той же ссылкой. В SharedFlow этот пласт проблем почти исчезает, потому что подписка — это не «элемент в списке», а выполняющаяся корутина.

То есть практический рецепт отписки выглядит так: сохранили Job от launch, потом сделали cancel() (иногда ещё join() — но сегодня нам достаточно просто отмены). Это приятно ещё и тем, что отмена корутины — привычный инструмент, который вы уже знаете по прошлым дням про structured concurrency.

Супер-короткий пример «подписались — получили одно событие — отписались»:

import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    val bus = EventBus()

    val job = launch {
        bus.events.collect { println("Audit: $it") }
    }

    bus.emit(TaskAdded("Read Kotlin docs")) // Audit: TaskAdded(text=Read Kotlin docs)
    job.cancel() // подписка остановлена
}

Если вам нужен «временный слушатель» (например, «логировать только пока идёт сценарий»), Job.cancel() — это прямо то, что доктор прописал.

3. replay и «поздние» подписчики

Самая частая учебная боль с потоками событий звучит так: «Я сделал emit, но подписчик ничего не получил. Kotlin сломался?». На самом деле нет: это вы случайно показали классический сценарий late subscriber — подписчик начал слушать после того, как событие уже улетело.

По определению обычный Flow производит значения только когда его собирают (collect), а SharedFlow раздаёт значения всем активным коллекторам. Но если коллекторов в момент события не было, то для SharedFlow с replay = 0 это событие, по сути, «никому не нужно», и оно не сохраняется.

Вот почему параметр replay так полезен именно в учебных примерах: он позволяет хранить последние N событий и «переигрывать» их новым подписчикам. Мы в этой лекции используем replay = 1, чтобы новый подписчик гарантированно увидел хотя бы «последний факт» и пример не зависел от таймингов.

Посмотрим на два поведения: сначала (мысленно) представьте replay = 0, затем — replay = 1. Мы оставим наш EventBus с replay = 1, но покажем, почему это важно.

Поздняя подписка: событие раньше, подписчик позже

import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    val bus = EventBus()

    bus.emit(TaskAdded("EARLY")) // подписчиков нет, но replay=1 запомнит последнее

    val job = launch {
        bus.events.collect { println("Late got: $it") }
    }

    // Благодаря replay=1 подписчик сразу получит "EARLY"
    job.cancel()
}

Если бы replay был равен 0, строка Late got: ... могла бы не появиться вообще, и студент (а иногда и преподаватель) начинал бы подозревать заговор таймера и планировщика.

Что делает replay = 1 в человеческих терминах

Можно воспринимать replay = 1 как «подушку безопасности для последнего события». Это не превращает события в «состояние» (мы не уходим в эту тему глубоко), но для демонстраций и простых сценариев «последний факт важен» делает поведение предсказуемым: подписчик подключился — сразу увидел последнее.

Важно договориться с собой, что именно вы строите: «чистые события» без памяти или «события, где последнее значение важно показать новому наблюдателю». В учебном коде часто выгоднее второе, потому что оно устойчивее к порядку запуска корутин.

4. Ошибки в обработчиках collect и «умирающие» подписки

Когда вы писали свой Event<T>, мы оборачивали каждый обработчик в try/catch, чтобы падение одного слушателя не ломало рассылку всем. В SharedFlow похожая проблема проявляется иначе: если внутри collect { ... } вы бросите исключение и не поймаете его, то корутина-коллектор завершится. А раз корутина завершилась — подписка закончилась.

То есть «одна ошибка в обработчике» может превратиться в «подписчик молча умер и больше никогда не слушает». Это особенно неприятно, потому что внешне кажется, что «события перестали приходить».

Поэтому, если обработчик может падать, имеет смысл изолировать ошибки прямо внутри collect:

import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    val bus = EventBus()

    val job = launch {
        bus.events.collect { e ->
            try {
                if (e is TaskAdded && e.text == "BAD") error("broken handler")
                println("OK: $e")
            } catch (ex: Exception) {
                println("handler error: ${ex.message}") // handler error: broken handler
            }
        }
    }

    bus.emit(TaskAdded("GOOD")) // OK: TaskAdded(text=GOOD)
    bus.emit(TaskAdded("BAD"))  // handler error: broken handler
    bus.emit(TaskAdded("GOOD2"))// OK: TaskAdded(text=GOOD2)

    job.cancel()
}

Тут есть приятный эффект: подписка не умирает из-за одного проблемного события. Вы как бы повторяете идею «изоляции ошибок обработчиков», только уже в корутинном стиле.

5. Пример: «команды → события → реакции»

Чтобы SharedFlow не выглядел как «магия ради магии», давайте встроим его в понятный бытовой сценарий. Пусть у нас есть список задач (MutableList<String>) и команда add. Внутри команды мы хотим: добавить задачу и опубликовать событие. Отдельные подписчики будут делать аудит и показывать пользователю уведомление.

Схема процесса может выглядеть так:

flowchart LR
    CLI[Команда в CLI] --> Domain[Логика add/remove]
    Domain --> Bus[EventBus / SharedFlow]
    Bus --> Audit[Audit subscriber]
    Bus --> Ui[UI subscriber]

Доменные функции публикуют событие и не знают про подписчиков

import kotlinx.coroutines.runBlocking

fun addTask(tasks: MutableList<String>, text: String, bus: EventBus) = runBlocking {
    tasks += text
    bus.emit(TaskAdded(text))
}

Да, тут runBlocking выглядит грубовато (и в реальном приложении вы бы сделали это аккуратнее), но для учебного консольного примера идея проста: действие произошло — событие опубликовано. Источник события не знает, кто слушает.

Подписчики подключаются в main и живут, пока живёт приложение

import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.flow.collect

fun main() = runBlocking {
    val tasks = mutableListOf<String>()
    val bus = EventBus()

    val auditJob = launch {
        bus.events.collect { println("AUDIT: $it") }
    }
    val uiJob = launch {
        bus.events.collect { println("UI: task list changed") }
    }

    addTask(tasks, "Learn SharedFlow", bus)
    // AUDIT: TaskAdded(text=Learn SharedFlow)
    // UI: task list changed

    auditJob.cancel()
    uiJob.cancel()
}

Смешной (и жизненный) момент: UI‑подписчик в консоли — это просто println. Но архитектурно это тот же приём, что и в больших приложениях: один и тот же факт «задача добавлена» может интересовать разные части системы.

6. Event<T> vs SharedFlow: когда что выбирать

Иногда правильный вывод из этой темы звучит не как «SharedFlow лучше», а как «это разные инструменты под разные ситуации». Самописный Event<T> отлично подходит для простой синхронной рассылки, где вы хотите полный контроль и минимум зависимостей. SharedFlow удобен, когда вы уже живёте в корутинах и хотите подписки/отписки через Job и единый стиль с Flow.

Сравним в табличке (без погружения в тонкие настройки):

Критерий
Event<T>
(самописный)
SharedFlow<T>
Модель подписки список функций (T) -> Unit корутина + collect { ... }
Отписка токен, удаление слушателя
job.cancel()
Ошибка слушателя ловим try/catch вокруг вызова если не ловить, корутина-коллектор завершится
Сохранение «последнего события» надо писать руками replay = 1 (и больше)
Связь с корутинами не обязана родная среда

Если где-то вам проще «обычный список слушателей» — берите Event<T>. Если вы уже в корутинном коде и хотите потоковую модель — берите SharedFlow. Kotlin как раз и хорош тем, что не заставляет выбирать один стиль навсегда.

7. Типичные ошибки при работе с SharedFlow как с шиной событий

Ошибка №1: делать MutableSharedFlow публичным.
Когда наружу торчит именно MutableSharedFlow, любой кусок кода может начать «эмитить» события. В результате события перестают быть «истиной от источника», а становятся «кто во что горазд». Гораздо устойчивее держать MutableSharedFlow приватно, а наружу отдавать SharedFlow, чтобы внешние части могли только подписываться.

Ошибка №2: забывать, что collect обычно не завершается сам.
Новичок часто пишет launch { flow.collect { ... } } и ждёт, что корутина «закончится» после одного события. Но collect — это подписка, и она будет ждать дальше. Если вы не отмените Job, подписчик останется активным. Это не баг: это контракт. Просто держите в голове жизненный цикл.

Ошибка №3: ожидать прошлые события при replay = 0.
Если replay равен 0, поздние подписчики не обязаны получать прошлые события — они увидят только будущие. Для учебных примеров и демонстраций лучше ставить replay = 1, иначе вы начинаете зависеть от таймингов: кто первый стартанул — тот и молодец.

Ошибка №4: не ловить исключения внутри обработчика collect.
Если обработчик внутри collect бросает исключение, корутина-коллектор может завершиться, и подписка «умрёт». Потом вы будете долго думать, почему события перестали приходить. Если обработчик потенциально «ломкий», изолируйте ошибки через try/catch прямо внутри collect, как мы делали раньше для отдельных слушателей в Event<T>.

Ошибка №5: путать «события» и «состояние» и ставить replay без понимания.
replay = 1 делает поведение стабильнее, но это уже «немножко память». Если вы бездумно увеличите replay или начнёте полагаться на него как на «хранилище», вы легко получите странные эффекты: новый подписчик видит старые события и реагирует так, будто они только что произошли. В этой лекции replay = 1 — осознанная учебная настройка, а не универсальная рекомендация на все случаи жизни.

Ошибка №6: думать, что SharedFlow автоматически делает обработчики безопасными от реентерабельности.
В ручной реализации событий мы решали проблему «модификация списка слушателей во время рассылки» через snapshot (toList()), и это работало, потому что мы контролировали структуру данных. В SharedFlow у вас уже другая модель: обработчики живут в корутинах, и важно следить за тем, что они делают внутри. Если обработчик по событию запускает новую публикацию событий, вы можете получить каскад. Это не обязательно плохо, но это должно быть осознанным дизайном, а не случайной реакцией «ой, давайте тут ещё emit сделаем».

1
Задача
Kotlin SELF, 57 уровень, 4 лекция
Недоступна
Два наблюдателя
Два наблюдателя
1
Задача
Kotlin SELF, 57 уровень, 4 лекция
Недоступна
Поздняя подписка
Поздняя подписка
1
Задача
Kotlin SELF, 57 уровень, 4 лекция
Недоступна
Шина событий
Шина событий
1
Задача
Kotlin SELF, 57 уровень, 4 лекция
Недоступна
Ошибка обработчика
Ошибка обработчика
1
Опрос
События/Observer, 57 уровень, 4 лекция
Недоступен
События/Observer
События/Observer
Комментарии (1)
ЧТОБЫ ПОСМОТРЕТЬ ВСЕ КОММЕНТАРИИ ИЛИ ОСТАВИТЬ КОММЕНТАРИЙ,
ПЕРЕЙДИТЕ В ПОЛНУЮ ВЕРСИЮ
kasnil Уровень 66
2 мая 2026
В Kotlin различаются две категории потоков - холодные и горячие. Холодные потоки представляют собой асинхронные потоки данных. Они начинют производить элементы, только когда их элементы потребляются отдельным коллектором. Примером холодного потока является: Flow. Горячие потоки работают в режиме трансляции и производят элементы независимо от того, потребляются ли они на самом деле. Примером горячего потока является: SharedFlow. Т.к. SharedFlow это горячий поток, то и появился replay, который создает кеш последних значений для новых подписчиков. Стоит отметить, что использование подчеркивания для приватной переменной в начале имени для MutableSharedFlow и имени без подчеркивания для публичной переменной типа SharedFlow, связано с тем, что Kotlin не поддерживает различные типы для приватных и публичных переменных. Есть issue Support having a "public" and a "private" type for the same property по добавлению возможности указывать тип свойства в зависимости от точки обращения - внутри или вне класса c версии Kotlin 2.4.0-Beta2.