Лекции

Senior

Kotlin Flow: основы

Холодные и горячие потоки Kotlin Flow, context, операторы, backpressure, sharing и UI state.

flowstateflowsharedflowbackpressurecoroutines
На этой странице

Это русскоязычная подготовленная версия предоставленного конспекта. Все исходные кодовые блоки сохранены как read-only примеры. Примеры не компилировались и не запускались в ходе подготовки.

Теория

Зачем Flow, если есть suspend

suspend:

один запрос → один результат
suspend fun getUser(): User

Flow:

0..N значений во времени
fun observeBookings(): Flow<List<Booking>>

Официальная документация: Документация Kotlin Coroutines.

Flow — асинхронный stream

val numbers: Flow<Int> = flow {
    emit(1)
    emit(2)
    emit(3)
}

numbers.collect { value ->
    println(value)
}

emit() производит значение, collect() потребляет.

Официальная документация: Flow API.

Suspend-based API

Концептуально:

interface Flow<out T> {
    suspend fun collect(
        collector: FlowCollector<T>
    )
}

interface FlowCollector<in T> {
    suspend fun emit(value: T)
}

Поэтому Flow естественно интегрирован с cancellation и CoroutineContext.

Официальная документация: Документация Kotlin Coroutines.

Cold Flow

val flow = flow {
    println("START")
    emit(1)
    emit(2)
}

Пока никто не делает collect, producer не запускается.

Каждый новый collector обычного cold Flow запускает upstream заново.

Официальная документация: Flow API.

Flow vs Sequence

Sequence<T> → lazy synchronous stream
Flow<T>     → lazy asynchronous suspend-aware stream

В Flow можно использовать suspend operations.

Официальная документация: Документация Kotlin Coroutines.

Context collector'а

Без специальных операторов upstream выполняется в coroutine context collector'а.

Отсюда необходимость flowOn.

flowOn

flow
    .flowOn(Dispatchers.Default)
    .collect {
        render(it)
    }

flowOn влияет на upstream.

flow
    .map { operationA(it) }
    .flowOn(Dispatchers.Default)
    .map { operationB(it) }
    .collect { ... }
Default:
flow
 ↓
map A
 ↓
flowOn
──────────────
collector context:
map B
 ↓
collect

Официальная документация: Dispatchers API.

Context preservation

Не следует произвольно делать:

flow {
    withContext(Dispatchers.IO) {
        emit(load())
    }
}

Для upstream context используется:

flow {
    emit(load())
}.flowOn(Dispatchers.IO)

Официальная документация: Dispatchers API.

Базовые operators

flow
    .map { it.toUiModel() }
    .filter { it.isActive }
    .collect { render(it) }

transform позволяет создать 0..N outputs на один input.

onStart

flow
    .map<Data, UiState> {
        UiState.Content(it)
    }
    .onStart {
        emit(UiState.Loading)
    }

catch

flow
    .catch { throwable ->
        emit(UiState.Error(throwable))
    }
    .collect { state ->
        render(state)
    }

catch ловит upstream exceptions относительно своей позиции в chain.

Cancellation

Flow построен на coroutines:

scope cancelled
      ↓
collect cancelled
      ↓
upstream cancelled

Официальная документация: Flow API.

Backpressure

Flow по умолчанию может suspend'ить producer, пока consumer обрабатывает значение.

emit(1)
   ↓
collector processes 1
   ↓
emit(2)

Официальная документация: Документация Kotlin Coroutines.

buffer

flow
    .buffer()
    .collect {
        slowOperation(it)
    }

Producer и consumer получают возможность работать более независимо.

conflate

Если нужны только более свежие значения, промежуточные можно пропускать.

flow
    .conflate()
    .collect {
        render(it)
    }

Хороший пример — progress.

collectLatest

flow.collectLatest { value ->
    process(value)
}

Новое значение отменяет обработку предыдущего.

A → process █████
              X
B →           process █████

mapLatest

queryFlow
    .mapLatest { query ->
        repository.search(query)
    }

При новом query предыдущая transformation отменяется.

Официальная документация: Flow API.

debounce

query
    .debounce(300)
    .distinctUntilChanged()
    .mapLatest { query ->
        repository.search(query)
    }

Подходит для search input, когда не нужно отправлять запрос на каждый символ.

distinctUntilChanged

A A A B B C
      ↓
A B C

StateFlow уже имеет equality-based conflation semantics.

Официальная документация: Flow API.

combine

combine(
    user,
    bookings,
) { user, bookings ->
    ScreenState(user, bookings)
}

Использует latest значения sources.

User:
   U1-------------U2---------

Bookings:
       B1----B2---------B3---

combine:
       U1B1--U1B2--U2B2--U2B3

zip

Сопоставляет emissions попарно.

A:
A1----A2----A3

B:
---B1----------B2----B3

zip:
---A1B1--------A2B2--A3B3

Ментальная модель:

