Kotlin 코루틴 시리즈 (4/5)

이전 편: 3 구조화된 동시성과 예외 처리

다음 편: 5 코루틴 테스트

suspend 함수는 하나의 값을 비동기로 반환한다. 하지만 실무에서는 "값이 시간에 걸쳐 여러 개 흘러오는" 상황이 훨씬 많다. 실시간 검색어 입력, WebSocket 메시지, 센서 데이터, DB 변경 감지 — 모두 스트림이다.

Java에서는 이런 문제를 RxJava의 Observable이나 Reactor의 Flux로 해결했다. 하지만 이 라이브러리들은 학습 곡선이 가파르고, 코루틴과 잘 맞지 않는다. Kotlin은 코루틴 위에 Flow라는 비동기 스트림 API를 제공한다. suspend 함수가 "비동기 단일 값"이라면, Flow는 "비동기 다중 값"이다.

Flow란

Flow는 순차적으로 값을 방출(emit)하는 비동기 스트림이다. Kotlin의 Sequence가 동기적으로 여러 값을 생산하듯, Flow는 비동기적으로 여러 값을 생산한다.

fun numbers(): Flow<Int> = flow {
    emit(1)
    delay(100)
    emit(2)
    delay(100)
    emit(3)
}
  • flow { } — Flow를 생성하는 빌더. 블록 안에서 emit()으로 값을 방출한다.
  • emit() — 값을 하나 내보낸다. suspend 함수이므로 중단될 수 있다.
  • delay() — Flow 안에서 suspend 함수를 자유롭게 호출할 수 있다.

이 코드만으로는 아무 일도 일어나지 않는다. Flow는 수집(collect)할 때까지 실행되지 않는다. 이것이 Flow의 가장 중요한 특성이다.

numbers().collect { value ->
    println(value)  // 1, 2, 3이 순차적으로 출력
}

비유하자면 Flow는 수도꼭지다. 꼭지를 열어야(collect) 물이 나오고, 꼭지를 잠그면(취소) 물이 멈춘다. 꼭지를 열기 전에는 물이 흐르지 않는다.


Cold Stream vs Hot Stream

스트림에는 두 종류가 있다. 이 구분을 이해하는 것이 Flow를 제대로 쓰는 출발점이다.

Cold Stream — 구독할 때 시작

val coldFlow = flow {
    println("Flow 시작!")
    emit(1)
    emit(2)
}

// 아직 "Flow 시작!"은 출력되지 않음

coldFlow.collect { println(it) }  // 여기서야 "Flow 시작!" 출력
coldFlow.collect { println(it) }  // 다시 처음부터 "Flow 시작!" 출력

Cold Stream은 구독자(collector)가 있을 때만 실행된다. 각 구독자가 독립적인 스트림을 받는다. 유튜브의 녹화 영상과 같다 — 누가 재생하든 처음부터 시작하고, 각자 독립적으로 본다.

flow { } 빌더로 만든 기본 Flow가 Cold Stream이다.

Hot Stream — 구독 여부와 관계없이 데이터가 흐름

Hot Stream은 구독자가 없어도 데이터가 생산된다. 유튜브의 라이브 방송과 같다 — 시청자가 없어도 방송은 계속되고, 중간에 들어온 시청자는 그 시점부터 본다.

StateFlowSharedFlow가 Hot Stream이다. 이건 뒤에서 자세히 다룬다.

graph LR subgraph "Cold Stream (Flow)" C1[Collector 1] -->|구독| P1[Producer] C2[Collector 2] -->|구독| P2[Producer] end subgraph "Hot Stream (SharedFlow)" P3[Producer] -->|방출| C3[Collector 1] P3 -->|방출| C4[Collector 2] end

Cold Stream에서는 각 Collector가 자신만의 Producer를 가진다. Hot Stream에서는 하나의 Producer가 여러 Collector에게 같은 데이터를 보낸다.


Flow 빌더

Flow를 생성하는 방법은 여러 가지다.

flow { } — 기본 빌더

가장 유연한 빌더다. 블록 안에서 emit()을 호출해 값을 방출한다.

fun fetchUsers(): Flow<User> = flow {
    val users = userRepository.findAll()  // suspend 함수
    users.forEach { user ->
        emit(user)
    }
}

flow 빌더 안에서는 suspend 함수를 자유롭게 호출할 수 있다. 네트워크 요청, DB 쿼리 등을 하면서 결과를 하나씩 방출하는 패턴이 일반적이다.


