여러 값을 시간에 걸쳐 흘려보내는 코루틴 기반 스트림.
fun observeUsers(): Flow<List<User>> = flow {
while (true) {
emit(api.getUsers()) // 값을 내보낸다
delay(5000)
}
}
viewModelScope.launch {
observeUsers().collect { users -> _state.value = users } // 수집해야 흐른다
}
cold — 수집할 때 비로소 시작한다
Flow 를 만드는 것만으로는 아무 일도 안 일어난다
- collect 하는 순간 블록이 실행된다
- 수집자가 둘이면 블록이 두 번 독립적으로 실행된다
"레시피를 만드는 것" 과 "요리하는 것" 의 차이
연산자
flow.map { it.toUiModel() }
.filter { it.isVisible }
.onEach { Log.d("TAG", "$it") }
.catch { emit(emptyList()) } // 업스트림 예외를 잡는다
.onStart { _loading.value = true }
.onCompletion { _loading.value = false }
.debounce(300) // 입력이 멈춘 뒤에만
.distinctUntilChanged() // 같은 값 연속이면 무시
.flowOn(Dispatchers.IO) // 업스트림을 IO 에서 실행
.collect { ... }
flowOn은 자기보다 위(업스트림)에만 적용된다 — 아래는 수집하는 쪽의 컨텍스트를 따른다. 순서가 의미를 바꾼다.
검색어 입력 처리 — 전형적 조합
queryFlow
.debounce(300) // 타이핑이 멈추길 기다린다
.distinctUntilChanged() // 같은 검색어면 다시 안 부른다
.flatMapLatest { q -> repository.search(q) } // 새 검색이 오면 이전 것을 취소
.collect { _results.value = it }
flatMapLatest가 이전 요청을 자동 취소한다 — 직접 관리하면 복잡한 경합 문제를 한 줄로 푼다.
백프레셔
flow.buffer() // 생산과 소비를 분리해 병렬로
.conflate() // 처리 중 들어온 값은 최신 것만 남긴다
.collectLatest { ... } // 새 값이 오면 이전 처리를 취소한다
취소와 예외
// try-catch 대신 catch 연산자를 쓴다 (취소를 안 삼킨다)
flow.catch { e -> emit(fallback) }.collect { ... }
// collect 안에서 던진 예외는 catch 가 안 잡는다 (다운스트림이므로)
면접 함정
- ❌ "Flow는 RxJava와 같다" → 코루틴 기반이라 구조적 동시성·취소가 자연스럽게 통합된다.
- ❌ "flowOn을 마지막에 붙이면 전체에 적용된다" → 위쪽에만 적용된다.
만드는 여러 방법
flowOf(1, 2, 3) // 고정 값
listOf(1,2,3).asFlow() // 컬렉션에서
flow { emit(api.get()) } // 빌더 (suspend 함수 호출 가능)
// 콜백 기반 API 를 Flow 로 감싼다
callbackFlow {
val listener = object : Listener {
override fun onEvent(e: Event) { trySend(e) }
}
api.register(listener)
awaitClose { api.unregister(listener) } // 반드시 정리한다
}
awaitClose를 빠뜨리면 수집이 끝나도 리스너가 남아 누수가 된다.
Room·Retrofit과의 조합
@Dao interface UserDao {
@Query("SELECT * FROM user")
fun observeAll(): Flow<List<User>> // DB 가 바뀌면 자동으로 새 값이 흐른다
}
// 여러 Flow 를 합친다
combine(userFlow, settingsFlow) { user, settings -> UiState(user, settings) }
.collect { ... }
// 순서대로 이어 붙인다
flowA.onCompletion { emitAll(flowB) }
테스트
@Test fun test() = runTest {
val values = flow.take(3).toList() // 유한하게 잘라 수집한다
assertEquals(listOf(1,2,3), values)
}
// 무한 Flow 를 그냥 toList() 하면 영영 안 끝난다 — take · first 로 자른다