bluetape4k Flow 확장 함수: 복잡한 흐름을 명시적인 운영 계약으로 바꾸기

Kotlin Flow 코드는 처음에는 단순합니다. MutableSharedFlow를 만들고, launch에서 수집하며, 필요하면 Job을 취소합니다.
그러나 검색 자동 완성, 대체 데이터 원본, 콜백 브리지, 메트릭 샘플링 같은 현실적인 요구가 들어오면 코드가 빠르게
“언제 누가 취소하는지”, “마지막 값만 필요한지”, “실패를 값으로 다룰지 종료 예외로 다룰지”를 숨기기 시작합니다.
bluetape4k-workshop의 Flow 확장 함수 예제는 이 숨은 결정을 연산자 이름으로 끌어올립니다. 이번 글은 여섯 예제를
하나씩 소개하되, API 목록이 아니라 어떤 문제를 어떤 운영 계약으로 바꾸는지에 초점을 둡니다.
| 예제 | 다루는 문제 | 핵심 Flow 확장 함수 |
|---|---|---|
| Search Pipeline | 빠른 입력, 설정 시점의 값, 최신 검색만 유지 | bufferingDebounce, withLatestFrom, flatMapLatest, takeUntil |
| Race and Fallback | 가장 빠른 정상 원본, 순차 대체 경로, 부분 병합 | race, amb, concat, concatArrayEager, concatMapEager, merge |
| Subject Bridge | 콜백을 이벤트·상태·이력·작업 큐 스트림으로 변환 | PublishSubject, BehaviorSubject, ReplaySubject, MulticastSubject, UnicastWorkSubject |
| Event Aggregation | 주문 이벤트 재생, 윈도우, 그룹화, 누적 상태 | chunked, windowed, groupBy, scanWith, bufferUntilChanged, zipWithNext |
| Parallel Enrichment | 주문별 독립 보강 작업을 병렬 레일로 처리 | parallel, sequential |
| Metrics Sampling | 연속 메트릭에서 빠른 미리보기와 안정된 대시보드 값을 분리 | throttleLeading, throttleTrailing, pairwise, takeUntil, mapResultCatching |
ReactiveX operator 문서는 buffer, groupBy, scan, window를
서로 다른 변환 operator로 분류하고, Reactor의
Three Sorts of Batching
문서는 그룹화, 윈도잉, 버퍼링을 함께 비교합니다. 다만 출력 형태는 서로 다릅니다. scan 계열은 배치 처리가 아니라
이전 상태와 현재 이벤트로 다음 상태를 만드는 누산 연산으로 봐야 합니다.
아래 그림은 RxJava 문서와 튜토리얼에서 자주 쓰는 마블 다이어그램 문법을 따릅니다. 기본적으로 위쪽 타임라인은 입력 Flow, 가운데 상자는 연산자, 아래쪽 타임라인은 출력 Flow를 나타냅니다. 병렬 레일처럼 타임라인 하나로 설명하기 어려운 예제는 같은 시각 언어로 분기·병합 지점을 드러냅니다.
입력값에서 출력값으로 먼저 보기
섹션 제목: “입력값에서 출력값으로 먼저 보기”아래 코드는 workshop 테스트가 고정하는 입출력을 글에서 읽기 쉽게 줄인 형태입니다. 실제 예제는 각 모듈의 README와 테스트에 더 많은 실패/취소/검증 케이스를 포함합니다.
// 1. 검색 파이프라인// 입력: 한 번의 연속 입력에서 "r", "re", "red"val results = pipeline.search( queries = flowOf("r", "re", "red"), settings = flowOf(settings(tenantId = "tenant-a")), sessionClosed = flowOf(), debounce = 100.milliseconds,).toList()// 출력: 요청 하나, query == "red", tenantId == "tenant-a"
// 2. 경쟁 / 대체 경로// 입력: cache(200ms), replica(20ms), remote(120ms)val winner = catalog.fastestHealthy(listOf(cache, replica, remote)).take(1).toList()// 출력: [CatalogSource.REPLICA], cache/remote 수집 작업은 취소
// 3. Subject 브리지// 입력: 구독 전 이벤트 하나, 구독 뒤 활성 이벤트 하나bridge.publishEvent(DeviceEvent("event-01", "device-01", CONNECTED, "early"))val received = async { bridge.events.take(1).toList() }bridge.awaitEventSubscribers()bridge.publishEvent(DeviceEvent("event-02", "device-01", TELEMETRY, "temperature=22"))// 출력: PublishSubject 구독자는 event-02만 수신
// 4. 이벤트 집계// 입력: sampleEvents() == 주문 이벤트 다섯 개val chunked = pipeline.chunkedActivity(sampleEvents().asFlow(), chunkSize = 2).toList()// 출력 청크 크기: [2, 2, 1]
val windows = pipeline.rollingActivity(sampleEvents().asFlow(), size = 3, step = 2).toList()// 출력 윈도우 크기: [3, 3, 1]
val groups = pipeline.groupedByOrder(sampleEvents().asFlow()).toList().associateBy { it.key }// 출력 그룹:// order-1 -> [OrderCreated, LineAdded, PaymentAuthorized]// order-2 -> [OrderCreated, ShipmentStarted]
val states = pipeline.readModels(sampleEvents().asFlow()).toList()// 출력: 상태 스냅숏, 최종 order-1 == PAID, 최종 order-2 == SHIPPED
// 5. 병렬 보강// 입력: O-1001 customer-1001, item suit-01 x2val enriched = pipeline.enrichInParallel(flowOf(order), parallelism = 3) { dispatchers[it] } .toList()// 출력: O-1001 -> loyaltyGrade=GOLD, discountPercent=10, fulfillable=true
// 6. 메트릭 샘플링// 입력: 200ms마다 cpu.usage 값 10,20,30,...,100val leading = pipeline.leadingPreview(highFrequencyCpuSamples(), 501.milliseconds).toList()val trailing = pipeline.dashboardSamples(highFrequencyCpuSamples(), 501.milliseconds).toList()// 출력 선두 값: [10.0, 40.0, 70.0, 100.0]// 출력 후미 값: [30.0, 60.0, 90.0, 100.0]검색 파이프라인: 연속 입력과 최신 조건을 분리한다
섹션 제목: “검색 파이프라인: 연속 입력과 최신 조건을 분리한다”검색 자동완성은 생각보다 많은 결정을 포함합니다. 빈 문자열은 버려야 하고, 사용자가 빠르게 입력하면 마지막 검색어만 의미가 있으며, 테넌트 설정이나 기능 플래그는 검색 실행 시점의 최신 값을 사용해야 합니다. 세션이 끝나면 대기 중인 검색도 멈춰야 합니다.
예제의 핵심 흐름은 다음과 같습니다.

