StateFlow и SharedFlow

StateFlow и SharedFlow – это API потоков, которые позволяют потокам оптимально передавать обновления состояния и значения нескольким потребителям.

StateFlow

StateFlow – это наблюдаемый поток, который хранит состояние и передает текущие и новые обновления состояния своим коллекторам. Текущее значение состояния также можно прочитать с помощью свойства value. Чтобы обновить состояние и отправить его в поток, присвойте новое значение свойству value класса MutableStateFlow.

В Android класс StateFlow отлично подходит для классов, которым необходимо поддерживать наблюдаемое изменяемое состояние.

Следуя примерам из раздела Потоки Kotlin, StateFlow можно получить из LatestNewsViewModel, чтобы View мог отслеживать изменения состояния интерфейса и сохранять состояние экрана при изменении конфигурации.

class LatestNewsViewModel(
    private val newsRepository: NewsRepository
) : ViewModel() {

    // Backing property to avoid state updates from other classes
    private val _uiState = MutableStateFlow(LatestNewsUiState.Success(emptyList()))
    // The UI collects from this StateFlow to get its state updates
    val uiState: StateFlow<LatestNewsUiState> = _uiState

    init {
        viewModelScope.launch {
            newsRepository.favoriteLatestNews
                // Update UI with the latest favorite news
                // Writes to the value property of MutableStateFlow,
                // adding a new element to the flow and updating all
                // of its collectors
                .collect { favoriteNews ->
                    _uiState.value = LatestNewsUiState.Success(favoriteNews)
                }
        }
    }
}

// Represents different states for the LatestNews screen
sealed class LatestNewsUiState {
    data class Success(val news: List<ArticleHeadline>) : LatestNewsUiState()
    data class Error(val exception: Throwable) : LatestNewsUiState()
}

Класс, отвечающий за обновление MutableStateFlow, является производителем, а все классы, получающие данные из StateFlow, – потребителями. В отличие от холодного потока, созданного с помощью конструктора flow, StateFlow является горячим: при сборе данных из потока не запускается код производителя. Объект StateFlow всегда активен и находится в памяти. Он становится доступен для сборки мусора только тогда, когда на него нет ссылок из корня сборки мусора.

Когда новый потребитель начинает получать данные из потока, он получает последнее состояние в потоке и все последующие состояния. Такое поведение можно наблюдать и в других классах, например в LiveData.

View ожидает StateFlow, как и в любом другом потоке:

class LatestNewsActivity : ComponentActivity() {
    private val latestNewsViewModel: LatestNewsViewModel = TODO() // getViewModel()

    override fun onCreate(savedInstanceState: Bundle?) {
        // ...
        // Start a coroutine in the lifecycle scope
        lifecycleScope.launch {
            // repeatOnLifecycle launches the block in a new coroutine every time the
            // lifecycle is in the STARTED state (or above) and cancels it when it's STOPPED.
            repeatOnLifecycle(Lifecycle.State.STARTED) {
                // Trigger the flow and start listening for values.
                // Note that this happens when lifecycle is STARTED and stops
                // collecting when the lifecycle is STOPPED
                latestNewsViewModel.uiState.collect { uiState ->
                    // New value received
                    when (uiState) {
                        is LatestNewsUiState.Success -> showFavoriteNews(uiState.news)
                        is LatestNewsUiState.Error -> showError(uiState.exception)
                    }
                }
            }
        }
    }
}

Чтобы преобразовать любой поток в StateFlow, используйте промежуточный оператор stateIn.

StateFlow, Flow и LiveData

StateFlow и LiveData похожи. Оба класса являются наблюдаемыми держателями данных и следуют схожей схеме при использовании в архитектуре приложения.

Однако обратите внимание, что StateFlow и LiveData работают по-разному:

  • StateFlow требует, чтобы в конструктор передавалось начальное состояние, а LiveData – нет.
  • LiveData.observe() автоматически отменяет регистрацию потребителя, когда представление переходит в состояние STOPPED, тогда как сбор данных из StateFlow или любого другого потока не останавливается автоматически. Чтобы добиться такого же поведения, вам нужно собрать поток из блока Lifecycle.repeatOnLifecycle.