combine = latest state
zip     = corresponding pairs

flatMapLatest

userId.flatMapLatest { id ->
    repository.observeUser(id)
}

При новом ID старая inner collection отменяется.

flatMapConcat

Inner flows выполняются последовательно.

A → █████
B       → █████
C             → █████

flatMapMerge

Inner flows могут выполняться concurrently.

A → █████████
B →   █████
C →      █████████

Итого:

flatMapConcat  → sequential
flatMapMerge   → concurrent
flatMapLatest  → cancel previous

Cold vs Hot

Cold:

Collector A → Producer A
Collector B → Producer B

Hot:

            Shared Producer
             /          \
Collector A              Collector B

flow {} обычно cold.

StateFlow и SharedFlow — hot.

Официальная документация: Flow API.

StateFlow

private val _state =
    MutableStateFlow(UiState())

val state: StateFlow<UiState> =
    _state.asStateFlow()

StateFlow:

  • всегда имеет current value;
  • требует initial value;
  • новый collector сразу получает current value;
  • имеет equality-based conflation;
  • хорошо подходит для state.

Официальная документация: Flow API.

Atomic update

_state.update {
    it.copy(isLoading = true)
}

Предпочтительно для read-modify-write обновлений при возможной concurrency.

SharedFlow

val events =
    MutableSharedFlow<Event>()

Более общий hot broadcast stream без обязательной semantics «текущего состояния».

Официальная документация: Flow API.

replay

MutableSharedFlow<Event>(
    replay = 2
)

Новый subscriber получает последние два replayed значения.

Но SharedFlow(replay = 1) всё равно не полностью равен StateFlow.

Официальная документация: Flow API.

State vs Event

State:

isLoading
user
bookings
selectedDate

Часто → StateFlow.

Event:

ShowSnackbar
NavigateBack
OpenUrl

Может быть SharedFlow/Channel, но сначала нужно определить delivery semantics.

Официальная документация: Flow API.

stateIn

val bookings =
    repository.observeBookings()
        .stateIn(
            scope = viewModelScope,
            started = SharingStarted.WhileSubscribed(5_000),
            initialValue = emptyList(),
        )
Cold Flow
   ↓
stateIn
   ↓
StateFlow

Официальная документация: Flow API.

shareIn

val shared =
    repository.events()
        .shareIn(
            scope = viewModelScope,
            started = SharingStarted.WhileSubscribed(),
            replay = 0,
        )
Cold Flow
   ↓
shareIn
   ↓
SharedFlow

Официальная документация: Flow API.

Зачем sharing

Без sharing два collectors cold Flow могут запустить дорогой upstream дважды.

После sharing:

             one upstream
                  │
          ┌───────┴───────┐
          ▼               ▼
     Collector A      Collector B

Официальная документация: Flow API.

SharingStarted

Частые варианты:

SharingStarted.Eagerly
SharingStarted.Lazily
SharingStarted.WhileSubscribed(...)

Для UI часто используется WhileSubscribed(timeout).

Timeout помогает не перезапускать upstream мгновенно при коротком configuration-change gap.

ViewModel pipeline

data class VenueUiState(
    val venue: Venue? = null,
    val bookings: List<Booking> = emptyList(),
    val selectedDate: LocalDate = LocalDate.now(),
)

private val selectedDate =
    MutableStateFlow(LocalDate.now())

val state: StateFlow<VenueUiState> =
    combine(
        repository.observeVenue(),
        repository.observeBookings(),
        selectedDate,
    ) { venue, bookings, date ->
        VenueUiState(
            venue = venue,
            bookings = bookings,
            selectedDate = date,
        )
    }.stateIn(
        scope = viewModelScope,
        started = SharingStarted.WhileSubscribed(5_000),
        initialValue = VenueUiState(),
    )

Официальная документация: Flow API.

Ментальная карта Flow

                  Flow
                    │
          ┌─────────┴─────────┐
          │                   │
        Cold                 Hot
          │                   │
      flow { }         ┌──────┴──────┐
                       │             │
                  StateFlow     SharedFlow
                       │             │
                     State        broadcast
Cold Flow ──stateIn──► StateFlow
Cold Flow ──shareIn──► SharedFlow

Официальная документация: Flow API.

Практика

Вопросы для разбора

Для каждого примера нарисуйте иерархию Job, определите CoroutineContext и Dispatcher, отметьте точки suspension/cancellation, затем проследите failure вверх и cancellation вниз. Для Flow дополнительно определите cold/hot semantics, владельца collection, поведение медленного consumer, replay и последствия lifecycle stop/start.

Проверка архитектурного решения

Спросите: кто владеет этой работой, кто может её отменить, где обрабатывается ошибка, нужна ли каждая emission или только latest, и переживает ли результат пересоздание UI. Coroutine не переживает process death; для гарантированной фоновой работы нужен durable scheduler, например WorkManager.

Дополнительные материалы