queries .map(String::trim) .filter(String::isNotBlank) .bufferingDebounce(debounce) .mapNotNull { burst -> burst.lastOrNull()?.let(::SearchQuery) } .withLatestFrom(settings) { query, latestSettings -> SearchRequest(query, latestSettings) } .flatMapLatest { request -> adapter.search(request) } .takeUntil(sessionClosed)bufferingDebounce는 연속 입력을 List<String>으로 묶습니다. 따라서 mapNotNull에서 묶음의 마지막 검색어를
SearchQuery로 바꾸는 단계가 필요합니다. withLatestFrom은 검색 요청을 만드는 순간의 최신 설정을 붙이고,
flatMapLatest는 이전 검색이 아직 끝나지 않았더라도 더 최신 검색어가 들어오면 이전 검색을 취소합니다.
여기서 중요한 것은 취소입니다. 검색 어댑터는 CancellationException을 삼키면 안 됩니다. 예제 테스트는
취소가 호출자에게 정상적으로 전파되는지 고정합니다. 스트림 연산자는 취소 경계를 숨기는 도구가 아니라, 어디에서
취소되어야 하는지 명확히 드러내는 도구입니다.
경쟁과 대체 경로: 빠른 원본과 우선순위가 높은 원본은 다르다
섹션 제목: “경쟁과 대체 경로: 빠른 원본과 우선순위가 높은 원본은 다르다”캐시, 로컬 복제본, 원격 API, 예비 엔드포인트가 있을 때 “먼저 오는 값”이 항상 정답은 아닙니다. 어떤 화면은 가장 빠른 정상 응답 하나면 충분하고, 어떤 작업은 우선순위에 따라 대체 경로를 선택해야 하며, 어떤 보강 작업은 여러 원본의 부분 결과를 합쳐야 합니다.
flow-extensions-race-fallback 예제는 이 차이를 연산자 선택으로 드러냅니다.