flowOf() — 고정된 값들

val numbers = flowOf(1, 2, 3, 4, 5)

이미 알고 있는 값들로 Flow를 만들 때 사용한다. listOf()의 Flow 버전이라고 생각하면 된다.


asFlow() — 컬렉션을 Flow로 변환

val listFlow = listOf(1, 2, 3).asFlow()
val rangeFlow = (1..10).asFlow()

기존 컬렉션이나 범위를 Flow로 변환한다. 동기 데이터를 비동기 파이프라인에 넣을 때 유용하다.


channelFlow — 여러 코루틴에서 방출

flow { } 빌더에는 제약이 하나 있다. 같은 코루틴에서만 emit()을 호출할 수 있다. 다른 코루틴에서 emit()을 호출하면 에러가 난다.

// 컴파일은 되지만 런타임 에러!
flow {
    launch {
        emit(1)  // 다른 코루틴에서 emit → 에러
    }
}

여러 코루틴에서 값을 방출해야 하면 channelFlow를 쓴다.

fun fetchAllData(): Flow<Data> = channelFlow {
    launch { send(fetchFromApi1()) }   // 병렬 요청 1
    launch { send(fetchFromApi2()) }   // 병렬 요청 2
}
  • channelFlow에서는 emit() 대신 send()를 사용한다.
  • 내부적으로 Channel을 사용하므로 여러 코루틴에서 안전하게 값을 보낼 수 있다.

중간 연산자

Flow의 진짜 강점은 연산자 체이닝이다. 컬렉션의 map, filter와 같은 연산을 비동기 스트림에서 할 수 있다.

중간 연산자(intermediate operator)는 Flow를 받아서 새로운 Flow를 반환한다. 중간 연산자 자체는 suspend 함수가 아니며, 체이닝만 정의할 뿐 값을 소비하지 않는다.

map — 변환

flowOf(1, 2, 3)
    .map { it * 2 }
    .collect { println(it) }  // 2, 4, 6

각 값을 변환한다. 컬렉션의 map과 동일하지만, 차이점이 있다. Flow의 map 안에서는 suspend 함수를 호출할 수 있다.

userIds.asFlow()
    .map { id -> userRepository.findById(id) }  // suspend 함수 호출 가능
    .collect { user -> println(user.name) }

컬렉션의 map에서는 이렇게 할 수 없다. 비동기 작업을 포함하는 변환이 자연스럽게 되는 것이 Flow의 장점이다.


filter, take, onEach 등

filter는 조건에 맞는 값만 통과시킨다. map과 마찬가지로 람다 안에서 suspend 함수를 호출할 수 있다.

flowOf(1, 2, 3, 4, 5)
    .filter { it % 2 == 0 }
    .collect { println(it) }  // 2, 4

take는 지정한 개수만큼만 값을 받고 Flow를 취소한다. 무한 스트림에서 유용하다. onEach는 값을 변경하지 않으면서 부수 효과(로깅, 추적 등)를 끼워넣는다.

이 연산자들의 공통점은 원본 Flow를 변경하지 않고 새 Flow를 반환한다는 것이다. 컬렉션의 함수형 연산과 동일한 패턴이지만, 비동기 스트림에서 동작한다는 점이 다르다.


transform — 유연한 변환

map은 1:1 변환이다. 하나의 입력에 하나의 출력. 하지만 실무에서는 "하나의 입력으로 여러 값을 방출"하거나 "조건에 따라 방출하지 않는" 경우가 있다. transform은 이 제약이 없어서, 하나의 입력에 0개, 1개, 또는 여러 개의 값을 방출할 수 있다.

flowOf("Kotlin", "Java", "Go")
    .transform { language ->
        emit("$language 시작")
        val result = fetchTutorial(language)  // suspend
        emit("$language 완료: $result")
    }
    .collect { println(it) }
  • 하나의 입력("Kotlin")에 두 개의 값("Kotlin 시작", "Kotlin 완료: ...")을 방출한다.
  • mapfilter로 표현하기 어려운 복잡한 변환에 적합하다.

종단 연산자

종단 연산자(terminal operator)는 Flow를 실제로 실행시키는 suspend 함수다. collect가 대표적이다.

val flow = flowOf(1, 2, 3)

// collect — 모든 값을 소비
flow.collect { println(it) }

// toList — Flow를 리스트로 변환
val list: List<Int> = flow.toList()

