
플로우 연산자에 대한 정리글로 돌아온 개발자 꿈나무 김조현입니다.
이번 글에서는 플로우 연산자가 무엇인지, 어떻게 사용하는지 등에 대해 정리해보겠습니다.
컬렉션을 조작하기 위해 다양한 연산자를 사용하는 것처럼 플로우를 변환할 때도 비슷한 연산자를 쓸 수 있습니다.
시퀀스와 마찬가지로 플로우도 중간 연산자와 최종 연산자를 구분합니다. 중간 연산자는 코드를 실행하지 않고 변경된 플로우를 반환하며, 최종 연산자는 컬렉션, 개별 원소, 계산된 값을 반환하거나 아무 값도 반환하지 않으면서 플로우를 수집하고 실제 코드를 실행합니다.
중간 연산자는 플로우에 적용돼 새로운 플로우를 반환합니다. 플로우는 업스트림과 다운스트림으로 구분할 수 있습니다.
연산자가 적용되는 플로우를 업스트림 플로우라고 하고, 중간 연산자가 반환하는 플로우를 다운스트림 플로우라고 부릅니다. 다운스트림 플로우는 또 다른 연산자의 업스트림 플로우로 작용할 수 있습니다. 시퀀스와 마찬가지로 중간 연산자가 호출되더라도 플로우 코드가 실제로 실행되지 않으며, 반환된 플로우는 콜드 상태입니다.
transform 함수는 업스트림 플로우의 각 원소에 대해 원하는 만큼의 원소를 다운스트림 플로우에 배출할 수 있게 해줍니다.
import kotlinx.coroutines.flow.*
fun main() {
val names = flow {
emit("Jo")
emit("May")
emit("Sue")
}
val upperAndLowercaseNames = names.transform {
emit(it.uppercase())
emit(it.lowercase())
}
runBlocking {
upperAndLowercaseNames.collect { print("$it") }
}
// JO jo MAY may SUE sue
}
take 함수는 수집자와 관련된 코루틴 스코프를 취소하는 방식 외에 플로우 수집을 제어된 방식으로 취소하는 또 다른 방법입니다.
import kotlinx.coroutines.flow.*
fun main() {
val temps = getTemperatures()
temps
.take(5)
.collect {
log(it)
}
}
onCompletion 연산자는 플로우가 정상 종료되거나, 취소되거나, 예외로 종료된 후에 호출되는 람다를 지정할 수 있게 해줍니다.
fun main() = runBlocking {
val temps = getTemperatures()
temps
.take(5)
.onCompletion { cause ->
if (cause != null) {
println("An error occurred! $cause")
} else {
println("Completed!")
}
}
.collect {
println(it)
}
}
onStart는 플로우의 수집이 시작될 때 첫 번째 배출이 일어나기 전에 실행됩니다.
onEach는 업스트림 플로우에서 배출된 각 원소에 대해 작업을 수행한 후 이를 다운스트림 플로우에 전달합니다.
onEmpty는 원소를 배출하지 않고 종료되는 플로우의 경우 로직을 추가로 수행하거나 기본값을 제공할 수 있습니다.
fun main() {
flow
.onEmpty {
println("Nothing - emitting default value!")
emit(0)
}
.onStart {
println("Starting!")
}
.onEach {
println("On $it!")
}
.onCompletion {
println("Done!")
}
.collect()
}
}
fun main() {
runBlocking {
process(flowOf(1, 2, 3))
// Starting!
// On 1!
// On 2!
// On 3!
// Done!
process(flowOf())
// Starting!
// Nothing - emitting default value
// On 0!
// Done!
}
}
buffer 연산자는 버퍼를 추가해서 다운스트림 플로우가 이미 배출된 원소를 처리하느라 바쁜 동안에도 업스트림 플로우가 원소를 배출할 수 있게 해줍니다. 즉, 업스트림 플로우의 실행을 다운스트림 플로우로부터 분리합니다.
fun getAllUserIds(): Flow<Int> {
return flow {
repeat(3) {
delay(200.milliseconds)
log("Emitting!")
emit(it)
}
}
}
suspend fun getProfileFromNetwork(id: Int): String {
delay(2.seconds)
return "Profile[$id]"
}
fun main() {
val ids = getAllUserIds()
runBlocking {
ids
.buffer(3)
.map{ getProfileFromNetwork(it) }
.collect{ log("Got $it") }
}
}
// 433 [main @coroutine#2] Emitting!
// 638 [main @coroutine#2] Emitting!
// 839 [main @coroutine#2] Emitting!
// 2439 [main @coroutine#1] Got Profile[0]
// 4440 [main @coroutine#1] Got Profile[1]
// 6441 [main @coroutine#1] Got Profile[2]
이 코드에서는 3개의 원소를 저장할 수 있는 버퍼를 추가하면 생성자는 새 사용자 식별자를 계속 생성해 버퍼에 넣을 수 있고, 수집자는 네트워크 요청을 계속 처리할 수 있습니다.
플로우에서 원소를 배출하고 처리하는 데 걸리는 시간이 변동될 때, 연산자 사슬에 버퍼를 도입하면 시스템 처리량을 늘리는 데 도움이 될 수 있습니다.
conflate 연산자는 수집자가 바쁜 동안 배출된 항목을 그냥 버립니다. buffer와 마찬가지로 conflate를 쓰면 업스트림 플로우의 실행을 다운스트림 연산자의 실행과 분리할 수 있습니다.
import kotlinx.coroutines.flow.*
fun main() {
runBlocking {
val temps = getTemperatures()
temps
.onEach {
log("Read $it from sensor")
}
.conflate()
.collect {
log("Collected $it")
delay(1.seconds)
}
}
}
플로우 값이 빠르게 ‘구식’이 되고 다른 배출된 원소로 대체되는 경우 느린 수집자가 플로우에서 최신 원소만 처리하게 함으로써 성능을 유지할 수 있습니다.
debounce 연산자는 업스트림에서 원소가 배출되지 않은 상태로 정해진 타임아웃 시간이 지나야만 항목을 다운스트림 플로우로 배출합니다.
val searchQuery = flow {
emit("K")
delay(100.milliseconds)
emit("Ko")
delay(200.milliseconds)
emit("Kotl")
delay(500.milliseconds)
emit("Kotlin")
}
fun main() = runBlocking {
searchQuery
.debounce(250.milliseconds)
.collect {
log("Searching for $it")
}
}
// 743 [main @coroutine#1] Searching for Kotl
// 995 [main @coroutine#1] Searching for Kotlin
flowOn 연산자는 withContext 함수와 비슷하게 코루틴 콘텍스트를 조정합니다.
import kotlinx.coroutines.*
fun main() {
runBlocking {
flowOf(1)
.onEach{ log("A") }
.flowOn(Dispatchers.Default)
.onEach{ log("B") }
.flowOn(Dispatchers.IO)
.onEach{ log("C") }
.collect()
}
}
// 159 [DefaultDispatcher-worker-3 @coroutine#3] A
// 165 [DefaultDispatcher-worker-1 @coroutine#2] B
// 177 [main @coroutine#1] C
flowOn 연산자는 업스트림 플로우의 디스패처에만 영향을 미칩니다. 다운스트림 플로우는 영향을 받지 않으므로 이 연산자를 ‘콘텍스트 보존’ 연산자라고도 부릅니다.
실행은 최종 연산자가 담당하며, 최종 연산자는 단일 값이나 값의 컬렉션을 계산하거나, 플로우의 실행을 촉발시켜 지정된 연산과 부소 효과를 수행합니다. 가장 일반적인 최종 연산자는 collect 입니다.
최종 연산자는 업스트림 플로우의 실행을 담당하기 때문에 항상 일시 중단 함수입니다.
이번 글에서는 플로우 연산자에 대한 내용을 정리해봤습니다.
플로우를 다룰 때 상황과 필요에 맞는 다양한 연산자가 있으며, 이를 활용하면 플로우를 더 유연하게 다룰 수 있겠다는 생각이 들었습니다.
다음에는 오류 처리와 테스트에 대한 글로 돌아오겠습니다.
읽어주셔서 감사합니다!🙂↕️