Flow
Flow
Flow는 비동기적으로 만들어질 수 있는 값의 순차적 스트림을 나타내요. 하나의 값을 반환하는 일시 중단 함수(suspending function)와 달리, Flow를 쓰면 시간에 따라 여러 개의 순차 값을 다룰 수 있어요.
출처: Kotlin 공식 문서
본문
Flow로 데이터를 점진적으로 로딩하는 플로우 파이프라인을 만들고, 이벤트 스트림에 반응하고, 구독(subscription) 스타일의 API를 모델링할 수 있어요.
플로우 파이프라인은 다음 역할이 얽힌 연산의 연속이에요.
- Emitter(생산자): 값을 만들어 내요.
- Intermediate operator(중간 연산자, 선택): Flow에서 값을 소비해 연산을 적용하고 다른 Flow를 반환해요.
- Collector(소비자): Flow에서 값을 소비해요.
이 파이프라인 역할들이 어떻게 함께 동작하는지 보여 주는 간단한 예시예요.
import kotlinx.coroutines.flow.*
//sampleStart
suspend fun main() {
// 생산자가 값을 만들어요
flowOf(0x4B, 0x6F, 0x74, 0x6C, 0x69, 0x6E)
// 중간 연산자가 값을 소비하고,
// 연산을 적용한 뒤 다른 Flow를 반환해요
.map { value -> value.toChar() }
// 소비자가 변환된 값을 받아요
.collect { updatedValue ->
println("Say '$updatedValue'!")
}
}
//sampleEnd
Flow에서 값은 생산자에서 소비자 쪽으로, 업스트림(upstream)에서 다운스트림(downstream)으로 이동해요. 중간 연산자는 업스트림 Flow를 수집(collect)하고, 그 값에 연산을 적용한 뒤 새 다운스트림 Flow를 반환해요. 그 다운스트림 Flow는 다음 소비자의 업스트림 Flow가 될 수 있어요.
Kotlin은 다음 Flow 타입을 제공해요.
- 콜드 Flow(Cold flows)는 수집(collect)될 때 값을 만들기 시작해요. 소비자마다 새롭고 독립적인 Flow 실행을 촉발해요.
- 핫 Flow(Hot flows)는 소비자와 무관하게 값을 방출하고, 모든 소비자와 같은 값 스트림을 공유해요.
참고: Turbine 라이브러리로 Kotlin Flow를 테스트할 수 있어요. 유닛 테스트에서 Flow 방출을 수집하고 단언하는 걸 단순하게 만들어 주는데, 완료와 실패 케이스까지 포함해요.
콜드 Flow
시퀀스처럼 콜드 Flow는 지연(lazy)돼요. 콜드 Flow 빌더의 코드 블록은 소비자가 수집할 때까지 실행되지 않아요. 새 소비자마다 Flow의 새 실행이 시작돼요.
콜드 Flow 만들기
콜드 Flow를 만들려면 flow() 빌더 함수를 사용해요. 그 블록 안에서 emit() 함수로 값을 소비자에게 방출해요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun main() {
// Flow를 만들어요
val pageFlow = flow {
for (page in 1..3) {
println("Loading page $page...")
// 각 페이지가 로딩될 때마다 방출해요
emit("Page $page")
}
}
println("Creating a cold flow doesn't run it!")
}
//sampleEnd
이 예시에서 flow() 빌더 함수는 Flow<T>를 반환하지만, 그 블록을 실행하기 시작하지는 않아요. 콜드 Flow는 요리 레시피 같아요. 값을 어떻게 만들지 정의하지만, 수집할 때만 값을 만들기 시작하죠.
다음 함수들로 콜드 Flow를 만들 수도 있어요.
예시를 볼게요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() {
// 제공된 값들로 Flow를 만들어요
val predefinedPageFlow = flowOf("Page 1", "Page 2", "Page 3")
// 범위로 Flow를 만들어요
val generatedPageFlow = (1..3).asFlow()
}
콜드 Flow 수집하기
콜드 Flow를 수집하려면 collect() 함수를 사용해요. 이 함수는 업스트림 Flow의 방출을 촉발해요. collect()에 람다를 넘기면 각 방출된 값을 받아요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
suspend fun main() {
withContext(Dispatchers.Default) {
val pageFlow = flow {
for (page in 1..3) {
println("Loading page $page...")
emit("Page $page")
}
}
// 각 방출된 페이지를 받는 람다로 Flow를 수집해요
pageFlow.collect { page ->
println("Processing $page...")
delay(100.milliseconds)
println("Done processing $page.")
}
}
}
//sampleEnd
collect()를 호출할 때마다 콜드 Flow 전체가 처음부터 실행돼요. 여러 소비자가 같은 콜드 Flow를 수집하면, 각 소비자가 자신만의 수집을 촉발해요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
suspend fun main() {
val pageFlow = flow {
// 현재 코루틴의 이름을 읽어요
val coroutineName = currentCoroutineContext()[CoroutineName]?.name
println("Starting emissions in $coroutineName")
for (page in 1..3) {
println("Loading page $page in $coroutineName")
emit("Page $page")
}
println("Done emitting in $coroutineName")
}
withContext(Dispatchers.Default) {
// 각 페이지를 느리게 처리하는 소비자를 실행해요
launch(CoroutineName("a slow coroutine")) {
pageFlow.collect {
println("Processing $it slowly")
delay(100.milliseconds)
println("Done processing $it slowly")
}
}
// 각 페이지를 빠르게 처리하는 소비자를 실행해요
launch(CoroutineName("a fast coroutine")) {
pageFlow.collect {
println("Processing $it quickly")
delay(10.milliseconds)
println("Done processing $it quickly")
}
}
}
}
//sampleEnd
이 예시에서 CoroutineName은 각 코루틴에 이름을 더해요. CoroutineName을 디버깅에 쓸 수 있는데, 여기서는 어떤 소비자가 각 수집을 실행하는지 보여 주는 데 도움을 줘요.
중간 Flow 연산자
중간 연산자는 업스트림 Flow에 연산을 적용하고 새 다운스트림 Flow를 반환해요. 이들은 콜드라서 반환된 Flow는 수집되기 전까지 값 처리를 시작하지 않아요. 업스트림 Flow가 핫이어도 마찬가지예요.
kotlinx.coroutines 라이브러리는 Flow를 변환하고 처리하는 폭넓은 중간 Flow 연산자를 제공해요. 내장 연산자가 제공하지 않는 동작이 필요하다면 직접 커스텀 연산자를 정의할 수도 있어요.
다음은 각 방출된 값에 변환을 적용하는 단순화된 커스텀 .map() 연산자의 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
// 기본 .map() 연산자의 단순화된 커스텀 구현
fun <T, R> Flow<T>.myMap(transform: suspend (value: T) -> R): Flow<R> = flow {
// 업스트림 Flow에서 값을 수집해요
[email protected] { value ->
// 각 수집된 값을 변환하고 결과를 방출해요
emit(transform(value))
}
}
suspend fun main() {
// Flow를 만들고 커스텀 map 연산자를 적용한 뒤 변환된 값을 수집해요
flowOf(1, 2, 3).myMap { 2 * it }.collect {
println("Collecting $it")
}
}
//sampleEnd
Flow 빌더 안에서 일시 중단 함수 호출하기
시퀀스와 달리, flow() 빌더 함수 안에서 일시 중단 함수를 호출할 수 있어요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
suspend fun loadPage(): Int {
delay(100)
return 3
}
suspend fun main() {
flow {
emit(loadPage())
}.collect {
println(it)
// 3
}
}
//sampleEnd
하지만 flow() 빌더 함수는 자신이 실행되는 같은 코루틴 컨텍스트에서 값을 방출해야 해요. 그 블록 안에서 emit()을 호출하는 다른 코루틴을 시작할 수 없고, withContext()로 코루틴 컨텍스트를 바꿀 수도 없어요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
suspend fun main() {
// 이 코드는 예외로 실패해요!
flow {
// withContext()로 코루틴 컨텍스트를 바꿔요
withContext(Dispatchers.IO) {
emit('a')
}
}.collect {
println("This never prints")
}
}
//sampleEnd
이 제약은 flow() 빌더 함수에 적용돼요. 업스트림 Flow가 다른 코루틴 컨텍스트에서 실행되길 원한다면 .flowOn() 연산자로 바꿀 수 있어요. 또는 channelFlow()를 사용해 여러 코루틴에서 값을 방출할 수 있어요.
콜드 Flow의 코루틴 컨텍스트를 .flowOn()으로 바꾸기
기본적으로 콜드 Flow는 소비자와 같은 코루틴 컨텍스트에서 실행돼요. Flow가 다른 코루틴 컨텍스트에서 실행되길 원하면 .flowOn() 연산자를 사용해요. 이 연산자는 컨텍스트를 보존해요. 다운스트림 Flow는 호출자의 컨텍스트에 두면서 오직 업스트림 Flow의 코루틴 컨텍스트만 바꾸죠.
다음은 한 코루틴 컨텍스트에서 값을 방출하고 다른 곳에서 수집하는 콜드 Flow의 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
suspend fun main() {
withContext(Dispatchers.Default + CoroutineName("downstream")) {
flow {
val coroutineName = currentCoroutineContext()[CoroutineName]?.name
// .flowOn()으로 적용된 코루틴 컨텍스트에서 방출해요
println("Emitting '1' in $coroutineName")
// Emitting '1' in upstream
emit(1)
// 업스트림 Flow의 코루틴 컨텍스트를 바꿔요
}.flowOn(Dispatchers.IO + CoroutineName("upstream"))
.collect {
val coroutineName = currentCoroutineContext()[CoroutineName]?.name
// 호출자의 코루틴 컨텍스트에서 수집해요
println("Collecting '$it' in $coroutineName")
// Collecting '1' in downstream
}
}
}
//sampleEnd
Flow에서 예외 처리하기
생산자와 소비자 둘 다 예외를 던질 수 있어요. Flow 수집 중에 예외를 처리하지 않으면, 그 예외는 소비자에서 업스트림으로 전파되어 collect() 함수의 호출자에게 던져져요.
collect() 함수를 try-catch 블록으로 감싸서 이런 예외를 처리할 수 있어요. 예를 들어:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
class MyFlowException(message: String) : Exception(message)
//sampleStart
suspend fun main() {
val myFlow = flow {
try {
// emit() 함수가 collect()로 전달된 람다를 호출해요
emit('a')
} catch (e: MyFlowException) {
println("Collector threw $e")
// 다운스트림 예외를 다시 던져요
throw e
}
}
// Flow 수집을 try-catch로 감싸요
try {
myFlow.collect {
// collect() 람다에서 예외를 던져요
throw MyFlowException("Can't process '$it'!")
}
} catch (e: MyFlowException) {
println("Flow collection failed with $e")
// 호출자에게 예외를 다시 던져요
throw e
}
}
//sampleEnd
이 예시에서 소비자는 emit() 함수에서 값을 받을 때 예외를 던져요. flow() 빌더 함수는 이 다운스트림 예외를 잡아요.
Flow 빌더 함수 안에서 소비자가 던진 예외를 잡았을 때는, 그 예외를 다시 던져야 해요. 그래야 예외 투명성(exception transparency)이 유지되고 collect()의 호출자가 예외를 처리할 수 있어요.
.catch() 연산자로 업스트림 예외 처리하기
예외가 소비자에게 도달하기 전에 처리하려면 .catch() 연산자를 사용해요.
.catch() 연산자로 업스트림 Flow의 예외를 처리할 수 있어요. 예를 들어 emit() 함수로 대체(fallback) 값을 다운스트림으로 방출하는 방식으로요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
suspend fun main() {
flow {
emit("a")
emit("b")
// 업스트림 Flow에서 예외를 던져요
throw UnsupportedOperationException(
"I am tired of listing letters"
)
}.catch { upstreamException ->
println("Upstream completed with $upstreamException!")
// 대체 값을 다운스트림으로 방출해요
emit("Upstream terminated with an exception!")
}.collect {
println("Got '$it'")
}
}
//sampleEnd
이 예시에서 업스트림 Flow는 예외를 던지기 전에 값을 방출해요. .catch() 연산자가 예외를 처리하고 대체 값으로 "Upstream terminated with an exception!"을 방출해요.
Flow가 정상 동작 중 일부 예외를 던질 것으로 예상된다면, 복구 가능한(recoverable) 예외는 .catch()에서 처리하고 예상치 못한 예외는 다시 던지세요.
데이터를 로딩하고 진행률을 보고하는 Flow의 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
sealed interface LoadingState {
sealed interface Terminal: LoadingState
object Started: LoadingState
data class Percentage(val percents: Int): LoadingState
object Failed: Terminal
object Done: Terminal
}
fun loadBlob(url: String) = flow {
emit(LoadingState.Started)
val failureChancePerStep = 1 - java.lang.Math.pow(0.99, 10.0)
repeat(10) { step ->
if (Random.nextDouble() < failureChancePerStep)
throw IOException("Failed to load!")
emit(LoadingState.Percentage((step + 1) * 10))
delay(10.milliseconds)
}
emit(LoadingState.Done)
}.catch { e ->
println("Loading data failed with $e")
if (e is IOException) {
// 예상된 예외를 처리해요
emit(LoadingState.Failed)
} else {
// 예상치 못한 예외를 다시 던져서 collect()가 그 예외로 실패하게 해요
throw e
}
}
suspend fun main() {
loadBlob("https://example.com/").collect {
println("Got '$it'")
}
}
//sampleEnd
이 예시에서 로딩이 예상된 예외로 실패하면, .catch() 연산자가 emit() 함수로 대체 상태를 방출해요. 예상치 못한 예외는 .catch() 연산자에서 다시 던져요. 이렇게 하면 collect() 함수의 호출자가 Flow가 처리하지 않는 예외를 받을 수 있어요.
.catch() 연산자는 소비자가 던진 예외는 처리하지 않아요. collect()에 전달된 람다가 예외를 던지면, collect() 함수 주위에 try-catch 블록으로 처리하세요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
suspend fun main() {
val myFlow = flow {
for (char in listOf('a', 'o', '5', 'c')) {
try {
emit(char)
} catch (e: IllegalArgumentException) {
println("Collector doesn't support character '$char': $e")
// 다운스트림 예외를 다시 던져요
throw e
}
}
}.catch { e ->
// 예외가 다운스트림에서 발생하므로 실행되지 않아요
println("Upstream threw an exception: $e")
}
try {
myFlow.collect {
require(!it.isDigit()) { "Digits are not allowed!" }
}
} catch (e: IllegalArgumentException) {
// collect() 람다의 예외를 처리해요
println("Flow collection failed with $e")
}
}
collect() 람다는 .catch() 다음에 실행되므로, .catch()로는 그 람다에서 발생한 예외를 처리할 수 없어요. .catch()로 방출된 값마다 실행되는 코드의 예외를 처리하려면, 그 코드를 .catch() 앞의 .onEach()에 두세요.
.onEach() 연산자는 각 값이 다운스트림으로 방출되기 전에 그 람다를 실행해요. .catch()가 .onEach()의 예외를 처리하면, Flow는 완료되고 다음 값을 방출하지 않아요.
예시를 볼게요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
suspend fun main() {
flowOf('a', 'o', '5', 'c')
// 각 값이 다운스트림으로 방출되기 전에 실행돼요
.onEach {
require(!it.isDigit()) { "Digits are not allowed!" }
println("Got '$it'")
}
.catch { e ->
println("Caught an exception: $e")
}
.collect()
}
이 예시에서 .onEach() 연산자는 .catch()보다 업스트림에서 실행되므로, '5'에서 require() 검사가 실패하면 .catch() 연산자가 그 예외를 처리해요.
예외 후 업스트림 Flow 다시 시작하기
네트워크 요청이 연결을 잃는 경우처럼, 일부 연산은 일시적으로 실패할 수 있어요. 이런 경우 .retry() 연산자로 예외 후 업스트림 Flow를 다시 시작할 수 있어요.
.retry() 연산자는 예외를 받고, 그 람다가 true를 반환하면 지정된 재시도 횟수만큼 수집을 다시 시작해요. 예를 들어 .retry(3)은 첫 번째 실패 시도 후 업스트림 Flow를 최대 세 번 재시도해요.
람다가 false를 반환하면 .retry()는 재시도를 멈추고 예외를 다시 던져요.
참고: 재시도 로직을 더 세밀하게 제어하려면
.retryWhen()연산자를 사용하세요..retry()처럼 예외를 받지만, 현재 시도 번호도 받고 재시도 전에 값을 방출할 수 있어요.
IOException 후 로딩을 최대 세 번 재시도하는 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
//sampleStart
sealed interface LoadingState {
sealed interface Terminal: LoadingState
object Started: LoadingState
data class Percentage(val percents: Int): LoadingState
object Failed: Terminal
object Done: Terminal
}
fun loadBlob(url: String) = flow {
emit(LoadingState.Started)
val failureChancePerStep = 1 - java.lang.Math.pow(0.99, 10.0)
repeat(10) { step ->
if (Random.nextDouble() < failureChancePerStep)
throw IOException("Failed to load!")
emit(LoadingState.Percentage((step + 1) * 10))
delay(10.milliseconds)
}
emit(LoadingState.Done)
}.retry(3) { e ->
if (e is IOException) {
// 예상된 오류예요
// 재시도 전에 1초를 기다려요
delay(1.seconds)
true
} else {
// 재시도를 멈추고 예상치 못한 예외를 다시 던져요
false
}
}
suspend fun main() {
loadBlob("https://example.org/").collect {
println("Got $it")
}
}
//sampleEnd
Flow 취소
Flow 취소는 요청이 타임아웃됐을 때처럼 더 이상 결과가 필요 없어지면 수집을 멈춰요.
Flow 수집은 collect() 함수를 호출하는 코루틴에 묶여 있어요. 그 코루틴이 취소되면 수집이 멈추고 업스트림 Flow도 함께 취소돼요.
Flow 수집을 취소하려면 수집하는 코루틴의 Job에서 cancel() 함수를 호출하면 돼요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
val myFlow = flow {
var i = 0
try {
while (true) {
println("Emitting $i")
emit(i)
println("Emitted $i")
++i
delay(10.milliseconds)
}
} catch (e: Throwable) {
println("Upstream finished with $e")
throw e
}
}
suspend fun main() {
coroutineScope {
val job = launch {
try {
myFlow.collect {
println("Processing $it")
delay(5.milliseconds)
}
} catch (e: Throwable) {
println("Collection finished with $e")
throw e
}
}
delay(100.milliseconds)
// Flow를 수집하는 코루틴을 취소해요
job.cancel()
}
}
//sampleEnd
소비자는 수집 코루틴이 활성 상태인 동안에도 업스트림 Flow를 취소할 수 있어요. 이렇게 하려면 소비자 쪽에서 CancellationException을 던지면 돼요.
.take() 연산자는 이 동작을 이용해 고정된 개수의 값 뒤에 수집을 멈춰요. 예를 들어 .take(3)은 업스트림 Flow의 처음 세 값만 수집한 뒤 그것을 취소해요.
.take() 연산자의 단순화된 버전을 사용한 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
// 기본 .take() 연산자의 단순화된 버전을 정의해요
fun <T> Flow<T>.myTake(count: Int): Flow<T> = flow {
require(count > 0)
val cancellationException = CancellationException()
var elementsRemaining = count
try {
[email protected] {
emit(it)
--elementsRemaining
if (elementsRemaining == 0) {
// 요청된 값 개수 뒤에 업스트림 Flow를 취소해요
throw cancellationException
}
}
} catch (e: Throwable) {
if (e === cancellationException) {
// 업스트림 Flow를 취소하는 데 사용된 CancellationException을 처리해요
// .myTake()에서 설정된 값 개수 뒤에 Flow를 완료해요
} else {
// 예상치 못한 예외를 다시 던져요
throw e
}
}
}
suspend fun main() {
(0..1000).asFlow().myTake(3).collect {
println("Got $it")
}
}
//sampleEnd
이 예시에서 .myTake() 함수는 요청된 모든 값이 방출될 때까지 업스트림 Flow에서 값을 방출해요. 그다음 CancellationException을 던져 업스트림 Flow를 취소해요.
channelFlow()로 값을 동시에 방출하기
flow() 빌더 함수는 한 코루틴에서 값을 방출하는 Flow에는 단순하고 효율적이에요. 여러 코루틴에서 같은 Flow로 값을 동시에 방출하고 싶다면 channelFlow() 빌더 함수를 사용해요. 여러 소스에서 데이터를 로딩하는 것처럼 결과를 점진적으로 보고하는 동시 작업에 쓸 수 있어요.
channelFlow() 빌더 함수는 채널(channel)을 사용해 여러 코루틴에서 값을 보내는 콜드 Flow를 만들어요. 빌더 안에서 값을 만들 때 emit() 함수 대신 send() 함수를 사용해요.
두 Flow를 동시에 수집하고 단순화된 .merge() 연산자 버전으로 그 값을 다시 방출하는 channelFlow() 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
// 기본 .merge() 연산자의 단순화된 버전을 정의해요
fun <T> Flow<T>.myMerge(other: Flow<T>): Flow<T> = channelFlow {
// 여기서는 CoroutineScope와 SendChannel을 리시버로 쓸 수 있어요
// 리시버 Flow를 수집하는 코루틴을 실행해요
launch {
// 리시버 Flow를 수집해요
[email protected] {
send(it)
}
}
launch {
// other Flow를 수집하는 코루틴을 실행해요
other.collect {
// SendChannel.send를 호출해요
send(it)
}
}
}
suspend fun main() {
val flow1 = (0..3).asFlow().onEach { delay(20.milliseconds) }
val flow2 = (6..9).asFlow().onEach { delay(50.milliseconds) }
flow1.myMerge(flow2).collect { println(it) }
}
//sampleEnd
channelFlow() 빌더 함수는 버퍼링된 채널을 사용해서, 버퍼가 가득 찰 때까지 생산자가 소비자보다 앞서 값을 보낼 수 있게 해요. 기본적으로 버퍼는 최대 64개의 값을 담을 수 있어요. 버퍼가 가득 차면 버퍼에 빈 공간이 생길 때까지 생산자가 일시 중단돼요.
.buffer() 연산자로 버퍼 용량을 바꿀 수 있어요. 예를 들어 .buffer(12)는 생산자가 소비자보다 최대 12개의 값을 앞서 보낼 수 있게 하고, .buffer(0)은 버퍼를 제거해서 소비자가 받을 수 있을 때만 각 값을 보내요.
예시를 볼게요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
suspend fun main() {
val oneHundredNumbers = channelFlow {
repeat(100) {
println("Sending $it")
send(it)
}
}
// 기본 버퍼 용량을 사용해요
oneHundredNumbers.collect {
println("Processing $it")
delay(10.milliseconds)
}
// 버퍼를 제거해서 보내기와 처리가 처음부터 번갈아 일어나게 해요
oneHundredNumbers.buffer(0).collect {
println("Processing $it")
delay(10.milliseconds)
}
}
//sampleEnd
이 예시에서 oneHundredNumbers Flow는 기본 버퍼 용량을 사용하고, oneHundredNumbers.buffer(0) Flow는 버퍼가 없어요.
기본 버퍼 용량일 때 생산자는 버퍼가 가득 찰 때까지 값을 빠르게 보내요. 그 후에는 버퍼에 빈 공간이 생길 때까지 send()가 일시 중단되므로, Sending과 Processing 메시지가 번갈아 나타나기 시작해요.
.buffer(0)에서는 각 send() 호출이 소비자가 값을 받을 수 있을 때까지 기다려서, Sending과 Processing이 처음부터 번갈아 나타나요.
핫 Flow
핫 Flow는 소비자와 무관하게 값을 방출하는 공유 스트림이에요. 활성 소비자가 없어도 값을 계속 방출하고, 여러 소비자가 새 실행을 시작하는 대신 이미 활성화된 스트림에서 같은 방출을 수집할 수 있어요.
핫 Flow의 소비자를 구독자(subscriber)라고 불러요.
애플리케이션의 여러 부분이 같은 업데이트 스트림에 반응해야 할 때 핫 Flow를 쓸 수 있어요. 채팅 메시지가 들어오는 것, 사용자 액션, UI 상태 변화 같은 거예요.
Kotlin은 두 가지 핫 Flow 타입을 제공해요.
SharedFlow는 여러 구독자에게 값을 브로드캐스트해요. 시간에 따라 발생하는 이벤트(메시지나 알림 같은 것)를 브로드캐스트할 때 사용해요.StateFlow는 항상 최신 상태 값을 유지하는 특수한SharedFlow예요. 시간에 따라 변하는 상태(UI 상태 같은 것)를 나타낼 때 사용해요.
SharedFlow 만들기
SharedFlow는 시간에 따라 발생하는 방출된 값을 구독자에게 브로드캐스트하는 핫 Flow예요.
MutableSharedFlow() 함수로 SharedFlow를 만들 수 있어요.
MutableSharedFlow는 값을 방출하는 함수를 노출해요. 그것을 직접 노출하면 클래스 밖의 코드가 Flow에 값을 방출할 수 있게 돼요.
이를 막으려면 가변 Flow를 private 백킹 프로퍼티(backing property)에 저장하고 .asSharedFlow() 함수로 읽기 전용 SharedFlow를 노출해요. 구독자에게 값을 방출하려면 MutableSharedFlow에서 emit() 함수를 사용해요.
data class Message(
val senderId: Int,
val time: Instant,
val text: String,
)
class Chatroom {
// SharedFlow를 private 백킹 프로퍼티에 저장해요
private val _messages = MutableSharedFlow<Message>()
// 구독자에게 읽기 전용 SharedFlow를 노출해요
val messages: SharedFlow<Message>
get() = _messages.asSharedFlow()
suspend fun sendMessageToEveryone(message: Message) {
// 구독자에게 메시지를 방출해요
_messages.emit(message)
}
}
콜드 Flow 때처럼 collect() 함수로 SharedFlow에서 값을 수집할 수 있어요.
또한 SharedFlow가 이미 방출된 값을 새 구독자에게 즉시 재생(replay)하도록 설정할 수 있어요. 리플레이 캐시는 작은 기록 버퍼처럼 동작하며, 이전 방출의 고정된 개수를 저장해요.
새 구독자가 받을 이전 방출 개수를 설정하려면 MutableSharedFlow()의 replay 파라미터를 사용해요.
// 새 구독자가 구독 시 받을 이미 방출된 메시지 개수를 설정해요
const val MESSAGES_TO_REMEMBER = 10
class Chatroom {
private val _messages = MutableSharedFlow<Message>(
// 새 구독자에게 마지막으로 방출된 메시지의 설정된 개수를 재생해요
replay = MESSAGES_TO_REMEMBER
)
val messages: SharedFlow<Message>
get() = _messages.asSharedFlow()
suspend fun sendMessageToEveryone(message: Message) {
// messages Flow의 구독자에게 메시지를 방출해요
_messages.emit(message)
}
}
핫 Flow 수집은 스스로 완료되지 않으므로, 더 이상 필요 없어지면 수집하는 코루틴을 취소해야 해요.
참고: 핫 Flow에는 닫기나 취소 연산이 없어요. 수집을 취소하는 것은 해당 구독자가 수집하는 것만 멈출 뿐이에요. 새 방출을 멈추려면 핫 Flow에 값을 만드는 코루틴이나 스코프를 취소하세요.
SharedFlow로 채팅방을 모델링하는 예를 볼게요. 각 새 메시지를 활성 구독자에게 보내고 최근 메시지를 나중에 들어온 구독자에게 재생해요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.*
data class Message(
val senderId: Int,
val time: Instant,
val text: String,
)
// 새 구독자가 구독 시 받을 이미 방출된 메시지 개수를 설정해요
const val MESSAGES_TO_REMEMBER = 10
class Chatroom {
// SharedFlow를 private 백킹 프로퍼티에 저장해요
private val _messages = MutableSharedFlow<Message>(
// 새 구독자에게 마지막으로 방출된 메시지의 설정된 개수를 재생해요
replay = MESSAGES_TO_REMEMBER
)
// 구독자에게 읽기 전용 SharedFlow를 노출해요
val messages: SharedFlow<Message>
get() = _messages.asSharedFlow()
// 구독자에게 메시지를 방출해요
suspend fun sendMessageToEveryone(message: Message) {
_messages.emit(message)
}
}
suspend fun main() {
val nUsers = 3
val chatroom = Chatroom()
withContext(Dispatchers.Default) {
// 각 사용자별 메시지 리더를 시작해요
val messageReaders = List(nUsers) { userId ->
// 메시지가 방출되기 전에 수집을 시작해요
launch(start = CoroutineStart.UNDISPATCHED) {
chatroom.messages.collect { message ->
println("User $userId received $message")
}
}
}
// 각 사용자로부터 인사말을 보내요
repeat(nUsers) { userId ->
chatroom.sendMessageToEveryone(
Message(
userId,
Clock.System.now(),
"Hello from $userId!"
)
)
}
// 사람들이 대화할 시간이 충분하도록 지연해요
delay(100.milliseconds)
// SharedFlow 수집은 스스로 끝나지 않으므로 리더를 취소해요
messageReaders.forEach { it.cancel() }
}
}
이 예시에서 CoroutineStart.UNDISPATCHED는 각 수집 코루틴을 즉시 시작해요. 이렇게 하면 각 코루틴이 collect()에 도달해 messages를 구독하고, sendMessageToEveryone()이 메시지를 방출하기 전에 일시 중단되는 것이 보장돼요. 이것이 없으면 수집 코루틴이 더 늦게 시작되어, 리플레이 캐시가 너무 작을 때 이전 방출을 놓칠 수 있어요.
명시적 백킹 필드로 핫 Flow 노출하기
명시적 백킹 필드(explicit backing fields)로 클래스 안에 가변 백킹 필드를 유지하면서 읽기 전용 SharedFlow를 노출할 수 있어요.
명시적 백킹 필드는 field 선언에서 구현 타입을 정의해요. 클래스 안에서 컴파일러가 프로퍼티를 백킹 필드 타입으로 스마트 캐스트하므로, 별도의 private 백킹 프로퍼티 없이 emit() 함수를 호출할 수 있어요.
참고: 명시적 백킹 필드는
.asSharedFlow()가 제공하는 읽기 전용 래퍼를 만들지 않아요. 노출된 Flow를 다운캐스트하는 것이 문제가 되지 않을 때만 이 패턴을 사용하세요.
예시를 볼게요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Clock
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.ExperimentalTime
import kotlin.time.Instant
data class Message(
val senderId: Int,
val time: Instant,
val text: String,
)
const val MESSAGES_TO_REMEMBER = 10
//sampleStart
class Chatroom {
// 가변 백킹 필드로 읽기 전용 SharedFlow를 노출해요
val messages: SharedFlow<Message>
field = MutableSharedFlow<Message>(
replay = MESSAGES_TO_REMEMBER
)
suspend fun sendMessageToEveryone(message: Message) {
// Chatroom 안의 가변 백킹 필드를 통해 방출해요
messages.emit(message)
}
}
//sampleEnd
suspend fun main() {
val nUsers = 3
val chatroom = Chatroom()
withContext(Dispatchers.Default) {
val messageReaders = List(nUsers) { userId ->
launch(start = CoroutineStart.UNDISPATCHED) {
chatroom.messages.collect { message ->
println("User $userId received $message")
}
}
}
repeat(nUsers) { userId ->
chatroom.sendMessageToEveryone(
Message(
senderId = userId,
time = Clock.System.now(),
text = "Hello from $userId!"
)
)
}
delay(100.milliseconds)
messageReaders.forEach { it.cancel() }
}
}
StateFlow 만들기
StateFlow는 단일 상태 값을 저장하고, 그 값이 새 값으로 교체될 때 업데이트를 방출하는 핫 Flow예요. 새 구독자는 수집을 시작하자마자 현재 값을 받고, 그다음 상태가 업데이트될 때마다 새 값을 받아요.
로딩 진행률, UI 상태, 객체의 상태처럼 시간에 따라 변하는 상태를 나타낼 때 StateFlow를 쓸 수 있어요.
StateFlow를 만들려면 초기 값으로 MutableStateFlow() 함수를 사용해요.
// LoadingState.Started를 초기 값으로 MutableStateFlow를 만들어요
val result = MutableStateFlow<LoadingState>(LoadingState.Started)
현재 상태를 설정하려면 value 프로퍼티를 사용해요.
fun loadBlob(url: String): StateFlow<LoadingState> {
val result = MutableStateFlow<LoadingState>(LoadingState.Started)
DownloadManager.startLoading(
url,
onPercentageLoaded = { percentage ->
// 현재 상태를 최신 진행률로 교체해요
result.value = LoadingState.Percentage(percentage)
},
onCompletion = {
// 현재 상태를 완료 상태로 교체해요
result.value = LoadingState.Done
},
onFailure = {
// 현재 상태를 실패 상태로 교체해요
result.value = LoadingState.Failed
}
)
}
참고:
value설정은 스레드 안전하고 현재 상태를 교체하지만, 이전 값에 기반해value를 갱신하는 것은 원자적이지 않아요. 새 상태가 이전 상태에 의존할 때는.update()를 대신 사용하세요.
MutableSharedFlow와 비슷하게, MutableStateFlow는 업데이트를 방출하는 API를 노출해요. 직접 노출하면 그걸 받는 어떤 코드든 MutableStateFlow로 다운캐스트해서 상태를 갱신할 수 있어요.
이를 막으려면 .asStateFlow() 함수로 가변 Flow를 읽기 전용 StateFlow로 노출하세요.
fun loadBlob(url: String): StateFlow<LoadingState> {
val result = MutableStateFlow<LoadingState>(LoadingState.Started)
DownloadManager.startLoading(
url,
onPercentageLoaded = { percentage ->
result.value = LoadingState.Percentage(percentage)
},
onCompletion = {
result.value = LoadingState.Done
},
onFailure = {
result.value = LoadingState.Failed
}
)
// 로딩 상태를 읽기 전용 StateFlow로 노출해요
return result.asStateFlow()
}
콜백 기반 API에서 로딩 진행률을 StateFlow로 보고하는 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random
//sampleStart
sealed interface LoadingState {
sealed interface Terminal: LoadingState
object Started: LoadingState
data class Percentage(val percents: Int): LoadingState
object Failed: Terminal
object Done: Terminal
}
fun loadBlob(url: String): StateFlow<LoadingState> {
// 초기 로딩 상태로 가변 StateFlow를 만들어요
val result = MutableStateFlow<LoadingState>(LoadingState.Started)
DownloadManager.startLoading(
url,
onPercentageLoaded = { percentage ->
// 현재 상태를 최신 진행률로 교체해요
result.value = LoadingState.Percentage(percentage)
},
onCompletion = {
// 현재 상태를 완료 상태로 교체해요
result.value = LoadingState.Done
},
onFailure = {
// 현재 상태를 실패 상태로 교체해요
result.value = LoadingState.Failed
}
)
// 로딩 상태를 읽기 전용 StateFlow로 노출해요
return result.asStateFlow()
}
// 데이터를 비동기로 다운로드하는 콜백 기반 API를 정의해요
object DownloadManager {
// url 로딩을 비동기로 시작해요
fun startLoading(
url: String,
onPercentageLoaded: (Int) -> Unit,
onCompletion: () -> Unit,
onFailure: (Throwable) -> Unit
) {
// GlobalScope는 이 예시를 자체 완결적으로 유지하기 위한 설명용으로만 사용해요
GlobalScope.launch {
val failureChancePerStep = 1 - java.lang.Math.pow(0.99, 10.0)
repeat(10) { step ->
if (Random.nextDouble() < failureChancePerStep) {
onFailure(IOException("Failed to load!"))
return@launch
}
onPercentageLoaded((step + 1) * 10)
delay(10.milliseconds)
}
onCompletion()
}
}
}
suspend fun main() {
loadBlob("https://example.com/").onEach { state ->
when (state) {
is LoadingState.Started -> {
// 진행률 업데이트를 기다려요
}
is LoadingState.Percentage ->
println("Loaded ${state.percents}...")
is LoadingState.Failed ->
println("Loading failed.")
is LoadingState.Done ->
println("Finished loading!")
}
}.takeWhile { it !is LoadingState.Terminal }.collect()
}
//sampleEnd
참고: 이 예시는 콜백 기반 API를 짧게 유지하려고
GlobalScope를 사용했어요. 실제 애플리케이션에서는 작업을 시작하는 함수(이 예시의startLoading()같은)에CoroutineScope를 전달하고, 호출자가 더 이상 필요 없어지면 작업을 취소할 수 있도록 그 스코프에서 코루틴을 실행하세요.
StateFlow는 핫 Flow이므로 수집이 스스로 끝나지 않아요. 이 예시에서 .takeWhile() 연산자는 로딩이 종료 상태에 도달하면 수집을 멈춰요.
StateFlow는 새 값이 현재 값과 다를 때만 업데이트를 방출해요.
참고:
StateFlow에 가변 객체를 저장하지 마세요. 객체 자체를 변경하는 것은 현재 값을 교체하지 않으므로 소비자가 업데이트를 받지 못해요.
현재 상태에서 새 상태를 계산해서 StateFlow를 갱신할 수도 있어요. 이런 갱신에는 .update() 함수를 사용해요. .update() 함수는 값을 원자적으로 갱신해요. 여러 코루틴이 같은 MutableStateFlow를 갱신할 때 도움이 돼요.
참고: 공유 값을 갱신하기만 하면 되고 시간에 따른 상태 변화를 관찰할 필요가 없다면,
AtomicInt나AtomicReference같은 Kotlin Atomics API를 사용하세요.
좋아요(likes)가 StateFlow에 저장되고 각 새 상태가 이전 상태에서 계산되는 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random
//sampleStart
class Post(val id: Long) {
// 현재 좋아요 수를 StateFlow로 저장해요
private val _numberOfLikes = MutableStateFlow<Int>(
// 초기 좋아요 수를 설정해요
0
)
// 현재 좋아요 수를 담은 읽기 전용 StateFlow를 노출해요
val numberOfLikes: StateFlow<Int>
get() = _numberOfLikes.asStateFlow()
// 좋아요를 더해요
fun like() {
// 동시·다중 스레드 호출을 위해 좋아요 수를 원자적으로 증가시켜요
_numberOfLikes.update { it + 1 }
}
}
suspend fun drawUpdatedNumberOfLikes(likes: Int) {
// 최신 좋아요 수를 표시해요
println("${Clock.System.now()}: the number of likes is $likes")
}
suspend fun main() {
withContext(Dispatchers.Default) {
val post = Post(15)
val notifyingJob = launch {
post.numberOfLikes.collect {
drawUpdatedNumberOfLikes(it)
}
}
// 게시물에 좋아요를 누르는 사용자를 흉내 내요
coroutineScope {
repeat(10) {
launch {
delay(Random.nextInt(100).milliseconds)
post.like()
}
}
}
// 모든 시뮬레이션 사용자가 끝난 뒤 수집을 취소해요
notifyingJob.cancelAndJoin()
}
}
//sampleEnd
이 예시에서 .update() 함수는 좋아요 수를 원자적으로 증가시켜요. 이렇게 하면 여러 코루틴이 동시에 like() 함수를 호출할 때 발생하는 갱신 손실(lost updates)을 막아요.
누적된 상태를 StateFlow에 저장하기
때로 구독자가 최신 방출 값만이 아니라 이전 방출 전체의 결과를 받길 원할 수 있어요.
예를 들어 채팅방은 메시지 기록을 하나의 상태 값으로 유지할 수 있어요. 새 사용자가 채팅방에 들어오면 먼저 현재 메시지 기록을 받아요. 그다음 새 메시지가 도착할 때 계속 업데이트를 받아요.
이 동작을 StateFlow로 모델링할 수 있어요.
그러려면 각 채팅 메시지를 SharedFlow<Message>로 개별 이벤트로 브로드캐스트하는 대신, 전체 메시지 기록을 StateFlow<List<Message>>의 현재 값으로 저장하면 돼요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random
//sampleStart
data class Message(
val senderId: Int,
val time: Instant,
val text: String,
)
class Chatroom {
// 전체 메시지 기록을 저장해요
private val _messageHistory = MutableStateFlow<List<Message>>(emptyList())
// 현재 메시지 기록을 담은 읽기 전용 StateFlow를 노출해요
val messageHistory: StateFlow<List<Message>>
get() = _messageHistory.asStateFlow()
// messageHistory Flow의 모든 구독자에게 메시지를 보내요
suspend fun sendMessageToEveryone(message: Message) {
// 새 메시지를 현재 기록에 원자적으로 더해요
_messageHistory.update {
it + message
}
}
}
suspend fun main() {
val nUsers = 3
val chatroom = Chatroom()
withContext(Dispatchers.Default) {
// 각 사용자별 메시지 리더를 시작해요
val messageReaders = List(nUsers) { userId ->
launch(start = CoroutineStart.UNDISPATCHED) {
chatroom.messageHistory.collect { currentHistory ->
println("User $userId sees the history as $currentHistory")
}
}
}
// 각 사용자로부터 인사말을 보내요
repeat(nUsers) { userId ->
chatroom.sendMessageToEveryone(
Message(
userId,
Clock.System.now(),
"Hello from $userId!"
)
)
}
// 사용자가 업데이트를 받을 시간이 충분하도록 지연해요
delay(100.milliseconds)
// StateFlow 수집은 스스로 끝나지 않으므로 리더를 취소해요
messageReaders.forEach { it.cancel() }
}
}
//sampleEnd
이 예시에서 messageHistory는 이전 메시지 전체 목록을 현재 상태로 저장해요. 새 메시지가 보내지면 .update() 함수가 이전 기록에서 새 리스트를 만들고 새 메시지를 원자적으로 더해요.
참고: 새 컬렉션을 만들어 불변 컬렉션을 갱신하는 것은 컬렉션이 커질수록 더 많은 시간이 걸릴 수 있어요. Experimental
kotlinx.collections.immutable라이브러리로 영속 컬렉션을 만들어 불변 컬렉션 갱신을 더 효율적으로 할 수 있어요.
messageHistory는 StateFlow이므로 구독자는 수집을 시작할 때 현재 메시지 기록을 받아요. 그 후 메시지가 보내질 때마다 새 리스트를 받는데, 그때마다 채팅 기록이 바뀌어요.
콜드 Flow를 핫 Flow로 변환하기
콜드 Flow는 소비자마다 업스트림 연산을 따로 실행해요. 여러 구독자가 같은 업스트림 수집의 방출이 필요할 때, 콜드 Flow를 그 수집을 구독자와 공유하는 핫 Flow로 변환할 수 있어요.
다음의 .shareIn() 단순화 버전은 이 아이디어를 보여 줘요. 콜드 Flow를 한 번 수집하고, 그 값을 MutableSharedFlow에 방출하고, 읽기 전용 SharedFlow로 노출하죠.
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.*
fun <T> Flow<T>.simpleShareIn(scope: CoroutineScope): SharedFlow<T> {
val sharedFlow = MutableSharedFlow<T>()
scope.launch {
[email protected] {
sharedFlow.emit(it)
}
}
return sharedFlow.asSharedFlow()
}
suspend fun main() {
}
이 예시에서 simpleShareIn()은 제공된 스코프에서 새 코루틴을 시작해요. 업스트림 Flow 수집을 멈추려면 수집 코루틴을 실행하는 스코프를 취소하면 돼요.
업스트림 Flow가 예외를 던지면 이 수집 코루틴이 실패해요. 공유 전에 .catch()나 .retry() 같은 연산자를 사용해 수집 코루틴이 실패하기 전에 업스트림 예외를 처리하세요.
내장 .shareIn() 함수는 직접 MutableSharedFlow를 만들지 않고도 이 패턴을 제공해요. 또한 업스트림 수집이 시작·멈출 때를 제어하고, 새 구독자가 받을 이전 방출 개수를 설정하는 옵션도 더해 줘요.
내장 .shareIn() 함수를 사용하려면 다음 인자를 제공하세요.
- 업스트림 Flow가 수집되는 코루틴 스코프
- 업스트림 수집의 시작·멈춤을 제어하는
SharingStarted전략. 예를 들어SharingStarted.Eagerly는 어떤 구독자도 수집을 시작하기 전에 제공된 스코프에서 업스트림 수집을 즉시 시작해요. - 새 구독자가 받을 이전 방출 개수를 제어하는 선택적
replay값
.shareIn() 함수는 제공된 코루틴 스코프에서 업스트림 Flow를 수집하고 그 방출을 구독자에게 브로드캐스트해요.
.shareIn()으로 콜드 Flow를 여러 구독자와 직렬화된 채팅 메시지를 공유하는 핫 Flow로 변환하는 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random
//sampleStart
data class Message(
val senderId: Int,
val time: Instant,
val text: String,
)
class Chatroom {
// 메시지 Flow를 저장해요
private val _messages = MutableSharedFlow<Message>()
// 방출된 메시지가 담긴 읽기 전용 SharedFlow를 노출해요
// 새 구독자는 이미 방출된 메시지를 받지 않아요
val messages: SharedFlow<Message>
get() = _messages.asSharedFlow()
// messages Flow의 모든 구독자에게 메시지를 보내요
suspend fun sendMessageToEveryone(message: Message) {
_messages.emit(message)
}
}
suspend fun main() {
val nUsers = 3
val chatroom = Chatroom()
withContext(Dispatchers.Default) {
// 현재 실행 중인 코루틴의 자식 스코프를 만들어요
val derivedFlowsScope = CoroutineScope(
currentCoroutineContext() + Job(currentCoroutineContext()[Job])
)
// 구독자 간에 직렬화된 메시지를 공유해요
val serializedMessages: SharedFlow<String> =
chatroom
.messages
.map {
// 공유 Flow를 위해 각 메시지를 한 번 직렬화해요
"senderId: ${it.senderId}, time: ${it.time}, text: " +
Base64.Default.encode(it.text.encodeToByteArray())
}
.shareIn(
// 이 스코프에서 공유 코루틴을 시작해요
// .map()을 포함한 업스트림 Flow가 그 코루틴에서 실행돼요
derivedFlowsScope,
// 첫 구독자가 나타나기 전에 업스트림 Flow 수집을 즉시 시작해요
SharingStarted.Eagerly,
// 새 구독자에게 이전 직렬화된 메시지를 재생하지 않아요
replay = 0,
)
// 각 사용자별 메시지 리더를 시작해요
val messageReaders = List(nUsers) { userId ->
launch(start = CoroutineStart.UNDISPATCHED) {
serializedMessages.collect { serializedMessage ->
println("User $userId observes the message $serializedMessage")
}
}
}
// 각 사용자로부터 인사말을 보내요
repeat(nUsers) { userId ->
chatroom.sendMessageToEveryone(
Message(
userId,
Clock.System.now(),
"Hello from $userId!"
)
)
}
// 사용자가 업데이트를 받을 시간이 충분하도록 지연해요
delay(100.milliseconds)
// SharedFlow 수집은 스스로 끝나지 않으므로 리더를 취소해요
messageReaders.forEach { it.cancel() }
// 파생 핫 Flow를 실행하는 스코프를 취소해요
derivedFlowsScope.cancel()
}
}
//sampleEnd
이 예시에서 .map() 연산자는 각 메시지를 직렬화하는 콜드 Flow를 만들어요. .shareIn() 함수가 없다면 소비자마다 그 직렬화를 따로 실행해요. .shareIn() 함수는 업스트림 수집 하나를 공유하므로 각 메시지는 한 번 직렬화되고 모든 구독자와 공유돼요.
SharingStarted.Eagerly가 업스트림 수집을 즉시 시작하기 때문에, .shareIn()이 호출되는 즉시 파생 핫 Flow가 chatroom.messages 수집을 시작해요.
비슷하게, 콜드 Flow를 StateFlow로 변환하려면 .stateIn() 함수를 사용해요.
.shareIn()과 달리 .stateIn()은 초기 값이 필요해요. StateFlow는 항상 현재 값이 있어야 하거든요.
예를 들어:
val lastUpdateFlow: StateFlow<Instant?> =
chatroom
.messageHistory
.map { currentHistory -> currentHistory.lastOrNull()?.time }
.stateIn(
// 이 스코프에서 공유 코루틴을 시작해요
// .map()을 포함한 업스트림 Flow가 그 코루틴에서 실행돼요
derivedFlowsScope,
// 첫 구독자가 나타나면 수집을 시작하고
// 마지막 구독자가 사라지면 멈춰요
SharingStarted.WhileSubscribed(),
// 첫 업스트림 방출 전에 초기 상태를 설정해요
null,
)
핫 Flow 취소하기
핫 Flow는 구독자가 취소돼도 멈추지 않아요.
핫 Flow를 수집하는 코루틴을 취소하면 그 구독자만 취소돼요. 핫 Flow는 다른 구독자에게 계속 값을 방출할 수 있고, 값을 만드는 코루틴도 계속 실행될 수 있어요.
핫 Flow 자체에는 취소 연산이 없어요. 핫 Flow를 취소하려면 그것의 값을 만드는 코루틴이나 스코프를 취소하세요.
.shareIn()이나 .stateIn() 확장 함수로 만든 핫 Flow는 공유 코루틴이 취소될 때까지 업스트림 Flow를 계속 수집해요. 업스트림 Flow 수집을 멈추려면 공유 코루틴을 실행하는 스코프를 취소하세요.
참고:
SharingStarted.WhileSubscribed()로 구독자가 없을 때 업스트림 수집을 자동으로 멈출 수도 있어요.
.stateIn()에 전달한 스코프를 취소하면 파생 핫 Flow가 새 값을 수집하는 것을 멈추는 예시예요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random
data class Message(
val senderId: Int,
val time: Instant,
val text: String,
)
class Chatroom {
// 메시지 기록을 저장해요
private val _messageHistory = MutableStateFlow<List<Message>>(emptyList())
// 현재 메시지 기록을 담은 읽기 전용 StateFlow를 노출해요
val messageHistory: StateFlow<List<Message>>
get() = _messageHistory.asStateFlow()
// messageHistory Flow의 모든 구독자에게 메시지를 보내요
suspend fun sendMessageToEveryone(message: Message) {
_messageHistory.update {
it + message
}
}
}
//sampleStart
suspend fun main() {
val chatroom = Chatroom()
withContext(Dispatchers.Default) {
// 현재 실행 중인 코루틴의 자식 스코프를 만들어요
val derivedFlowsScope = CoroutineScope(
currentCoroutineContext() + Job(currentCoroutineContext()[Job])
)
val totalMessages = chatroom.messageHistory
.map { currentHistory ->
currentHistory.size
}.onEach {
println("There are currently $it messages")
}.stateIn(
// 이 스코프에서 공유 코루틴을 시작해요
derivedFlowsScope
)
// messageHistory를 갱신해요
chatroom.sendMessageToEveryone(
Message(0, Clock.System.now(), "We are shutting down soon!")
)
delay(100.milliseconds)
// 파생 핫 Flow를 실행하는 스코프를 취소해요
derivedFlowsScope.cancel()
// messageHistory를 갱신하지만, totalMessages는 더 이상 그 갱신을 받지 않아요
chatroom.sendMessageToEveryone(
Message(0, Clock.System.now(), "We have shut down.")
)
println("Last collected history size: ${totalMessages.value}")
println("Actual history size: ${chatroom.messageHistory.value.size}")
}
}
//sampleEnd
이 예시에서 derivedFlowsScope.cancel() 함수를 호출하면 totalMessages가 messageHistory에서 업데이트를 수집하는 것을 멈춰요.
sendMessageToEveryone() 함수는 여전히 messageHistory를 갱신해요. 그 함수를 호출하는 코루틴이 취소되지 않았기 때문이에요. 결과적으로 totalMessages.value는 마지막에 수집된 크기를 유지하고, chatroom.messageHistory.value.size는 실제 메시지 개수를 보여 줘요.
핫 Flow에서 예외 처리하기
콜드 Flow에서는 .catch() 같은 연산자로 먼저 처리하지 않으면 업스트림 예외가 collect()의 호출자에게 전파돼요.
핫 Flow는 생산자에서 구독자로 예외를 전파하지 않아요. MutableSharedFlow에 방출하거나 MutableStateFlow를 갱신하는 코드가 예외를 던지면, 그 코드를 실행하는 코루틴에서 처리하세요. 구독자가 수집 중 예외를 던지면 수집 코루틴에서 처리하세요.
.shareIn()이나 .stateIn() 확장 함수로 만든 핫 Flow는 공유 코루틴에서 업스트림 Flow를 수집해요. 업스트림 Flow가 예외를 던지면, 그 예외가 공유 코루틴을 취소해요.
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.*
//sampleStart
suspend fun main() {
withContext(Dispatchers.Default) {
launch {
flow<Int> {
error("An upstream failure")
}.stateIn(
this@launch
)
}
}
}
//sampleEnd
실패 후 업스트림 수집을 다시 시작할 수 있어요. 그러려면 .shareIn()이나 .stateIn() 앞에 .retry() 연산자를 두세요.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Duration.Companion.milliseconds
//sampleStart
suspend fun main() {
coroutineScope {
launch {
var currentAttempt = 0
val stateFlow = flow {
delay(10.milliseconds)
if (currentAttempt++ < 5) {
println("An error happened!")
error("An upstream failure")
} else {
println("Success.")
emit(10)
}
}
// 복구 가능한 실패 후 업스트림 Flow를 다시 시작해요
.retry(retries = 5)
.stateIn(
// 이 스코프에서 공유 코루틴을 시작해요
this@launch
)
stateFlow.collect {
println("Observed $it")
// 수집과 공유 코루틴을 취소해요
[email protected]()
}
}
}
}
//sampleEnd
이 예시에서 Flow는 값을 방출하기 전에 다섯 번 실패해요. .retry()가 .stateIn()보다 먼저 실행되므로, 각 업스트림 실패가 공유 코루틴에 도달하기 전에 처리해요.
업스트림 Flow가 10을 방출하면 수집 코루틴이 값을 받고 자기 자신을 취소해요. 같은 코루틴이 공유 코루틴의 부모이기도 하므로, 이것이 파생 핫 Flow를 멈춰요.
더 알아보기 (Learn more)
Flow는 콜드와 핫이라는 두 축으로 이해하는 게 핵심이에요. 데이터를 점진적으로 로딩하거나 상태를 관찰할 때 StateFlow, 일회성 이벤트를 여러 곳에 전달할 때 SharedFlow를 고르는 기준이 자연스럽게 잡혀요. 예외 처리와 취소, 그리고 .shareIn()·.stateIn()으로 흐름을 공유하는 패턴을 직접 코드로 만들어 보면서 익히면, 코루틴 기반 비동기 코드에서 Flow가 얼마나 강력한 도구인지 실감할 수 있어요.