// first — 첫 번째 값만 받고 취소
val first: Int = flow.first()

// reduce — 누적 연산
val sum: Int = flow.reduce { acc, value -> acc + value }
  • collect — 모든 값을 하나씩 소비한다. 가장 많이 쓰는 종단 연산자.
  • toList() / toSet() — Flow의 모든 값을 컬렉션으로 모은다.
  • first() — 첫 번째 값만 받고 나머지는 취소한다.
  • reduce() / fold() — 값을 누적해서 하나의 결과를 만든다.

종단 연산자가 호출되어야 Flow가 실행된다는 점을 잊지 말자. 중간 연산자만 체이닝하면 아무 일도 일어나지 않는다.


Flow의 Context

Flow는 기본적으로 collect를 호출한 코루틴의 Context에서 실행된다. 이것을 Context 보존(context preservation)이라고 한다.

fun numbers(): Flow<Int> = flow {
    println("Flow: ${Thread.currentThread().name}")  // Main
    emit(1)
}

withContext(Dispatchers.Main) {
    numbers().collect { value ->
        println("Collect: ${Thread.currentThread().name}")  // Main
    }
}

Flow 빌더 안의 코드도, collect의 람다도 모두 같은 Context(여기서는 Main)에서 실행된다.

하지만 실무에서는 "데이터는 IO 스레드에서 가져오고, UI는 Main 스레드에서 업데이트" 같은 패턴이 필요하다. 이때 flowOn을 쓴다.

fun fetchUsers(): Flow<User> = flow {
    // 이 블록은 Dispatchers.IO에서 실행됨
    val users = repository.findAll()
    users.forEach { emit(it) }
}
.flowOn(Dispatchers.IO)  // 위쪽(upstream)의 Context를 변경
  • flowOn자기보다 위쪽(upstream)의 실행 Context를 변경한다.
  • 아래쪽(downstream, collect 쪽)의 Context는 변경하지 않는다.
graph LR A["flow { emit() }"] -->|"flowOn(IO)"| B["map { }"] -->|"flowOn(Default)"| C["collect { }"] A -.- D["Dispatchers.IO"] B -.- E["Dispatchers.Default"] C -.- F["Collector의 Context"]
flow 안에서 withContext 금지

flow { } 빌더 안에서 withContext로 Context를 바꾸면 안 된다. Flow의 Context 보존 원칙을 위반하기 때문에 런타임 에러가 발생한다. 반드시 flowOn을 사용해야 한다.


StateFlow와 SharedFlow

flow { } 빌더로 만든 Flow는 Cold Stream이다. 구독할 때마다 처음부터 실행된다. 하지만 실무에서는 "현재 상태를 유지하면서 여러 구독자에게 공유하는" Hot Stream이 필요할 때가 많다.

StateFlow — 상태를 가진 Flow

StateFlow항상 하나의 현재 값을 가지고 있는 Hot Stream이다. 새 구독자가 들어오면 가장 최근 값을 즉시 받는다.

class UserViewModel {
    private val _uiState = MutableStateFlow(UiState.Loading)
    val uiState: StateFlow<UiState> = _uiState.asStateFlow()

    fun loadUser() {
        viewModelScope.launch {
            val user = repository.fetchUser()
            _uiState.value = UiState.Success(user)  // 상태 업데이트
        }
    }
}
  • MutableStateFlow — 값을 변경할 수 있는 StateFlow. 내부(ViewModel 등)에서 사용.
  • asStateFlow() — 읽기 전용 StateFlow로 변환. 외부에 노출할 때 사용.
  • .value — 현재 값을 즉시 읽거나 변경할 수 있다.

StateFlow의 핵심 특성은 세 가지다.

  1. 항상 값이 있다 — 초기값이 필수이므로, null 체크 없이 .value로 접근할 수 있다.
  2. 중복 방출 제거 — 같은 값을 연속으로 설정하면 구독자에게 알리지 않는다. distinctUntilChanged가 내장되어 있다.
  3. 마지막 값만 유지 — 버퍼가 1개다. 이전 값은 사라진다.

Android의 LiveData와 비슷하지만, 코루틴 기반이고 플랫폼 독립적이라는 차이가 있다.


SharedFlow — 이벤트를 위한 Flow

StateFlow는 "상태"에 적합하다. 하지만 일회성 이벤트(토스트 메시지, 네비게이션, 에러 알림 등)에는 맞지 않다. 같은 이벤트를 두 번 보내면 StateFlow는 중복으로 판단해서 무시하기 때문이다.

