Kotlin Coroutine: Channel vs Flow, 언제 쓰나 | 비교 실전 가이드

이 글의 핵심

Flow를 collect할 때마다 업스트림 로직이 처음부터 다시 실행되어 같은 작업이 중복되는 문제는 콜드 스트림 특성을 모를 때 자주 생깁니다. 워커 결과를 한 곳으로 모으는 Channel과 선언적 변환에 강한 Flow를 실무 시나리오로 비교하고, 버퍼가 가득 찼을 때의 동작과 트러블슈팅 포인트를 정리합니다.

들어가며

Kotlin에서 여러 값을 시간에 따라 다룰 때 후보는 크게 둘입니다. Channel은 보통 핫(hot)에 가깝으며, 생산자가 소비자와 독립적으로(또는 강하게 결합해) 이벤트를 밀어 넣습니다. Flow는 기본적으로 콜드(cold)이며, 수집(collect)이 시작될 때 업스트림이 실행됩니다. 이 글은 “둘 다 스트림 같은데 뭐가 다르냐”는 질문에 백프레셔·소유권·테스트 관점에서 답을 정리합니다. Kotlin 코루틴에서 채널과 플로우의 차이를 핫/콜드와 수집 시점으로 나누면 선택이 단순해집니다. 코루틴과 스레드의 큰 그림은 코루틴 vs 스레드, 기본기는 코루틴 가이드와 함께 보면 좋습니다.


개념 설명

  • Channel: send/receive로 동시에 실행 중인 생산자·소비자를 연결합니다. 버퍼가 없으면 한쪽이 준비될 때까지 서로를 맞춥니다(동기화 채널). “이미 돌아가는 파이프라인”을 표현하기 좋습니다.
  • Flow: 중단 가능한 일련의 값을 표현합니다. collect가 호출되기 전까지는 업스트림 로직이 필요할 때만 돌아갑니다(콜드). “같은 소스를 여러 번 구독할 수 있는” 리액티브 시퀀스에 가깝습니다.
  • SharedFlow / StateFlow: Flow 계열이지만 핫에 가깝게 동작합니다. “Channel vs Flow” 질문은 종종 콜드 Flow vs SharedFlow까지 확장됩니다.

“핫”과 “콜드”를 가장 실용적으로 구분하는 질문은 “아무도 수집하지 않을 때 생산 코드가 도는가?” 입니다. flow { } 블록은 collect가 호출되기 전에는 한 줄도 실행되지 않고, 수집자가 둘이면 두 번 실행됩니다. Channel은 이미 실행 중인 코루틴이 send하는 통로이므로 수집자가 몇이든 생산은 한 번이고, 대신 각 값은 수신자 중 한 명에게만 전달됩니다. 이 “한 명에게만” 특성이 Channel을 작업 큐로, SharedFlow를 방송으로 구분하는 기준입니다. 워커 세 개가 같은 Channel에서 receive하면 일을 나눠 가지고, 세 개가 같은 SharedFlow를 collect하면 모두 같은 값을 받습니다.


실전 구현 (단계별 코드)

의존성(2026년 기준 Kotlin 2.x + Coroutines 1.10+)

// build.gradle.kts
dependencies {
    implementation("org.jetbrains.kotlinx:kotlinx-coroutines-core:1.10.2")
}

Channel: 워커가 결과를 모아 전달

기본 Channel 사용

import kotlinx.coroutines.*
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
    val ch = Channel<Int>(capacity = 4)
    val producer = launch {
        repeat(5) { i ->
            println("Sending $i")
            ch.send(i)  // 버퍼 가득 차면 중단
        }
        ch.close()
        println("Producer done")
    }
    val consumer = launch {
        for (x in ch) {
            println("Received $x")
            delay(100)  // 느린 소비자 시뮬레이션
        }
        println("Consumer done")
    }
    producer.join()
    consumer.join()
}

출력:

Sending 0
Sending 1
Sending 2
Sending 3
Sending 4
Received 0
Producer done
Received 1
Received 2
Received 3
Received 4
Consumer done

