16장 플로우
16.1 플로우는 연속적인 값의 스트림을 모델링한다
16.1.1 플로우를 사용하면 배출되자마자 원소를 처리할 수 있다
일시 중단 함수의 한계:
// 일시 중단 함수는 단일 값만 반환 가능
suspend fun createValues(): List<Int> {
return buildList {
add(1)
delay(1.seconds)
add(2)
delay(1.seconds)
add(3)
delay(1.seconds)
}
}
// → 3초 후에 모든 값이 한번에 출력됨
플로우의 장점:
// 플로우는 값이 계산되자마자 즉시 처리 가능
fun createValues(): Flow<Int> {
return flow {
emit(1)
delay(1000.milliseconds)
emit(2)
delay(1000.milliseconds)
emit(3)
delay(1000.milliseconds)
}
}
// 사용법
myFlowOfValues.collect { log(it) }
// → 값이 배출되자마자 바로바로 출력됨
일시 중단 함수는 "한 번에 하나의 값만" 돌려줄 수 있으나
실제로 "시간에 따라 계속 변하는 값"이 필요한 경우가 많다.
- 온도계에서 계속 측정되는 온도
- 사용자가 검색창에 타이핑하는 글자들
- 주식 가격의 실시간 변화
이런 "연속된 값들의 스트림"을 다루는 게 플로우
- flow { }:값들을 만드는 곳
- emit(): 하나씩 떨어뜨리기
- collect { }: 값들을 받아서 처리
16.1.2 코틀린 플로우의 여러 유형
콜드 플로우(Cold Flow):
- 값이 실제로 소비되기 시작할 때만 값을 배출
- 비동기 데이터 스트림
핫 플로우(Hot Flow):
- 상태 플로우와 공유플로우 가 있다 (StateFlow와 SharedFlow -StateFlow는 SharedFlow의 특수한 형태)
- 값이 소비되고 있는지와 상관없이 독립적으로 값을 배출
- 브로드캐스트 방식으로 동작
16.2 콜드 플로우
16.2.1 flow 빌더 함수를 사용해 콜드 플로우 생성
기본 사용법:
val letters = flow {
log("Emitting A!")
emit("A")
delay(200.milliseconds)
log("Emitting B!")
emit("B")
}
// flow 빌더만 호출해도 아무 일이 일어나지 않음 → "콜드"라고 불리는 이유
무한 플로우 생성 가능:
val counterFlow = flow {
var x = 0
while (true) {
emit(x++)
delay(200.milliseconds)
}
}
16.2.2 콜드 플로우는 수집되기 전까지 작업을 수행하지 않는다
수집(collect) 시에만 실행:
letters.collect {
log("Collecting $it")
delay(500.milliseconds)
}
// collect를 여러 번 호출하면 코드가 여러 번 실행됨
letters.collect { log("(1) Collecting $it") }
letters.collect { log("(2) Collecting $it") }
16.2.3 플로우 수집 취소
코루틴 취소로 플로우 수집 중단:
val collector = launch {
counterFlow.collect {
println(it)
}
}
delay(5.seconds)
collector.cancel() // 플로우 수집이 중단됨
- job.cancel(): 수동 취소
- take(n): n개 후 자동 취소
- takeWhile { }: 조건이 false가 되면 취소
- first(): 첫 번째 값 후 취소
예시 : 사용자가 화면을 떠날 떄, 다른 작업으로 전환할 때 ,
16.2.4 콜드 플로우의 내부 구현
핵심 인터페이스:
interface Flow<T> {
suspend fun collect(collector: FlowCollector<T>)
}
interface FlowCollector<T> {
suspend fun emit(value: T)
}
동작 원리:
- collect 호출 → 플로우 빌더 함수 본문 실행
- emit 호출 → collect에 전달된 람다 호출
- 람다 실행 완료 → 빌더 함수로 돌아가 계속 실행
16.2.5 채널 플로우를 사용한 동시성 플로우
일반 플로우의 한계:
// 순차적으로 실행되어 느림
val randomNumbers = flow {
repeat(10) {
emit(getRandomNumber()) // 각각 500ms 소요 → 총 5초
}
}
채널 플로우로 동시성 처리:
val randomNumbers = channelFlow {
repeat(10) {
launch {
send(getRandomNumber()) // 동시에 실행 → 총 500ms
}
}
}
선택 기준:
- 기본적으로 일반 콜드 플로우 사용 (더 간단하고 성능이 좋음)
- 플로우 내에서 새로운 코루틴을 시작해야 하는 경우에만 채널 플로우 사용
16.3 핫 플로우
16.3.1 공유 플로우는 값을 구독자에게 브로드캐스트한다
공유 플로우 생성:
class RadioStation {
private val _messageFlow = MutableSharedFlow<Int>()
val messageFlow = _messageFlow.asSharedFlow()
fun beginBroadcasting(scope: CoroutineScope) {
scope.launch {
while(true) {
delay(500.milliseconds)
val number = Random.nextInt(0..10)
log("Emitting $number!")
_messageFlow.emit(number)
}
}
}
}
특징:
- 구독자 유무와 상관없이 배출 발생
- 여러 구독자가 같은 값을 받음
- 구독 시작 이후의 값만 수신
값 재생(Replay):
private val _messageFlow = MutableSharedFlow<Int>(replay = 5)
// 새 구독자에게 최근 5개 값을 재생
콜드 플로우를 공유 플로우로 변환:
val temps = getTemperatures()
val sharedTemps = temps.shareIn(this, SharingStarted.Lazily)
// SharingStarted 옵션:
// - Eagerly: 즉시 시작
// - Lazily: 첫 번째 구독자가 나타날 때 시작
// - WhileSubscribed: 구독자가 있을 때만 실행
16.3.2 시스템 상태 추적: 상태 플로우
상태 플로우 생성:
class ViewCounter {
private val _counter = MutableStateFlow(0) // 초기값 필요
val counter = _counter.asStateFlow()
fun increment() {
_counter.update { it + 1 } // 원자적 업데이트
}
}
// 현재 값에 즉시 접근 가능
println(vc.counter.value) // 일시 중단 없음
안전한 상태 업데이트:
// 위험: 비원자적 연산
fun increment() {
_counter.value++
}
//안전: 원자적 연산
fun increment() {
_counter.update { it + 1 }
}
동등성 기반 통합:
switch.turn(Direction.LEFT)
switch.turn(Direction.LEFT) // 같은 값이므로 배출되지 않음
switch.turn(Direction.RIGHT) // 다른 값이므로 배출됨
콜드 플로우를 상태 플로우로 변환:
val temps = getTemperatures()
val tempState = temps.stateIn(this) // 항상 즉시 시작
println(tempState.value) // 최신 값에 즉시 접근
16.3.3 상태 플로우와 공유 플로우의 비교
상태 플로우 장점:
- 더 간단한 API
- 현재 값에 즉시 접근 가능
- 동등성 기반 통합으로 불필요한 업데이트 방지
- 상태 관리에 특화
예시 - 메시지 브로드캐스터:
// 공유 플로우 방식 (복잡)
private val _messages = MutableSharedFlow<String>()
// 구독 전 메시지는 손실됨
// 상태 플로우 방식 (간단)
private val _messages = MutableStateFlow<List<String>>(emptyList())
fun addMessage(msg: String) {
_messages.update { it + msg } // 모든 메시지 기록 유지
}
16.3.4 언제 어떤 플로우를 사용할까?
플로우 선택 가이드:
플로우 타입 특징 사용 시기
| 콜드 플로우 | - 수집자에 의해 활성화 |
- 하나의 수집자
- 모든 배출을 받음
- 보통 완료됨 | - 네트워크 요청
- 데이터베이스 읽기
- 파일 처리
- 서비스 함수 | | 공유 플로우 SharedFlow | - 기본적으로 활성
- 여러 구독자
- 구독 시점부터 받음
- 완료되지 않음 | - 이벤트 브로드캐스트
- 알림 시스템
- 여러 구독자가 필요한 경우 | | 상태 플로우 StateFlow | - 현재 상태 표현
- 동등성 기반 통합
- 즉시 값 접근 가능 | - 상태 관리
- UI 상태
- 설정 값
- 카운터, 스위치 등 |
일반적인 패턴:
- 서비스 함수는 콜드 플로우로 선언
- 필요시 shareIn() 또는 stateIn()으로 핫 플로우로 변환
- 상태 관리가 필요하면 상태 플로우 우선 고려
- 단순 이벤트 브로드캐스트는 공유 플로우 사용
예시:
*// 콜드: API 호출 (매번 새로 요청)*
fun getUserData() = flow {
emit(api.getUser()) *// collect할 때마다 새로 호출*
}
*// 핫: 버튼 클릭 이벤트*
val buttonClicks = MutableSharedFlow<Unit>()
*// 핫: 로그인 상태*
val loginState = MutableStateFlow("로그아웃")
핵심
- 플로우는 시간에 따라 배출되는 값의 스트림을 처리하는 코루틴 기반 추상화
- 콜드 플로우는 "게으른" 특성으로 수집할 때만 실행
- 핫 플로우는 "능동적"으로 구독자 유무와 관계없이 배출
- 상태 플로우는 상태 관리에 최적화된 특별한 공유 플로우
- 적절한 플로우 선택이 성능과 코드 품질에 중요한 영향
17장 플로우 연산자
17.1 플로우 연산자로 플로우 조작
플로우는 컬렉션과 시퀀스처럼 다양한 연산자를 통해 조작하고 변환할 수 있습니다.
연산자 분류:
- 중간 연산자: 코드를 실행하지 않고 변경된 플로우를 반환
- 최종 연산자: 실제 코드를 실행하여 컬렉션, 개별 원소, 계산된 값을 반환
flow {} → map {} → filter {} → onEach {} → collect {}
↑ 중간 연산자들 ↑ 최종 연산자
업스트림 다운스트림
17.2 중간 연산자는 업스트림 플로우에 적용되고 다운스트림 플로우를 반환한다
용어 정리:
- 업스트림(상류): 연산자가 적용되는 플로우
- 다운스트림(하류): 연산자가 반환하는 플로우
17.2.1 업스트림 원소별로 임의의 값을 배출: transform 함수
map의 한계:
val names = flowOf("Jo", "May", "Sue")
val uppercasedNames = names.map {
it.uppercase() // 1개 입력 → 1개 출력만 가능
}
// DO MAY SUE
transform의 자유도:
val upperAndLowercasedNames = names.transform {
emit(it.uppercase()) // 1개 입력 →
emit(it.lowercase()) // 여러 개 출력 가능
}
// DO jo MAY may SUE sue
17.2.2 take나 관련 연산자는 플로우를 취소할 수 있다
자동 취소 기능:
val temps = getTemperatures() // 무한 플로우
temps
.take(5) // 5개만 가져온 후 자동으로 업스트림 플로우 취소
.collect { log(it) }
17.2.3 플로우의 각 단계 후킹: onStart, onEach, onCompletion, onEmpty
생명주기 연산자들:
fun process(flow: Flow<Int>) = flow
.onStart {
println("작업 시작!")
}
.onEach {
println("처리 중: $it")
}
.onEmpty {
println("비어있음 - 기본값 추가")
emit(0)
}
.onCompletion { cause ->
if (cause != null) {
println("문제로 종료: $cause")
} else {
println("성공적으로 완료!")
}
}
.collect {
println("최종 결과: $it")
}
// 일반적인 플로우
process(flowOf(1, 2, 3))
// 작업 시작!
// 처리 중: 1
// 최종 결과: 1
// 처리 중: 2
// 최종 결과: 2
// 처리 중: 3
// 최종 결과: 3
// 성공적으로 완료!
// 빈 플로우
process(flowOf())
// 작업 시작!
// 비어있음 - 기본값 추가
// 처리 중: 0
// 최종 결과: 0
// 성공적으로 완료!
17.2.4 다운스트림 연산자와 수집자를 위한 원소 버퍼링: buffer 연산자
문제: 생산자-수집자 동기화
fun getAllUserIds(): Flow<Int> {
return flow {
repeat(3) {
delay(200.milliseconds) // 데이터베이스 지연
emit(it)
}
}
}
// 문제가 있는 코드: 순차적 실행으로 느림
ids.map { getProfileFromNetwork(it) } // 각각 2초 소요
.collect { log("Got $it") }
// 총 실행시간: 약 7.2초
해결: buffer 연산자
ids.buffer(3) // 3개 원소를 버퍼에 저장
.map { getProfileFromNetwork(it) }
.collect { log("Got $it") }
// 총 실행시간: 약 6초 (버퍼링으로 인한 속도 향상)
버퍼 오버플로우 처리:
.buffer(capacity = 10, onBufferOverflow = BufferOverflow.SUSPEND) // 생산자 대기
.buffer(capacity = 10, onBufferOverflow = BufferOverflow.DROP_OLDEST) // 오래된 값 버림
.buffer(capacity = 10, onBufferOverflow = BufferOverflow.DROP_LATEST) // 최신 값 버림
17.2.5 중간값을 버리는 연산자: conflate 연산자
최신 값만 처리:
val temps = getTemperatures() // 500ms마다 온도 배출
temps
.onEach { log("Read $it from sensor") }
.conflate() // 중간값 버림
.collect {
log("Collected $it")
delay(1.seconds) // 수집자가 느림
}
// 결과:
// Read 20 from sensor → Collected 20
// Read -10 from sensor (버려짐)
// Read 3 from sensor → Collected 3
// Read 13 from sensor (버려짐)
// Read 26 from sensor → Collected 26
사용 시기:
- 플로우 값이 빠르게 구식이 되는 경우
- 최신 값만 중요한 경우 (예: 실시간 주가, 센서 데이터)
17.2.6 일정 시간 동안 값을 필터링하는 연산자: debounce 연산자
검색 입력 최적화 예제:
val searchQuery = flow {
emit("K")
delay(100.milliseconds)
emit("Ko")
delay(200.milliseconds)
emit("Kotl")
delay(500.milliseconds)
emit("Kotlin")
}
// 일반 수집: 매 입력마다 검색 실행
searchQuery.collect { log("Searching for $it") }
// Searching for K, Ko, Kotl, Kotlin (4번 검색)
// debounce 사용: 입력이 멈춘 후에만 검색 실행
searchQuery
.debounce(250.milliseconds) // 250ms 동안 입력이 없으면 배출
.collect { log("Searching for $it") }
// Searching for Kotl, Kotlin (2번만 검색)
사용 시기:
- 검색창 입력 최적화
- API 호출 횟수 제한
- 사용자 입력 디바운싱
17.2.7 플로우가 실행되는 코루틴 콘텍스트를 바꾸기: flowOn 연산자
디스패처 전환:
runBlocking {
flowOf(1)
.onEach { log("A") } // Default 디스패처에서 실행
.flowOn(Dispatchers.Default)
.onEach { log("B") } // IO 디스패처에서 실행
.flowOn(Dispatchers.IO)
.onEach { log("C") } // Main 디스패처에서 실행 (collect 콘텍스트)
.collect()
}
// 결과:
// [DefaultDispatcher-worker-3] A
// [DefaultDispatcher-worker-1] B
// [main] C
중요한 특징:
- flowOn은 업스트림 플로우에만 영향을 미침
- '콘텍스트 보존' 연산자라고도 불림
- 다운스트림은 영향 받지 않음
잘못된 생각: "flowOn을 쓰면 전체 플로우가 다른 스레드에서 실행된다"
올바른 이해: "flowOn은 그보다 위쪽(업스트림)에 있는 연산들만 다른 스레드로 보낸다"
flowOf(1, 2, 3)
.map { it * 2 } // 이것만 IO 스레드
.flowOn(Dispatchers.IO) // 경계선
.filter { it > 2 } // 이건 원래 스레드
.collect { println(it) } // 이것도 원래 스레드
17.3 커스텀 중간 연산자 만들기
중간 연산자의 역할:
- 동시에 수집자와 생산자 역할 수행
- 업스트림에서 수집 → 변환/처리 → 다운스트림에 배출
예제: 마지막 n개 원소의 평균 계산 연산자
fun Flow<Double>.averageOfLast(n: Int): Flow<Double> =
flow {
val numbers = mutableListOf<Double>()
collect { value ->
if (numbers.size >= n) {
numbers.removeFirst()
}
numbers.add(value)
emit(numbers.average())
}
}
// 사용법
flowOf(1.0, 2.0, 30.0, 121.0)
.averageOfLast(3)
.collect { print("$it ") }
// 1.0 1.5 11.0 51.0
17.4 최종 연산자는 업스트림 플로우를 실행하고 값을 계산한다
최종 연산자의 특징:
- 실제로 플로우 코드를 실행
- 항상 일시 중단 함수
- 단일 값, 컬렉션, 또는 부수 효과 수행
주요 최종 연산자:
collect 계열:
// 기본 collect
flow.collect { log(it) }
// 파라미터 없는 collect (부수 효과만)
flow.onEach { log(it) }.collect()
단일 값 반환:
val firstTemp = getTemperatures().first() // 첫 번째 값만
val firstOrNull = getTemperatures().firstOrNull() // 안전한 첫 번째 값
17.4.1 프레임워크는 커스텀 연산자를 제공한다
Jetpack Compose 통합 예제:
@Composable
fun TemperatureDisplay(temps: Flow<Int>) {
val temperature = temps.collectAsState(null) // Flow → State 변환
Box {
temperature.value?.let {
Text("The current temperature is $it!")
}
}
}
플로우 연산자 활용 패턴 정리
성능 최적화 연산자
연산자 목적 사용 시기
| buffer | 생산자-수집자 분리 | 처리 속도가 다를 때 |
| conflate | 중간값 버림 | 최신 값만 중요할 때 |
| debounce | 입력 지연 처리 | 연속 입력 최적화 |
생명주기 연산자
연산자 실행 시점
| onStart | 수집 시작 전 |
| onEach | 각 원소 처리 시 |
| onEmpty | 빈 플로우일 때 |
| onCompletion | 완료/취소/오류 시 |
변환 연산자
연산자 특징
| map | 1:1 변환 |
| transform | 1:N 변환 (자유도 높음) |
| take | 지정된 개수만 취소 후 취소 |
실무 활용 팁
1. 검색 기능 최적화:
searchInput
.debounce(300.milliseconds) // 입력 완료 대기
.filter { it.length > 2 } // 최소 길이 필터
.map { searchAPI(it) } // API 호출
.collect { updateUI(it) } // UI 업데이트
2. 실시간 데이터 처리:
sensorData
.buffer(capacity = 100) // 버퍼링으로 성능 향상
.conflate() // 최신값만 처리
.flowOn(Dispatchers.IO) // I/O 스레드에서 실행
.collect { processData(it) }
3. UI 상태 관리:
userActions
.onStart { showLoading() }
.map { processAction(it) }
.onEmpty { showEmptyState() }
.onCompletion { hideLoading() }
.collect { updateUI(it) }
핵심
- 중간 연산자는 플로우 체인을 구성하지만 실행하지는 않음
- 최종 연산자가 호출되어야 비로소 전체 체인이 실행됨
- buffer, conflate, debounce로 성능 최적화 가능
- flowOn으로 실행 콘텍스트 제어 가능 (업스트림에만 영향)
- 생명주기 연산자로 각 단계별 부수 효과 처리
- transform으로 유연한 1:N 변환 가능
- 커스텀 연산자 작성 시 flow { collect { emit() } } 패턴 사용
18장 코틀린 오류 처리와 테스트
18.1 코루틴 내부에서 던져진 오류 처리
18.1.1 코루틴 빌더를 try-catch로 감싸는 것은 효과가 없다
잘못된 방법:
fun main(): Unit = runBlocking {
try {
launch {
throw UnsupportedOperationException("Ouch!")
}
} catch (u: UnsupportedOperationException) {
println("Handled $u") // 실행되지 않음!
}
}
// Exception in thread "main" java.lang.UnsupportedOperationException: Ouch!
올바른 방법:
fun main(): Unit = runBlocking {
launch {
try {
throw UnsupportedOperationException("Ouch!")
} catch (u: UnsupportedOperationException) {
println("Handled $u") // 정상 실행됨
}
}
}
// Handled java.lang.UnsupportedOperationException: Ouch!
코루틴 빌더는 "새로운 코루틴을 만드는 것"이지 "바로 실행하는 것"이 아니기 때문
핵심 개념:
- 코루틴 빌더(launch, async)는 새로운 코루틴을 생성
- 새 코루틴에서 발생한 예외는 코루틴 빌더를 감싸는 catch 블록에서 잡히지 않음
- 예외 처리는 코루틴 내부에서 해야 함
18.1.2 async의 예외 처리
async와 await():
fun main(): Unit = runBlocking {
val myDeferredInt: Deferred<Int> = async {
throw UnsupportedOperationException("Ouch!")
}
try {
val i: Int = myDeferredInt.await()
println(i)
} catch (u: UnsupportedOperationException) {
println("Handled: $u")
}
}
특징:
- async에서 발생한 예외는 await() 호출 시 다시 던져짐
- await()를 try-catch로 감싸서 처리 가능
- 원래 예외도 여전히 부모 코루틴에 전파됨
18.2 코틀린 코루틴에서의 오류 전파
18.2.1 자식이 실패하면 모든 자식을 취소하는 코루틴
구조적 동시성의 오류 전파 규칙:
- 실패한 자식 코루틴이 부모에게 예외 전파
- 부모는 다른 모든 자식 취소
- 부모도 같은 예외로 실패
- 예외를 상위 계층으로 전파
예제:
fun main(): Unit = runBlocking {
launch {
try {
while (true) {
println("Heartbeat!")
delay(500.milliseconds)
}
} catch (e: Exception) {
println("Heartbeat terminated: $e")
throw e
}
}
launch {
delay(1.seconds)
throw UnsupportedOperationException("Ow!")
}
}
// Heartbeat!
// Heartbeat!
// Heartbeat terminated: kotlinx.coroutines.JobCancellationException: Parent job is Cancelling
// Exception in thread "main" java.lang.UnsupportedOperationException: Ow!
18.2.2 구조적 동시성은 코루틴 스코프를 넘는 예외에만 영향을 미친다
코루틴 안에서 예외를 잡으면 가족에게 피해 안 줌
: try-catch가 코루틴 내부에 있으면 예외가 "코루틴 경계"를 넘지 않아서 다른 형제들에게 피해를 주지 않음
스코프 내 예외 처리:
fun main(): Unit = runBlocking {
launch {
// 하트비트 코루틴 (변경사항 없음)
}
launch {
try {
delay(1.seconds)
throw UnsupportedOperationException("Ow!")
} catch(u: UnsupportedOperationException) {
println("Caught $u") // 예외를 잡아서 전파되지 않음
}
}
}
// Heartbeat!
// Heartbeat!
// Caught java.lang.UnsupportedOperationException: Ow!
// Heartbeat! (계속 실행됨)
// Heartbeat!
// 잘못된 위치 (코루틴 밖)
try {
launch { throw Exception() } // 소용없음
} catch (e: Exception) { }
// 올바른 위치 (코루틴 안)
launch {
try { throw Exception() } catch (e: Exception) { } // 효과적
}
18.2.3 슈퍼바이저는 부모와 형제가 취소되지 않게 한다
SupervisorJob의 특징:
- 자식이 실패해도 부모와 형제 코루틴은 생존
- 예외를 상위 계층으로 전파하지 않음
- 구조적 동시성의 "경계" 역할
supervisorScope 사용:
fun main(): Unit = runBlocking {
supervisorScope {
launch {
try {
while (true) {
println("Heartbeat!")
delay(500.milliseconds)
}
} catch (e: Exception) {
println("Heartbeat terminated: $e")
throw e
}
}
launch {
delay(1.seconds)
throw UnsupportedOperationException("Ow!")
}
}
}
// Heartbeat!
// Heartbeat!
// Exception in thread "main" java.lang.UnsupportedOperationException: Ow!
// Heartbeat! (계속 실행됨)
// Heartbeat!
18.3 CoroutineExceptionHandler: 예외 처리를 위한 마지막 수단
18.3.1 예외 핸들러 정의와 사용
커스텀 예외 핸들러:
class ComponentWithScope(dispatcher: CoroutineDispatcher = Dispatchers.Default) {
private val exceptionHandler = CoroutineExceptionHandler { _, e ->
println("[ERROR] ${e.message}")
}
private val scope = CoroutineScope(
SupervisorJob() + dispatcher + exceptionHandler
)
fun action() = scope.launch {
throw UnsupportedOperationException("Ouch!")
}
}
fun main() = runBlocking {
val supervisor = ComponentWithScope()
supervisor.action()
delay(1.seconds)
}
// [ERROR] Ouch!
중요한 규칙:
- 예외 핸들러는 루트 코루틴에서만 호출됨
- 중간 계층의 핸들러는 무시됨
- launch로 시작된 최상위 코루틴에만 적용됨
18.3.2 launch와 async에 적용할 때의 차이점
launch vs async:
// launch: 예외 핸들러가 처리
scope.launch {
throw Exception("문제!")
} // 예외 핸들러 실행
// async: await하는 곳에서 처리
val result = scope.async {
throw Exception("문제!")
}
try {
result.await() // 여기서 예외 발생
} catch (e: Exception) {
// 여기서 직접 처리
}
이유:
- async의 예외는 await() 호출자가 처리해야 함 : "결과를 돌려줘야 하는" 작업이라서, 예외도 await하는 사람이 직접 처리해야
- launch의 예외는 "발사 후 망각(fire-and-forget)" 방식으로 핸들러가 처리 : 문제가 생기면 스스로 예외 핸들러에게 알려줌
18.4 플로우에서 예외 처리
18.4.1 기본적인 플로우 예외 처리
예외를 던지는 플로우:
class UnhappyFlowException: Exception()
val exceptionalFlow = flow {
repeat(5) { number ->
emit(number)
}
throw UnhappyFlowException()
}
fun main() = runBlocking {
val transformedFlow = exceptionalFlow.map { it * 2 }
try {
transformedFlow.collect {
print("$it ")
}
} catch (u: UnhappyFlowException) {
println("\\nHandled: $u")
}
}
// 0 2 4 6 8
// Handled: UnhappyFlowException
18.4.2 catch 연산자로 업스트림 예외 처리
catch 연산자 사용:
fun main() = runBlocking {
exceptionalFlow
.catch { cause ->
println("\\nHandled: $cause")
emit(-1) // 기본값 배출
}
.collect {
print("$it ")
}
}
// 0 1 2 3 4
// Handled: UnhappyFlowException
// -1
중요한 특징:
- catch는 업스트림 예외만 처리
- 다운스트림에서 발생한 예외는 처리하지 않음
- 취소 예외는 자동으로 인식해서 catch 블록 호출 안 함
업스트림 vs 다운스트림:
exceptionalFlow
.map { it + 1 } // 업스트림
.catch { cause -> // 여기서 위의 예외들 처리
println("\\nHandled $cause")
}
.onEach { // 다운스트림
throw UnhappyFlowException() // 이 예외는 catch되지 않음!
}
.collect()
18.4.3 술어가 참일 때 플로우의 수집 재시도: retry 연산자
retry 연산자:
val unstableNetworkFlow = flow {
println("네트워크 요청 중...")
if (Random.nextDouble() < 0.7) { // 70% 실패 확률
throw Exception("네트워크 오류!")
}
emit("데이터 받음")
}
unstableNetworkFlow
.retry(3) { cause ->
println("재시도: $cause")
cause is Exception // Exception이면 재시도
}
.collect {
println("최종 결과: $it")
}
특징:
- 업스트림 플로우를 처음부터 다시 시작
- 모든 중간 연산자도 다시 실행됨
- 부수 효과가 있는 작업은 멱등성 확보 필요
18.5 코루틴과 플로우 테스트
18.5.1 코루틴을 사용하는 테스트를 빠르게 만들기: 가상 시간과 테스트 디스패처
문제: 실시간 테스트의 느림
// runBlocking 사용: 실제 20초 대기
@Test
fun testDelay() = runBlocking {
delay(20.seconds) // 실제로 20초 기다림
// 테스트가 느려짐
}
해결: runTest 사용
@Test
fun testDelay() = runTest {
val startTime = System.currentTimeMillis()
delay(20.seconds) // 가상 시간으로 즉시 완료
println(System.currentTimeMillis() - startTime)
// 11 (밀리초만 소요)
}
18.5.2 가상 시간 제어
TestCoroutineScheduler 함수들:
@Test
fun testDelay() = runTest {
var x = 0
launch {
delay(500.milliseconds)
x++
}
launch {
delay(1.second)
x++
}
println(currentTime) // 0
delay(600.milliseconds) // 가상 시간 진행
assertEquals(1, x)
println(currentTime) // 600
delay(500.milliseconds)
assertEquals(2, x)
println(currentTime) // 1100
}
스케줄러 제어 함수들:
- runCurrent(): 현재 예약된 모든 코루틴 실행
- advanceUntilIdle(): 예약된 모든 코루틴 실행 (미래 포함)
- currentTime: 가상 시간 조회
- delay(): 가상 시계 진행
18.5.3 터빈(Turbine)으로 플로우 테스트 : 플로우 전용 테스트 도구
기본적인 플로우 테스트:
val myFlow = flow {
emit(1)
emit(2)
emit(3)
}
// 일반적인 방법
@Test
fun doTest() = runTest {
val results = myFlow.toList()
assertEquals(3, results.size)
}
// ---
val infiniteFlow = flow {
var count = 0
while (true) {
emit(count++)
delay(1.seconds)
}
}
// toList()로는 테스트 불가 (무한히 실행됨)
터빈 라이브러리 사용:
@Test
fun infiniteFlowTest() = runTest {
infiniteFlow.test {
assertEquals(0, awaitItem()) // 첫 번째 값
assertEquals(1, awaitItem()) // 두 번째 값
assertEquals(2, awaitItem()) // 세 번째 값
cancelAndIgnoreRemainingEvents() // 테스트 종료
}
}
터빈의 장점:
- 더 세밀한 플로우 동작 검증
- 무한 플로우 테스트 지원
- 복잡한 플로우 시나리오 처리
- 모든 원소가 적절히 소비되도록 보장
시간 제어 함수들:
runTest {
delay(1.seconds) *// 가상으로 1초 진행*
advanceTimeBy(500.ms) *// 500ms 진행*
runCurrent() *// 현재 스케줄된 작업 실행*
advanceUntilIdle() *// 모든 작업 완료까지 진행*
currentTime *// 현재 가상 시간 확인*
}
터빈 주요 함수들:
flow.test {
awaitItem() *// 다음 값 대기*
awaitComplete() *// 완료 대기*
awaitError() *// 에러 대기*
cancelAndIgnoreRemainingEvents() *// 테스트 종료*
}
실무 활용 팁
1. 예외 계층 설계:
- 애플리케이션 최상위: SupervisorJob + ExceptionHandler
- 작업 단위: 일반 Job (형제 취소 허용)
- 개별 작업: try-catch로 로컬 처리
2. 플로우 오류 대응:
- 네트워크 요청: retry + catch 조합
- 사용자 입력: catch로 기본값 제공
- 센서 데이터: conflate + catch
3. 테스트 최적화:
- runTest로 가상 시간 활용
- 디스패처를 주입 가능하게 설계
- 터빈으로 플로우 동작 검증
중요한 주의사항
- CancellationException은 삼키지 말 것 - 코루틴 생명주기의 자연스러운 부분
- 중간 계층의 ExceptionHandler는 동작하지 않음 - 루트 코루틴에만 설치
- async의 예외는 await()에서 처리 - ExceptionHandler 호출되지 않음
- catch는 업스트림만 처리 - 다운스트림 예외는 별도 처리 필요
- 테스트 디스패처 외의 디스패처는 가상 시간 영향 없음 - 명시적 디스패처 주입 필요
'language > Kotilin In Action' 카테고리의 다른 글
| 7주차(14,15) (1) | 2025.09.01 |
|---|---|
| 6주차 (13장) (2) | 2025.09.01 |
| 5주차(11~12장) (1) | 2025.09.01 |
| 4주차(9~10장) (1) | 2025.09.01 |
| 3주차 (6~8장) (1) | 2025.09.01 |