Как сделать холодные потоки горячими с помощью shareIn

StateFlow – это горячий поток. Он остается в памяти, пока собирается или пока на него есть ссылки из корня сборки мусора. Вы можете сделать холодные потоки горячими, используя оператор shareIn.

Например, если вы используете callbackFlow, созданный в потоках Kotlin, то вместо того, чтобы каждый сборщик создавал новый поток, вы можете передавать данные, полученные из Firestore, между сборщиками с помощью shareIn. Вам нужно передать следующие данные:

  • CoroutineScope, с помощью которого можно поделиться потоком. Область действия должна быть больше, чем у любого потребителя, чтобы общий поток оставался активным столько, сколько нужно.
  • Количество элементов, которые нужно воспроизвести для каждого нового коллектора.
  • Правило поведения при запуске.

class NewsRemoteDataSource(
    private val externalScope: CoroutineScope
) {
    val latestNews: Flow<List<ArticleHeadline>> = flow {
        // ...
    }.shareIn(
        externalScope,
        replay = 1,
        started = SharingStarted.WhileSubscribed()
    )
}

В этом примере поток latestNews воспроизводит последний отправленный элемент для нового сборщика и остается активным, пока externalScope работает и есть активные сборщики. Правило запуска SharingStarted.WhileSubscribed() сохраняет активность вышестоящего производителя, пока есть активные подписчики. Также доступны другие правила запуска, например SharingStarted.Eagerly для немедленного запуска продюсера или SharingStarted.Lazily для начала трансляции после появления первого подписчика и ее постоянной активности.

SharedFlow

Функция shareIn возвращает SharedFlow – горячий поток, который передает значения всем потребителям, собирающим данные из него. SharedFlow – это обобщенная версия StateFlow с широкими возможностями настройки.

Вы можете создать объект SharedFlow без использования shareIn. Например, вы можете использовать SharedFlow, чтобы отправлять сигналы остальным частям приложения, и тогда весь контент будет периодически обновляться одновременно. Помимо получения последних новостей, вы также можете обновить раздел с информацией о пользователе, добавив в него его любимые темы. В приведенном ниже фрагменте кода TickHandler предоставляет доступ к SharedFlow, чтобы другие классы знали, когда обновлять его контент. Как и в случае с StateFlow, чтобы отправлять объекты в поток, используйте в классе вспомогательное свойство типа MutableSharedFlow:

// Class that centralizes when the content of the app needs to be refreshed
class TickHandler(
    private val externalScope: CoroutineScope,
    private val tickIntervalMs: Long = 5000
) {
    // Backing property to avoid flow emissions from other classes
    private val _tickFlow = MutableSharedFlow<Unit>(replay = 0)
    val tickFlow: SharedFlow<Unit> = _tickFlow

    init {
        externalScope.launch {
            while (true) {
                _tickFlow.emit(Unit)
                delay(tickIntervalMs)
            }
        }
    }
}

class NewsRepository(
    // ...
    private val tickHandler: TickHandler,
    private val externalScope: CoroutineScope
) {
    init {
        externalScope.launch {
            // Listen for tick updates
            tickHandler.tickFlow.collect {
                refreshLatestNews()
            }
        }
    }

    suspend fun refreshLatestNews() { /* ... */ }
    // ...
}

Вы можете настроить поведение SharedFlow следующими способами:

  • replay позволяет повторно отправить ряд ранее переданных значений новым подписчикам.
  • onBufferOverflow позволяет задать правила для случаев, когда буфер заполнен объектами, которые нужно отправить. Значение по умолчанию – BufferOverflow.SUSPEND, при котором вызывающий объект приостанавливается. Другие варианты: DROP_LATEST или DROP_OLDEST.

У MutableSharedFlow также есть свойство subscriptionCount, в котором указано количество активных сборщиков. Это позволяет оптимизировать бизнес-логику. MutableSharedFlow также содержит функцию resetReplayCache, если вы не хотите воспроизводить последнюю информацию, отправленную в поток.

Дополнительные ресурсы для работы с потоками