SharedFlow는 이런 이벤트 처리를 위한 Hot Stream이다.

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

    suspend fun emit(event: Event) {
        _events.emit(event)
    }
}

SharedFlow는 StateFlow와 다른 점이 있다.

  • 초기값이 없다 — 구독 시점 이후의 이벤트만 받는다.
  • 중복 값도 방출한다 — 같은 이벤트를 여러 번 보낼 수 있다.
  • replay 설정 가능MutableSharedFlow(replay = 3)으로 최근 3개 값을 새 구독자에게 재생할 수 있다.
StateFlowSharedFlow
초기값필수없음
중복 방출무시허용
용도UI 상태일회성 이벤트
replay1 (고정)설정 가능

Flow의 에러 처리

Flow에서 예외가 발생하면 Flow가 종료된다. 예외 처리 방법은 여러 가지가 있다.

catch — 선언적 에러 처리

flow {
    emit(1)
    throw RuntimeException("에러 발생!")
    emit(2)
}
.catch { e ->
    println("잡힘: ${e.message}")
    emit(-1)  // 대체 값 방출 가능
}
.collect { println(it) }
// 출력: 1, 잡힘: 에러 발생!, -1
  • catch자기보다 위쪽(upstream)에서 발생한 예외만 잡는다.
  • catch 블록 안에서 emit()으로 대체 값을 방출할 수 있다.
  • catch 아래쪽(downstream)의 예외는 잡지 않는다.

이 점이 중요하다. collect 안에서 발생한 예외는 catch로 잡을 수 없다.

flow { emit(1) }
    .catch { println("여기서 잡히지 않음") }
    .collect { throw RuntimeException("collect에서 에러") }  // catch를 통과함!

collect는 Flow 파이프라인의 가장 끝(downstream)에 있으므로, 그 위에 선언된 catch로는 잡을 수 없다. 이 문제를 해결하려면 값 처리 로직을 collect에서 빼서 catch보다 위쪽인 중간 연산자로 옮겨야 한다. onEach가 이 역할에 적합하다.

flow { emit(1) }
    .onEach { value ->
        // 값 처리 로직을 여기로 옮김
        process(value)  // 여기서 예외 발생하면 catch에서 잡힘
    }
    .catch { e -> println("잡힘: ${e.message}") }
    .collect()  // 빈 collect

onEach는 중간 연산자이므로 catch보다 위쪽(upstream)이다. 따라서 onEach 안에서 발생한 예외는 catch에서 잡힌다. 실무에서는 이 패턴이 자주 쓰이므로, "collect에서 로직을 처리하지 말고, onEach로 옮긴 뒤 catch로 감싸라"는 관용구를 기억해두면 좋다.


onCompletion — 완료/에러 감지

flowOf(1, 2, 3)
    .onCompletion { cause ->
        if (cause != null) {
            println("에러로 종료: ${cause.message}")
        } else {
            println("정상 완료")
        }
    }
    .collect { println(it) }

onCompletion은 Flow가 완료되었을 때(정상이든 에러든) 호출된다. Java의 finally와 비슷한 역할이다. 다만 예외를 "처리"하지는 않는다. 예외는 여전히 하류로 전파된다.


실전 패턴

debounce — 검색어 입력 최적화

사용자가 타이핑할 때마다 API를 호출하면 낭비다. 입력이 멈춘 후 일정 시간이 지나면 그때 호출하는 것이 좋다.

searchQueryFlow
    .debounce(300)              // 300ms 동안 새 입력 없으면 통과
    .filter { it.isNotBlank() } // 빈 문자열 제외
    .distinctUntilChanged()     // 같은 쿼리 반복 방지
    .flatMapLatest { query ->   // 이전 검색 취소, 최신 쿼리만 실행
        searchApi(query)
    }
    .collect { results ->
        updateUi(results)
    }
  • debounce(300) — 300ms 동안 새 값이 들어오지 않아야 통과시킨다.
  • distinctUntilChanged() — 이전 값과 같으면 무시한다.
  • flatMapLatest — 새 값이 들어오면 이전 Flow를 취소하고 새 Flow를 시작한다.

이 패턴은 검색, 자동완성, 필터링 등에서 거의 표준처럼 사용된다.


combine — 여러 Flow 합치기