race와 amb는 먼저 정상 값을 낸 원본을 고르고 나머지를 취소합니다. concat은 우선순위를 지키고, merge는 여러 원본의 부분 결과를 함께 사용합니다.| 선택 | 의미 |
|---|---|
race / amb | 먼저 정상 값을 내는 원본을 선택하고 나머지는 취소 |
concat | 엄격한 대체 경로 순서를 유지 |
concatArrayEager | 원본은 미리 시작하되 출력 순서는 유지 |
concatMapEager | 동적 원본을 즉시 시작하고 외부 순서를 보존 |
merge | 부분 보강처럼 여러 원본의 기여를 함께 사용 |
materialize / dematerialize | 실패를 값으로 관찰한 뒤 마지막에 종료 오류로 변환 |
이 예제는 select나 async를 직접 조립하는 방식을 없애자는 이야기가 아닙니다. “지연 시간 우선인가, 원본 우선순위
우선인가, 부분 결과 병합인가”라는 정책을 코드 표면에서 읽을 수 있게 만들자는 뜻입니다.
Subject 브리지: 콜백 스트림도 종류가 있다
섹션 제목: “Subject 브리지: 콜백 스트림도 종류가 있다”외부 SDK 콜백을 Flow로 바꿀 때 callbackFlow 하나로 모든 것을 처리하면 늦게 참여한 구독자의 의미가 흐려집니다.
과거 이벤트를 받을 수 있는지, 최신 상태 하나만 받으면 되는지, 분기 전달을 기다려야 하는지, 작업 항목을 한 번만 소비해야
하는지가 서로 다르기 때문입니다.
Subject 브리지 예제는 콜백을 다섯 가지 스트림 성격으로 나눕니다.

| Subject | 쓰임 |
|---|---|
PublishSubject | 현재 참여한 구독자에게만 이벤트 전달 |
BehaviorSubject | 늦게 참여한 구독자에게 최신 상태를 먼저 전달 |
ReplaySubject | 제한된 이력을 재생 |
MulticastSubject | 예상한 구독자 수가 모일 때까지 분기 전달을 조율 |
UnicastWorkSubject | 작업 항목을 단일 소비자 큐처럼 한 번만 소비 |
이 구분은 테스트에서도 중요합니다. 예제는 늦게 참여한 구독자가 과거 이벤트를 받는지, 최신 상태를 받는지, 작업 항목이 중복 소비되지 않는지를 각각 확인합니다. 콜백 브리지의 핵심은 “Flow로 바꿨다”가 아니라 “어떤 구독 계약으로 바꿨는가”입니다.
집계와 샘플링: 윈도우는 메모리 계약이다
섹션 제목: “집계와 샘플링: 윈도우는 메모리 계약이다”이벤트 집계 예제는 재생 스트림을 chunked, windowed, groupBy, scanWith로 다룹니다. 메트릭 샘플링
예제는 연속 메트릭을 throttleLeading과 throttleTrailing로 나눕니다.
둘의 공통점은 윈도우가 단순한 편의 API가 아니라 메모리와 지연 시간의 계약이라는 점입니다.
chunked: 이벤트를 구체적인 배치로 묶는다
섹션 제목: “chunked: 이벤트를 구체적인 배치로 묶는다”
chunked(2)는 타임라인을 두 개씩 잘라 List 형태의 배치를 내보냅니다. 마지막 부분 배치를 허용하면 [e5]도 출력됩니다.events .chunked(chunkSize, partialWindow = true) .map { chunk -> summarize(chunk) }위 예제 입력이 e1, e2, e3, e4, e5라면 출력 배치 크기는 [2, 2, 1]입니다. 이 연산자를 선택할 때는 배치 하나가
메모리에 얼마나 쌓일 수 있는지 같이 설명해야 합니다.
windowed: 겹치는 윈도우를 만든다
섹션 제목: “windowed: 겹치는 윈도우를 만든다”
windowed(size = 3, step = 2)는 겹치는 윈도우를 만들 수 있습니다. 같은 이벤트가 둘 이상의 윈도우에 들어갈 수 있다는 점이 chunked와 다릅니다.sampleEvents() 기준으로 첫 번째 윈도우는 [e1, e2, e3], 두 번째 윈도우는 [e3, e4, e5], 마지막 부분 윈도우는
[e5]입니다. 출력 크기는 [3, 3, 1]이지만, 이는 배치 크기가 아니라 윈도우 크기입니다.
groupBy: 시간이나 크기가 아니라 키로 스트림을 나눈다
섹션 제목: “groupBy: 시간이나 크기가 아니라 키로 스트림을 나눈다”
groupBy(orderId)는 타임라인을 자르지 않습니다. 같은 키를 가진 이벤트가 별도의 스트림으로 흐르도록 분할합니다.order-1은 Created, LineAdded, PaymentAuthorized 스트림이 되고, order-2는 Created, ShipmentStarted
스트림이 됩니다. 따라서 groupBy는 “배치가 몇 개 생성되었는가”보다 “키별 스트림 생명주기를 어떻게 닫을 것인가”가 더
중요합니다.
scanWith: 배치가 아니라 상태 스냅숏을 누적한다
섹션 제목: “scanWith: 배치가 아니라 상태 스냅숏을 누적한다”
scanWith(seed, accumulator)는 이벤트를 묶지 않습니다. 이전 상태와 현재 이벤트로 다음 읽기 모델 스냅숏을 만듭니다.readModels(sampleEvents())는 배치 스트림이 아니라 상태 스냅숏 스트림입니다. 마지막 상태를 보면 order-1은
PAID, order-2는 SHIPPED가 됩니다.
metrics .throttleLeading(window) .filter { it.value >= alertThreshold }
throttleLeading은 window의 첫 값을 빠르게 내보내고, throttleTrailing은 window가 안정된 뒤 마지막 값을 내보냅니다.throttleLeading은 경보 미리보기에 적합합니다. 윈도우 안의 첫 값을 빠르게 보여주기 때문입니다. 반대로
throttleTrailing은 대시보드 타일에 적합합니다. 윈도우가 안정된 뒤 마지막 값을 보여주기 때문입니다.
예제는 입력 검증도 스트림 바깥으로 밀어내지 않습니다. 유효하지 않은 ID, 음수 금액, 제어 문자, 유한하지 않은 메트릭은
수집 전에 실패해야 합니다. 로그용 toString()과 Flow<T>.log()도 민감한 값을 그대로 노출하지 않도록 정제를
전제로 둡니다.
병렬 보강: 병렬 레일은 독립 작업에만 사용한다
섹션 제목: “병렬 보강: 병렬 레일은 독립 작업에만 사용한다”flow-extensions-parallel-enrichment는 주문별 고객 조회, 재고 확인, 할인 계산을 병렬 레일로
나눕니다. 각 주문의 보강 작업이 독립적이기 때문에 Flow<T>.parallel(...)로 레일을 나누고, 마지막에 sequential()로
다시 단일 스트림으로 병합합니다.