출력 순서를 따라가 보면 버퍼의 역할이 보입니다. runBlocking은 한 스레드에서 코루틴을 번갈아 실행하므로 먼저 시작한 생산자가 0~3을 버퍼(용량 4)에 채우고, Sending 4에서 버퍼가 가득 차 send가 중단(suspend) 됩니다. 스레드를 막는 것이 아니라 코루틴만 멈추는 것이라 그 사이 소비자가 실행되어 0을 꺼내고, 빈자리가 생기자 생산자가 깨어나 4를 넣고 채널을 닫습니다. for (x in ch)는 채널이 닫히고 버퍼가 모두 비워진 뒤에 끝나므로, close()를 호출해도 이미 들어간 값은 유실되지 않습니다. 반대로 생산자가 close()를 빠뜨리면 소비자의 for 루프가 영원히 다음 값을 기다려 consumer.join()이 끝나지 않습니다. 채널 코드에서 프로그램이 멈춘다면 가장 먼저 확인할 부분입니다.

무버퍼 Channel (Rendezvous)

fun main() = runBlocking {
    val ch = Channel<Int>()  // capacity = 0 (기본값)
    val producer = launch {
        repeat(3) { i ->
            println("Sending $i")
            ch.send(i)  // 소비자가 받을 때까지 중단
            println("Sent $i")
        }
        ch.close()
    }
    val consumer = launch {
        delay(200)  // 소비자가 늦게 시작
        for (x in ch) {
            println("Received $x")
            delay(100)
        }
    }
    producer.join()
    consumer.join()
}

출력:

Sending 0
(200ms 대기)
Received 0
Sent 0
Sending 1
(100ms 대기)
Received 1
Sent 1
...

다중 생산자/소비자

fun main() = runBlocking {
    val ch = Channel<Int>(capacity = 10)
    // 생산자 3개
    repeat(3) { producerId ->
        launch {
            repeat(5) { i ->
                ch.send(producerId * 100 + i)
            }
        }
    }
    // 소비자 2개
    repeat(2) { consumerId ->
        launch {
            for (x in ch) {
                println("Consumer $consumerId received $x")
            }
        }
    }
    delay(500)  // 모든 작업 완료 대기
    ch.close()
}

이 예제의 delay(500) 후 close()는 개념 설명용이고 실무에서는 쓰면 안 되는 패턴입니다. 생산자가 500ms 안에 끝난다는 보장이 없으므로, 느린 환경에서는 아직 보내는 중인 생산자가 ClosedSendChannelException으로 실패합니다. 생산자들의 Job을 모아 joinAll()로 모두 끝난 것을 확인한 뒤 닫는 것이 정확합니다. 여러 소비자가 for (x in ch)로 같은 채널을 읽으면 값은 먼저 요청한 소비자에게 하나씩 분배되므로, 어떤 소비자가 어떤 값을 받을지는 정해져 있지 않습니다. 순서가 중요한 처리라면 소비자를 하나로 두거나 키별로 채널을 나눠야 합니다.

채널 생산자는 produce 빌더로 쓰는 편이 안전합니다. val ch = produce { repeat(5) { send(it) } }처럼 쓰면 블록이 끝날 때 채널이 자동으로 닫히고, 블록이 예외로 끝나면 그 예외로 채널이 닫혀 소비자 쪽에 전달됩니다. close() 누락 버그가 원천적으로 없어집니다.

Flow: 수집할 때마다(또는 각 collect마다) 로직 실행

기본 Flow (콜드)

import kotlinx.coroutines.flow.*
import kotlinx.coroutines.*
fun numbers(): Flow<Int> = flow {
    println("Flow started")
    repeat(5) { i ->
        emit(i)
        delay(100)
    }
    println("Flow completed")
}
fun main() = runBlocking {
    println("=== First collect ===")
    numbers().collect { println("a: $it") }
    
    println("\n=== Second collect ===")
    numbers().collect { println("b: $it") }  // 콜드: 다시 실행됨
}

출력:

=== First collect ===
Flow started
a: 0
a: 1
a: 2
a: 3
a: 4
Flow completed
=== Second collect ===
Flow started
b: 0
b: 1
b: 2
b: 3
b: 4
Flow completed

Flow 연산자 체이닝

fun main() = runBlocking {
    flow {
        repeat(10) { emit(it) }
    }
    .filter { it % 2 == 0 }  // 짝수만
    .map { it * it }         // 제곱
    .take(3)                 // 처음 3개
    .collect { println(it) }
}