val searchQuery = MutableStateFlow("")
val sortOrder = MutableStateFlow(SortOrder.NEWEST)

combine(searchQuery, sortOrder) { query, sort ->
    repository.search(query, sort)
}
.collect { results ->
    updateUi(results)
}

combine은 여러 Flow 중 어느 하나라도 새 값을 방출하면 가장 최신 값들을 조합해서 방출한다. 검색어가 바뀌어도, 정렬이 바뀌어도 결과가 갱신된다. UI 상태를 여러 소스에서 조합할 때 자주 쓴다.


심화 분석

Flow와 Sequence의 차이

Kotlin의 Sequence도 지연 평가(lazy evaluation)를 지원한다. Flow와 비슷해 보이지만, 핵심적인 차이가 있다.

Sequence동기적이다. 값을 생산하는 동안 스레드를 점유한다. yield()로 값을 내보내지만, suspend 함수를 호출할 수 없다.

Flow비동기적이다. 값을 생산하는 동안 스레드를 반환할 수 있다. emit()으로 값을 내보내고, suspend 함수를 자유롭게 호출할 수 있다.

정리하면 이렇다.

  • 데이터가 이미 메모리에 있고, 변환만 하면 되는 경우 → Sequence
  • 네트워크, DB 등 비동기 작업이 포함된 스트림 → Flow

stateIn과 shareIn

Cold Flow를 Hot Flow로 변환하는 함수다. 하나의 Flow를 여러 곳에서 구독할 때, 매번 처음부터 실행되는 것을 방지한다.

val userFlow: StateFlow<User> = repository.observeUser()
    .stateIn(
        scope = viewModelScope,
        started = SharingStarted.WhileSubscribed(5000),
        initialValue = User.Empty
    )
  • stateIn — Cold Flow를 StateFlow로 변환한다. initialValue가 필수다.
  • shareIn — Cold Flow를 SharedFlow로 변환한다.
  • SharingStarted.WhileSubscribed(5000) — 마지막 구독자가 사라진 후 5초 동안 공유를 유지한다. 화면 회전 등 짧은 구독 해제에도 Flow를 재시작하지 않으려는 설정이다.

buffer와 conflate

기본적으로 Flow는 순차 처리다. 방출과 수집이 번갈아 실행된다. 방출이 빠르고 수집이 느리면 병목이 생긴다.

flow {
    repeat(3) {
        delay(100)   // 100ms마다 생산
        emit(it)
    }
}
.buffer()  // 방출과 수집을 별도 코루틴에서 실행
.collect {
    delay(300)  // 300ms 처리
    println(it)
}
  • buffer() — 방출과 수집을 병렬로 실행한다. 방출 측이 수집을 기다리지 않고 버퍼에 쌓는다.
  • conflate() — 수집이 느릴 때 중간 값을 건너뛴다. 가장 최신 값만 수집한다.

자주 하는 실수

collect를 빼먹기

// 아무 일도 일어나지 않는다!
flowOf(1, 2, 3)
    .map { it * 2 }
    .filter { it > 2 }
// collect가 없으면 Flow는 실행되지 않음

Flow는 Cold Stream이다. 종단 연산자(collect, toList, first 등)가 없으면 어떤 코드도 실행되지 않는다. 컬렉션 체이닝과 헷갈리기 쉬운 부분이다.

flow 안에서 withContext 사용

// 잘못된 사용 — 런타임 에러
flow {
    withContext(Dispatchers.IO) {
        emit(fetchData())
    }
}

// 올바른 사용 — flowOn으로 Context 전환
flow {
    emit(fetchData())
}
.flowOn(Dispatchers.IO)

flow { } 빌더는 Context 보존 원칙을 강제한다. 내부에서 withContext로 Context를 바꾸면 IllegalStateException이 발생한다. Context를 바꾸려면 반드시 flowOn을 사용해야 한다.

StateFlow에서 이벤트 처리

// 나쁜 예 — 같은 에러 메시지가 두 번 발생하면 두 번째는 무시됨
private val _error = MutableStateFlow<String?>(null)

fun onError(message: String) {
    _error.value = message  // 같은 메시지면 구독자에게 전달 안 됨
}

StateFlowdistinctUntilChanged가 내장되어 있어서, 같은 값을 연속으로 설정하면 구독자에게 알리지 않는다. 일회성 이벤트(에러 메시지, 네비게이션 등)에는 SharedFlow를 사용해야 한다.