parallel(3)은 독립적인 보강 작업을 레일로 나누고, sequential()은 결과를 다시 단일 출력 스트림으로 병합합니다.orders .parallel(parallelism, runOn) .map { rail -> enrich(rail) } .sequential()여기서 병렬화의 기준은 “빨라질 것 같다”가 아닙니다. 각 레일이 공유 가변 상태 없이 독립적으로 처리될 수 있고, 마지막에 다시 순서를 합치는 지점이 명확해야 합니다. 이 조건이 없으면 연산자가 간결해 보여도 경쟁 상태를 감춘 코드가 됩니다.
운영에서 조심할 점
섹션 제목: “운영에서 조심할 점”Flow 확장 함수를 사용한다고 코루틴 운영 계약이 바뀌지는 않습니다. 오히려 더 명확히 드러나야 합니다.
- suspend 호출 주변에서
CancellationException을 삼키지 않는다. - 생명주기 종료 신호는
takeUntil같은 명시적인 경계로 둔다. - 대체 경로와 경쟁 연산자를 함께 사용할 때 지연 시간, 원본 우선순위, 부분 병합 중 무엇을 우선할지 정한다.
- 윈도우·청크·그룹화에는 메모리와 지연 시간의 상한을 함께 설명한다.
- 로그와 디버그 출력에 비밀값, 토큰, 비밀번호, 원본 본문을 그대로 노출하지 않는다.
- 병렬 레일은 독립성을 검증한 작업에만 적용하고, 다시 단일 스트림으로 병합하는 위치를 명확히 둔다.
이 예제의 장점은 Flow를 단순히 “더 세련된 연산자 체인”으로 보여주지 않는다는 점입니다. 각 연산자가 실패, 취소, 구독자, 윈도우, 배압 가운데 어떤 계약을 담당하는지 테스트와 README에서 함께 드러냅니다.
댓글
GitHub 계정으로 의견을 남기거나 reaction을 남길 수 있습니다.