출력:

0
4
16

연산자 체인에서 알아 둘 점은 값이 하나씩 끝까지 흘러간다는 것입니다. 0이 방출되면 filter → map → take → collect를 거쳐 출력되고, 그다음에 1이 방출됩니다. 리스트처럼 단계마다 중간 컬렉션을 만들지 않습니다. 또 take(3)이 세 번째 값을 받으면 업스트림을 취소하므로, repeat(10)의 나머지 반복은 실행되지 않습니다. 이 취소는 업스트림 블록 안에서 AbortFlowException으로 전달되는데, flow { } 안에서 try { emit(x) } catch (e: Exception) { }로 모든 예외를 삼키면 이 취소까지 막혀 IllegalStateException: Flow exception transparency is violated가 날 수 있습니다. Flow 안에서 예외를 처리할 때는 catch 연산자를 쓰는 것이 규칙입니다.

Flow 백프레셔 (buffer, conflate, collectLatest)

fun main() = runBlocking {
    // 빠른 생산자
    val fastFlow = flow {
        repeat(10) { i ->
            emit(i)
            delay(10)  // 빠름
        }
    }
    // 느린 소비자
    println("=== No buffer ===")
    fastFlow.collect { 
        delay(100)  // 느림
        println(it) 
    }
    println("\n=== With buffer ===")
    fastFlow.buffer(5).collect { 
        delay(100)
        println(it) 
    }
    println("\n=== Conflate (latest only) ===")
    fastFlow.conflate().collect { 
        delay(100)
        println(it) 
    }
    println("\n=== collectLatest (cancel previous) ===")
    fastFlow.collectLatest { 
        delay(100)
        println(it) 
    }
}

네 가지 방식은 “소비자가 느릴 때 무엇을 희생할지”가 다릅니다. 버퍼가 없으면 emit이 수집 블록이 끝날 때까지 기다리므로 생산과 소비가 순차로 실행되어 전체 시간이 (10ms + 100ms) × 10이 됩니다. buffer(5)는 생산자와 소비자를 별도 코루틴으로 나눠 동시에 실행하므로, 값은 하나도 버리지 않으면서 전체 시간이 소비자 속도에 가까워집니다. conflate()는 소비자가 바쁜 동안 들어온 값 중 가장 최신 것만 남기고 나머지를 버리며, collectLatest는 새 값이 오면 처리 중이던 블록을 취소하고 새 값을 처리합니다. 센서 값이나 검색어 자동완성처럼 최신 값만 의미가 있으면 뒤의 두 방식이, 결제 이벤트처럼 하나도 버리면 안 되면 buffer가 맞습니다.

핫에 가까운 이벤트: SharedFlow

기본 SharedFlow

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val scope = CoroutineScope(Dispatchers.Default + SupervisorJob())
    val hot = MutableSharedFlow<Int>(
        replay = 0,              // 새 구독자에게 재전송할 개수
        extraBufferCapacity = 16 // 버퍼 크기
    )
    // 구독자 1
    val sub1 = scope.launch { 
        hot.collect { println("Subscriber 1: $it") } 
    }
    // 구독자 2
    val sub2 = scope.launch { 
        hot.collect { println("Subscriber 2: $it") } 
    }
    delay(100)  // 구독자 준비 대기
    // 이벤트 발행
    repeat(3) { 
        hot.emit(it)
        delay(50)
    }
    delay(100)
    sub1.cancel()
    sub2.cancel()
    scope.cancel()
}

출력:

Subscriber 1: 0
Subscriber 2: 0
Subscriber 1: 1
Subscriber 2: 1
Subscriber 1: 2
Subscriber 2: 2

여기서 delay(100) // 구독자 준비 대기는 장식이 아니라 필수입니다. replay = 0인 SharedFlow는 그 순간 구독 중인 수집자에게만 값을 전달하고, 구독자가 없으면 값을 그냥 버립니다. launch는 코루틴을 예약만 할 뿐 즉시 collect까지 실행하지 않으므로, 대기 없이 바로 emit하면 첫 이벤트들이 조용히 사라집니다. 저는 화면 진입 시 “초기 로딩 완료” 이벤트를 SharedFlow로 보냈다가, 화면의 수집 코드가 조금 늦게 시작되는 기기에서만 이벤트가 사라지는 버그를 겪은 적이 있습니다. 시간 지연으로 맞추는 대신 hot.subscriptionCount.first { it >= 2 }처럼 구독자 수를 기다리거나, 반드시 전달되어야 하는 일회성 이벤트라면 버퍼를 가진 Channel + receiveAsFlow()를 쓰는 편이 안전합니다.

