코틀린 언어 용어 사전
Flowcold 스트림 · collect · flowOn

Flow

코루틴 위의 비동기 스트림. 수집을 시작해야 흐르는 cold 스트림이다.

여러 값을 시간에 걸쳐 흘려보내는 코루틴 기반 스트림.

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 로 자른다

함께 보면 좋은 용어

노트에서 맥락과 함께 보기 — Flow·StateFlow·SharedFlow — 비동기 스트림