extraBufferCapacity = 16은 느린 구독자 때문에 emit이 중단되기 전까지 쌓아 둘 수 있는 개수입니다. 버퍼가 차면 기본 정책(BufferOverflow.SUSPEND)에 따라 emit이 가장 느린 구독자를 기다리므로, 구독자 하나가 느리면 모든 구독자가 함께 느려집니다. 버퍼가 있을 때만 성공하는 tryEmit은 버퍼가 없거나 가득 차면 false를 돌려주고 값을 버린다는 점도 주의해야 합니다.

StateFlow (상태 관리)

data class UiState(val count: Int, val message: String)
class ViewModel {
    private val _state = MutableStateFlow(UiState(0, "Initial"))
    val state: StateFlow<UiState> = _state.asStateFlow()
    fun increment() {
        _state.update { it.copy(count = it.count + 1) }
    }
    fun setMessage(msg: String) {
        _state.update { it.copy(message = msg) }
    }
}
fun main() = runBlocking {
    val vm = ViewModel()
    // 상태 구독
    val job = launch {
        vm.state.collect { state ->
            println("State: count=${state.count}, message=${state.message}")
        }
    }
    delay(100)
    // 상태 변경
    vm.increment()
    delay(50)
    vm.increment()
    delay(50)
    vm.setMessage("Updated")
    delay(100)
    job.cancel()
}

출력:

State: count=0, message=Initial
State: count=1, message=Initial
State: count=2, message=Initial
State: count=2, message=Updated

StateFlow는 SharedFlow의 특수한 형태로, 항상 현재 값 하나를 갖고(replay = 1), 새 구독자에게 즉시 그 값을 주며, 이전 값과 같으면(equals) 방출하지 않습니다. 그래서 data class로 상태를 만들고 copy로 갱신하는 것이 관용구입니다. 반대로 가변 객체를 바꾼 뒤 같은 인스턴스를 다시 대입하면 equals가 참이라 수집자에게 아무것도 전달되지 않습니다. update { }는 compare-and-set으로 동작해 여러 코루틴이 동시에 갱신해도 증가가 유실되지 않으므로, _state.value = _state.value.copy(...)보다 안전합니다.

StateFlow는 또 느린 수집자에게 중간 값을 건너뛰게 합니다(conflation). 위 예제에서 delay(50) 없이 increment()를 연달아 호출하면 수집자는 count=1을 못 보고 count=2만 볼 수 있습니다. 상태는 “지금 어떤지”만 중요하므로 이것이 올바른 동작이지만, 모든 변화를 처리해야 하는 이벤트를 StateFlow로 보내면 이벤트가 사라집니다.

SharedFlow vs StateFlow

fun main() = runBlocking {
    // SharedFlow: 이벤트 스트림
    val events = MutableSharedFlow<String>()
    launch {
        events.collect { println("Event: $it") }
    }
    events.emit("Click")
    events.emit("Scroll")
    // StateFlow: 상태 (항상 최신 값 유지)
    val state = MutableStateFlow("Initial")
    launch {
        state.collect { println("State: $it") }
    }
    state.value = "Loading"
    state.value = "Success"
    delay(100)
    coroutineContext.cancelChildren()  // collect는 끝나지 않으므로 직접 취소해야 main이 종료됨
}

이 예제를 실제로 실행해 보면 기대와 다른 결과가 나옵니다. launch 직후 바로 emit("Click")을 호출하므로 수집 코루틴이 아직 시작되지 않은 상태이고, 구독자가 없는 SharedFlow는 두 이벤트를 모두 버립니다. StateFlow 쪽은 수집자가 시작되기 전에 값이 "Success"까지 바뀌므로 State: Success 한 줄만 출력됩니다. 즉 SharedFlow는 구독 전 이벤트를 잃고, StateFlow는 중간 상태를 건너뛴다는 두 성질이 한 예제에 모두 드러납니다. 또 SharedFlow와 StateFlow의 collect는 절대 스스로 끝나지 않으므로, 마지막 줄처럼 자식 코루틴을 취소하지 않으면 runBlocking이 영원히 반환되지 않습니다. 테스트 코드에서 이 점을 놓치면 테스트가 타임아웃으로 멈춥니다.


고급 활용: 백프레셔와 버퍼

Channel 백프레셔 전략

버퍼 용량 제한

fun main() = runBlocking {
    // 버퍼 크기 2
    val ch = Channel<Int>(capacity = 2)
    val producer = launch {
        repeat(5) { i ->
            println("Sending $i")
            ch.send(i)  // 버퍼 가득 차면 여기서 중단
            println("Sent $i")
        }
        ch.close()
    }
    delay(500)  // 소비자 늦게 시작
    val consumer = launch {
        for (x in ch) {
            println("Received $x")
            delay(200)
        }
    }
    producer.join()
    consumer.join()
}

출력:

Sending 0
Sent 0
Sending 1
Sent 1
Sending 2
(버퍼 가득, 대기)
(500ms 후 소비자 시작)
Received 0
Sent 2
Sending 3
Received 1
...

무제한 버퍼 (위험)

import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
fun main() = runBlocking {
    val ch = Channel<Int>(capacity = UNLIMITED)
    val producer = launch {
        repeat(1_000_000) { i ->
            ch.send(i)  // 절대 중단 안 됨 → 메모리 폭증 위험
        }
        ch.close()
    }
    val consumer = launch {
        for (x in ch) {
            delay(10)  // 느린 소비
        }
    }
    producer.join()
    consumer.join()
}

문제점: 생산자가 빠르면 메모리 무한 증가

Conflated Channel (최신 값만)

import kotlinx.coroutines.channels.Channel.Factory.CONFLATED
fun main() = runBlocking {
    val ch = Channel<Int>(capacity = CONFLATED)
    val producer = launch {
        repeat(10) { i ->
            ch.send(i)  // 이전 값 덮어씀
            delay(10)
        }
        ch.close()
    }
    delay(150)  // 소비자 늦게 시작
    val consumer = launch {
        for (x in ch) {
            println("Received $x")
        }
    }
    producer.join()
    consumer.join()
}

출력: 최신 값만 받음 (중간 값 손실)

Flow 백프레셔 전략

buffer() - 버퍼 추가

fun main() = runBlocking {
    flow {
        repeat(5) { i ->
            emit(i)
            println("Emitted $i")
        }
    }
    .buffer(2)  // 버퍼 크기 2
    .collect { 
        delay(100)  // 느린 소비자
        println("Collected $it") 
    }
}

conflate() - 최신 값만

fun main() = runBlocking {
    flow {
        repeat(10) { i ->
            emit(i)
            delay(10)
        }
    }
    .conflate()  // 소비자가 바쁘면 중간 값 스킵
    .collect { 
        delay(100)
        println("Collected $it") 
    }
}

출력: 0과, 처리가 끝날 때마다 그 시점의 최신 값 몇 개만 출력 (중간 값 스킵, 정확한 값은 타이밍에 따라 다름)

collectLatest() - 이전 수집 취소

fun main() = runBlocking {
    flow {
        repeat(5) { i ->
            emit(i)
            delay(50)
        }
    }
    .collectLatest { value ->
        println("Collecting $value")
        delay(200)  // 느린 처리
        println("Processed $value")
    }
}

출력:

Collecting 0
Collecting 1  (0 처리 취소)
Collecting 2  (1 처리 취소)
Collecting 3  (2 처리 취소)
Collecting 4  (3 처리 취소)
Processed 4   (마지막만 완료)

구조화된 동시성 (Structured Concurrency)

coroutineScope로 자식 묶기

suspend fun fetchUserData(userId: Int): UserData = coroutineScope {
    // 병렬로 3개 API 호출
    val nameDeferred = async { fetchName(userId) }
    val ordersDeferred = async { fetchOrders(userId) }
    val profileDeferred = async { fetchProfile(userId) }
    // 하나라도 실패하면 나머지 자동 취소
    UserData(
        name = nameDeferred.await(),
        orders = ordersDeferred.await(),
        profile = profileDeferred.await()
    )
}
suspend fun fetchName(userId: Int): String {
    delay(100)
    return "User-$userId"
}
suspend fun fetchOrders(userId: Int): List<String> {
    delay(150)
    return listOf("Order1", "Order2")
}
suspend fun fetchProfile(userId: Int): String {
    delay(80)
    return "Profile-$userId"
}
fun main() = runBlocking {
    try {
        val data = fetchUserData(123)
        println("Fetched: $data")
    } catch (e: Exception) {
        println("Failed: ${e.message}")
    }
}

Channel 파이프라인 + 구조화된 동시성

fun main() = runBlocking {
    coroutineScope {
        val numbers = Channel<Int>(capacity = 10)
        val squares = Channel<Int>(capacity = 10)
        // Stage 1: 숫자 생성
        launch {
            repeat(10) { i ->
                numbers.send(i)
            }
            numbers.close()
        }
        // Stage 2: 제곱 계산
        launch {
            for (n in numbers) {
                squares.send(n * n)
            }
            squares.close()
        }
        // Stage 3: 출력
        launch {
            for (s in squares) {
                println("Square: $s")
            }
        }
    }  // 모든 자식 작업 완료 대기
    println("Pipeline completed")
}

핫 Flow: SharedFlow와 StateFlow

SharedFlow (이벤트 스트림)

class EventBus {
    private val _events = MutableSharedFlow<String>(
        replay = 0,              // 새 구독자에게 재전송 안 함
        extraBufferCapacity = 64 // 버퍼 크기
    )
    val events: SharedFlow<String> = _events.asSharedFlow()
    suspend fun emit(event: String) {
        _events.emit(event)
    }
}
fun main() = runBlocking {
    val bus = EventBus()
    // 구독자 1
    val job1 = launch {
        bus.events.collect { println("Sub1: $it") }
    }
    // 구독자 2
    val job2 = launch {
        bus.events.collect { println("Sub2: $it") }
    }
    delay(100)
    // 이벤트 발행
    bus.emit("Event1")
    bus.emit("Event2")
    bus.emit("Event3")
    delay(100)
    job1.cancel()
    job2.cancel()
}

StateFlow (상태 관리)

class Counter {
    private val _count = MutableStateFlow(0)
    val count: StateFlow<Int> = _count.asStateFlow()
    fun increment() {
        _count.update { it + 1 }
    }
    fun decrement() {
        _count.update { it - 1 }
    }
}
fun main() = runBlocking {
    val counter = Counter()
    // 상태 구독
    val job = launch {
        counter.count.collect { count ->
            println("Count: $count")
        }
    }
    delay(100)
    counter.increment()  // Count: 1
    delay(50)
    counter.increment()  // Count: 2
    delay(50)
    counter.decrement()  // Count: 1
    delay(100)
    job.cancel()
}

replay 옵션 비교

fun main() = runBlocking {
    // replay = 0 (기본)
    val noReplay = MutableSharedFlow<Int>(replay = 0)
    noReplay.emit(1)
    noReplay.emit(2)
    
    launch {
        noReplay.collect { println("No replay: $it") }
    }
    // 출력: 없음 (이미 발행된 값)
    delay(100)
    // replay = 2
    val withReplay = MutableSharedFlow<Int>(replay = 2)
    withReplay.emit(1)
    withReplay.emit(2)
    withReplay.emit(3)
    
    launch {
        withReplay.collect { println("With replay: $it") }
    }
    // 출력: With replay: 2, With replay: 3 (마지막 2개)
    delay(100)
    coroutineContext.cancelChildren()  // SharedFlow의 collect는 끝나지 않으므로 취소
}

replay는 “늦게 온 구독자에게 과거 값 몇 개를 다시 보여 줄지”를 정합니다. replay = 0에서 구독자 없이 emit한 1, 2는 어디에도 저장되지 않고 사라지며, 이때 emit은 기다리지도 않고 즉시 반환됩니다. replay를 늘리면 늦은 구독자도 최근 이벤트를 받지만, 화면 회전처럼 구독이 다시 시작될 때마다 같은 이벤트(예: 토스트 메시지)가 다시 재생되는 부작용이 생깁니다. 일회성 UI 이벤트에 replay = 1을 쓰면 화면을 돌릴 때마다 같은 토스트가 뜨는 버그가 흔한 이유입니다.


성능·비교 요약

관점ChannelFlow (콜드)
시작 시점보통 생산자·소비자가 동시에 돌아감collect 시 업스트림 실행
재사용동일 채널 인스턴스를 여러 소비자와? 설계 필요collect마다 새 실행(기본)
백프레셔버퍼·블로킹 send연산자로 조절
테스트send/receive로 단계 분리runTest + 가상 시간
UI/상태채널만으로는 UI 모델이 부족한 경우 많음StateFlow로 상태 표현에 강함

실무 사례

  • 백그라운드 워커 풀 → 단일 집계기: Channel로 작업 큐를 두며, 소비자가 DB에 배치 기록.
  • REST 응답 스트리밍: 서버에서 청크를 밀 때는 프레임워크별로 다르지만, 동시성 경계는 Channel로 두기 쉽습니다.
  • UI 이벤트: 사용자 입력·네트워크 결과를 상태 한 방향으로 줄이려면 StateFlow/SharedFlow가 많이 사용됩니다.
  • 리포지토리 레이어: “한 번 호출해 여러 값”은 Flow, “동시에 실행되는 파이프라인”은 Channel 후보입니다.

트러블슈팅

증상: Flow가 두 번 실행된다
→ 콜드 Flow는 collect마다 처음부터 다시 돕니다. shareIn으로 핫으로 바꿀지, 의도인지 확인하세요.

증상: Channel에서 ClosedSendChannelException
→ 닫힌 채널에 send했는지, close() 순서가 맞는지 확인하세요.

증상: 생산이 너무 빨라 메모리가 급증
→ Channel.UNLIMITED, buffer(Channel.UNLIMITED), extraBufferCapacity를 크게 준 SharedFlow처럼 상한 없는 버퍼를 의심하세요. 용량과 onBufferOverflow(DROP_OLDEST/DROP_LATEST) 정책을 명시하세요.

증상: 프로그램이나 테스트가 끝나지 않는다
→ 닫지 않은 Channel을 for로 읽고 있거나, SharedFlow·StateFlow를 collect하는 코루틴이 남아 있는 경우입니다. 이들의 collect는 스스로 끝나지 않으므로 스코프를 취소하거나 take(n)·first()로 수집을 제한하세요.

증상: shareIn/stateIn으로 만든 Flow가 화면을 나가도 계속 돈다
→ SharingStarted.Eagerly나 Lazily는 구독자가 없어도 멈추지 않습니다. UI 상태라면 SharingStarted.WhileSubscribed(5_000)으로 마지막 구독자가 사라진 뒤 5초 후 업스트림을 멈추는 설정이 일반적입니다(화면 회전처럼 잠깐 끊기는 경우에는 다시 시작하지 않도록 여유를 둔 값입니다).

증상: 테스트가 불안정하다
→ runTest, StandardTestDispatcher, 가상 시간으로 결정적으로 만드세요.


마무리

Channel은 동시에 도는 생산·소비를 잇는 도구이며, Flow는 지연·중단 가능한 값의 나열을 표현하는 도구입니다. “핫/콜드·백프레셔·수집 시점”만 명확히 해도 선택이 단순해집니다. 팀 내 용어를 코루틴 vs 스레드와 맞춰 두면 리뷰 비용도 줄어듭니다.


자주 묻는 질문 (FAQ)

Q. Flow를 collect할 때마다 처음부터 다시 실행되는 이유는 무엇인가요?

A. flow { }로 만든 Flow는 콜드 스트림이라, collect를 호출할 때마다 빌더 블록이 새로 실행됩니다. 그래서 두 곳에서 collect하면 네트워크 요청이나 DB 조회도 두 번 일어납니다. 여러 수집자가 한 번의 실행 결과를 공유해야 한다면 shareIn이나 stateIn으로 핫 스트림(SharedFlow·StateFlow)으로 바꾸고, 그렇지 않다면 의도된 동작인지 먼저 확인합니다.


같이 보면 좋은 글