스트림 처리의 윈도잉과 시간 의미론(이벤트 시간·워터마크·정확히 한 번)
1. 개요
가. 정의
스트림 처리(Stream Processing)는 시작과 끝이 정해지지 않은 무한(unbounded) 데이터가 도착하는 즉시 연속적으로 변환·집계·분석하여 저지연으로 결과를 산출하는 데이터 처리 패러다임이다.
윈도잉(Windowing)은 무한 스트림을 유한한 처리 단위로 자르기 위해, 시간·개수·세션 등의 기준으로 이벤트를 묶는 경계 설정 기법이며, 시간 의미론(Time Semantics)은 이벤트가 "실제로 발생한 시각(event time)"과 "엔진이 처리한 시각(processing time)"을 구분해 집계의 정확성을 보장하는 모델이다.
스트림 처리의 본질적 어려움은 "끝나지 않는 데이터에 대해 언제, 무엇을, 얼마나 정확하게 계산을 확정할 것인가"라는 질문에 있다. 배치 처리는 입력이 모두 모인 뒤 계산하므로 경계가 자명하지만, 스트림은 데이터가 끊임없이 흐르고 네트워크 지연으로 순서가 뒤섞여(out-of-order) 도착하므로, 집계를 확정할 시점을 스스로 정의해야 한다. 윈도잉은 "무엇을 하나의 묶음으로 볼 것인가(where in event time)"를, 워터마크는 "언제 그 묶음을 닫고 결과를 낼 것인가(when in processing time)"를 결정한다. 정보관리기술사 관점에서 이 주제는 단순한 API 사용이 아니라, 지연·정확성·완전성(completeness)의 트레이드오프를 설계하는 실시간 데이터 파이프라인 거버넌스 문제로 이해해야 한다.
본질적으로 스트림 처리는 "데이터가 완전히 모일 때까지 기다리면 늦고, 기다리지 않으면 불완전하다"는 모순을 다룬다. 배치가 완전성을 위해 지연을 감수하고 순수 실시간이 지연을 없애려 완전성을 포기한다면, 현대 스트림 처리는 윈도우·워터마크·트리거라는 세 손잡이로 그 중간 지점을 파라미터화하여, 같은 코드로 "빠르지만 근사한 결과"와 "느리지만 정확한 결과" 사이를 업무 요구에 맞게 조정할 수 있게 한다는 데 그 특징이 있다.
나. 등장 배경과 필요성
첫째, 비즈니스 의사결정의 시간 지평이 "다음 날"에서 "지금 이 순간"으로 이동했다. 이상거래탐지(FDS), 실시간 추천, 사물인터넷(IoT) 설비 이상감지, 급변하는 재고·시세 반영은 야간 배치로는 가치를 잃는다. 수초~수분의 지연 안에 결정을 내려야 하므로, 데이터가 쌓이기를 기다리지 않고 흐르는 중에 처리하는 구조가 필요해졌다.
둘째, 분산 수집 환경에서 데이터는 본질적으로 지연되고 뒤섞여 도착한다. 모바일 기기가 오프라인이었다가 복구되면서 몇 분 전 이벤트를 늦게 전송하거나, 파티션별 처리 속도 차이로 순서가 역전된다. 처리 시각 기준으로 집계하면 "10시에 발생했으나 10시 5분에 도착한 주문"이 엉뚱한 윈도우에 들어가 결과가 왜곡된다. 이를 바로잡으려면 이벤트 자체가 가진 발생 시각을 기준으로 삼는 이벤트 시간 처리와, 지연을 허용하되 한계를 긋는 워터마크 개념이 필요했다.
셋째, 실시간이라고 해서 정확성을 포기할 수는 없다. 장애로 노드가 죽거나 재시작될 때 이벤트가 누락되거나 중복 집계되면, 과금·정산·규제 보고 같은 영역에서는 치명적이다. 따라서 장애가 나도 각 이벤트가 결과에 정확히 한 번만 반영되는 정확히 한 번(exactly-once) 처리 보장이 스트림 엔진의 핵심 요구로 자리 잡았다.
넷째, 운영·비용 측면의 요구도 커졌다. 24시간 멈추지 않고 돌아가는 파이프라인은 상태가 무한히 커지지 않도록 관리해야 하고, 로직을 바꾸면 과거 데이터를 다시 흘려 결과를 재생성(reprocessing)할 수 있어야 한다. 또 순간 폭주(트래픽 스파이크) 시 상류로 부하를 전파해 스스로를 보호하는 역압(backpressure) 제어가 없으면 메모리 고갈로 전체가 붕괴한다. 이처럼 스트림 처리는 정확성뿐 아니라 지속 운영성까지 함께 요구하는 기술로 진화했다.
다. 처리 보장 수준의 스펙트럼
스트림 처리의 신뢰성은 "장애 시 각 이벤트가 결과에 몇 번 반영되는가"로 규정된다. 최소 보장인 at-most-once는 재전송을 하지 않아 유실을 허용하는 대신 중복은 없다. at-least-once는 장애 시 재처리로 유실을 막지만 같은 이벤트가 두 번 반영될 수 있어, 합계·카운트가 과대 계상된다. exactly-once는 상태(state)와 출력이 장애 전후로 일관되게 정확히 한 번만 반영됨을 보장한다. 다만 이는 네트워크 전송 자체를 물리적으로 한 번만 하는 것이 아니라, 체크포인트와 트랜잭션으로 "효과(effect)가 한 번"이 되도록 만드는 것이므로 effectively-once로 부르기도 한다. 이 구분은 뒤의 심화 섹션에서 다룰 체크포인트·트랜잭션 메커니즘으로 구현된다.
보장 수준의 선택은 공짜가 아니다. 상위 보장일수록 상태 스냅샷, 트랜잭션 조율, 지연 커밋 등의 오버헤드가 커지고 지연이 늘어난다. 따라서 모든 파이프라인을 일괄적으로 exactly-once로 설계하기보다, 과금·정산·규제 보고처럼 정확성이 금전·법적 책임과 직결되는 경로만 exactly-once로 두고, 대시보드 지표나 추천 신호처럼 소량의 중복·유실이 허용되는 경로는 at-least-once로 두어 비용을 최적화하는 선택적 적용이 실무의 정석이다.
2. 전체 구조와 시간 의미론
가. 스트림 처리 파이프라인 구조
실시간 파이프라인은 일반적으로 불변 로그 기반 메시지 브로커를 중심으로, 상태를 가진 처리 연산자와 체크포인트 저장소, 그리고 멱등·트랜잭션 싱크로 구성된다.
flowchart LR
SRC["이벤트 소스(앱·IoT·로그)"] --> BROKER["메시지 브로커(불변 로그·파티션)"]
BROKER --> OP["스트림 연산자(윈도우·집계)"]
OP --> ST["상태 저장소(Keyed State)"]
OP --> CP["체크포인트(스냅샷)"]
CP --> DFS["내구성 저장소(분산 파일시스템)"]
OP --> SINK["싱크(멱등·트랜잭션)"]
SINK --> SERVE["서빙(DB·대시보드·알림)"]
메시지 브로커(예: Apache Kafka)는 파티션 단위로 순서가 보장되는 불변 로그를 제공하여, 장애 시 오프셋을 되감아 재생(replay)할 수 있는 재처리의 기반이 된다. 보존 기간(retention)을 예컨대 7일로 두면 그 기간 내 어느 시점으로든 소비 위치를 되돌려 로직을 다시 적용할 수 있고, 파티션 수는 병렬 처리 단위를 결정해 처리량 확장의 기준이 된다. 스트림 연산자는 키별로 상태를 유지하며 윈도우 집계를 수행하고, 주기적으로 상태 스냅샷을 체크포인트로 떠서 분산 파일시스템에 저장한다. 장애가 나면 가장 최근 체크포인트에서 상태를 복원하고 그 지점의 오프셋부터 다시 소비하므로, at-least-once 재처리가 기본이 된다. 여기에 멱등 또는 트랜잭션 싱크를 결합하면 중복 반영까지 제거되어 exactly-once가 완성된다. 이 구조에서 체크포인트 주기는 복구 지점의 조밀함을, 상태 백엔드의 성능은 집계 처리량을 좌우하므로, 두 값을 함께 조율하는 것이 운영 설계의 출발점이 된다.
나. 이벤트 시간·처리 시간·워터마크
시간 의미론의 출발점은 세 가지 시각을 구분하는 것이다. 이벤트 시간은 이벤트가 실제로 발생한 시각(예: 결제 승인 시각)으로 데이터 안에 타임스탬프로 박혀 있다. 수집 시간(ingestion time)은 브로커·엔진에 들어온 시각, 처리 시간은 연산자가 실제로 처리한 시각이다. 결과의 재현성과 정확성은 이벤트 시간 기준 처리에서 나온다. 같은 입력을 다시 흘리면(재처리) 처리 시간은 매번 달라지지만 이벤트 시간 기준 집계는 항상 같은 결과를 주기 때문이다.
문제는 이벤트 시간으로 윈도우를 닫으려면 "그 윈도우에 속할 이벤트가 더 이상 오지 않는다"는 판단이 필요하다는 점이다. 이를 추정하는 장치가 워터마크다. 워터마크 W(t)는 "이벤트 시간 t 이전의 데이터는 (대체로) 모두 도착했다"는 엔진의 진행 표시자이며, 이 워터마크가 윈도우 종료 시각을 넘어서는 순간 해당 윈도우의 집계를 확정·발사(trigger)한다. 실무에서는 관측된 최대 이벤트 시간에서 허용 지연만큼 뺀 유계 무질서(bounded out-of-orderness) 워터마크를 흔히 쓴다. 예를 들어 지연 허용을 5초로 두면 W = max_event_time - 5s가 된다.
sequenceDiagram
participant E as 이벤트 스트림
participant W as 워터마크 생성기
participant WIN as 윈도우 연산자
participant OUT as 결과 싱크
E->>W: 이벤트 도착(이벤트시간 타임스탬프 포함)
W->>WIN: 워터마크 전파("t 이전 도착 완료")
WIN->>WIN: 윈도우 경계와 워터마크 비교
alt 워터마크 > 윈도우 종료
WIN->>OUT: 윈도우 집계 확정·발사
else 지연 데이터 도착
WIN->>WIN: 허용 지연(allowed lateness) 내면 갱신
WIN->>OUT: 사이드 출력(side output)으로 분리
end
참고로 수집 시간은 이벤트 시간과 처리 시간의 절충으로, 소스에 신뢰할 타임스탬프가 없을 때 브로커 진입 시각을 기준으로 삼아 어느 정도의 순서 안정성을 얻는 실용적 대안이다. 다만 재처리 시 결과가 달라질 수 있어 완전한 재현성은 보장하지 못한다.
워터마크는 본질적으로 완전성과 지연의 트레이드오프를 조율하는 손잡이다. 허용 지연을 크게 잡으면 늦게 오는 데이터까지 포함해 결과가 정확해지지만 윈도우 확정이 늦어져 지연이 커진다. 반대로 작게 잡으면 빠르지만 지각 데이터를 놓쳐 결과가 불완전해진다. 그래서 많은 엔진은 윈도우 확정 후에도 허용 지연(allowed lateness) 기간 동안 지각 데이터로 결과를 갱신하고, 그보다 더 늦은 데이터는 사이드 출력으로 분리해 별도 보정 경로로 처리한다.
다. 윈도우의 유형
윈도우는 무한 스트림을 어떤 기준으로 묶느냐에 따라 나뉜다. 유형 선택은 질문의 성격(주기적 집계인지, 활동 구간 분석인지)에 따라 달라지며, 아래 표는 보조 요약일 뿐 각 유형의 선택 이유는 문단으로 설명한다.
| 윈도우 유형 | 정의 | 특징 | 대표 용도 |
|---|---|---|---|
| 텀블링(Tumbling) | 고정 크기·비중첩 | 경계가 겹치지 않아 각 이벤트가 정확히 1개 윈도우 | 1분 매출 집계, 시간별 리포트 |
| 슬라이딩(Sliding) | 고정 크기·일정 간격 이동 | 윈도우가 겹쳐 이벤트가 여러 윈도우에 속함 | 최근 5분 이동평균(1분마다) |
| 세션(Session) | 활동 간 공백(gap) 기준 | 크기 가변, 활동 구간 자동 분리 | 사용자 세션, 장비 가동 구간 |
| 전역(Global) | 경계 없음·커스텀 트리거 | 사용자가 트리거 정의 | 개수 기반 집계 |
텀블링 윈도우는 1분, 1시간처럼 고정 길이로 겹치지 않게 자르며, 각 이벤트가 정확히 하나의 윈도우에만 속한다. "매 1분 거래 건수"처럼 명확한 주기 리포트에 적합하다. 경계가 단순해 상태 관리 비용이 가장 낮고, 윈도우가 닫히면 해당 상태를 바로 폐기할 수 있어 메모리 예측이 쉽다. 반면 경계가 고정이라, 임계 현상이 두 윈도우에 걸쳐 나뉘면 각각은 임계치 미만이 되어 탐지를 놓칠 수 있다는 점을 설계 시 감안해야 한다.
슬라이딩 윈도우는 크기(예: 5분)와 이동 간격(예: 1분)을 따로 두어 윈도우가 겹치므로, "1분마다 갱신되는 최근 5분 이동평균" 같은 연속 추세 모니터링에 쓴다. 한 이벤트가 (크기÷간격)개의 윈도우에 중복 포함되므로 상태 크기와 연산량이 그만큼 커진다. 위 예에서는 한 이벤트가 5개 윈도우에 속하므로, 단순 텀블링 대비 약 5배의 상태·연산 부담이 생긴다. 추세의 민감도(간격을 짧게)와 자원 비용 사이의 절충이 핵심 설계 포인트다.
세션 윈도우는 고정 경계가 아니라 활동 사이의 공백(gap)으로 구간을 나눈다. 예컨대 공백 임계치를 30분으로 두면, 사용자의 연속 클릭이 30분 이상 끊기면 세션이 종료되고 새 세션이 시작된다. 사용자 행동 분석이나 설비 가동 세션 추출에 자연스럽다. 윈도우 길이가 데이터에 따라 동적으로 결정되므로, 지각 이벤트가 두 세션 사이의 공백을 메우면 이미 발사된 두 세션을 하나로 병합(merge)해야 하는 복잡성이 생긴다. 이 때문에 세션 윈도우는 상태·트리거 설계가 가장 까다로운 유형이다.
전역 윈도우는 시간 경계 없이 "100건마다"처럼 사용자 정의 트리거로 발사한다. 시간 기반 집계로 표현하기 어려운 개수·임계 조건 기반 처리에 유연하게 대응하되, 트리거와 상태 정리를 전적으로 사용자가 책임져야 하므로 남용하면 상태 누수와 메모리 증가의 원인이 된다.
라. 트리거와 누적 모드
같은 윈도우라도 "언제, 몇 번 결과를 낼 것인가"는 트리거(Trigger)로, "재발사 시 이전 결과와 어떻게 합칠 것인가"는 누적 모드(accumulation mode)로 분리해 설계한다. 지각 데이터가 중요한 환경에서는 워터마크 도달 시 1차 결과를 빠르게 내고(speculative), 이후 지각분이 오면 윈도우를 재발사해 결과를 정정하는 방식이 유용하다. 이때 누적(accumulating) 모드는 매번 전체 재계산 값을 내보내고, 폐기 후 누적(accumulating & retracting) 모드는 이전 값을 취소하는 보정 레코드와 새 값을 함께 내보내 다운스트림이 이중 집계하지 않도록 한다. 이 "조기 결과 + 사후 정정" 패턴은 Dataflow 모델의 핵심 기여로, 완전성을 희생하지 않으면서 저지연을 함께 얻는 열쇠다.
3. 비교와 적용 사례
가. 대표 엔진 비교와 차이의 이유
| 구분 | Apache Flink | Kafka Streams | Spark Structured Streaming |
|---|---|---|---|
| 처리 모델 | 진정한 레코드 단위 스트리밍 | 레코드 단위(라이브러리) | 마이크로배치(기본)·연속(실험적) |
| 상태·체크포인트 | 비동기 배리어 스냅샷(ABS) | 체인지로그 토픽 + RocksDB | 체크포인트 + WAL |
| exactly-once | 상태+트랜잭션 싱크로 보장 | exactly_once_v2(Kafka 간) |
멱등·트랜잭션 싱크 기반 |
| 지연 특성 | 밀리초급 저지연 | 밀리초급(앱 내장) | 배치 주기만큼(수백 ms~초) |
세 엔진의 차이는 "스트리밍을 무엇으로 보느냐"는 근본 관점에서 비롯된다. Flink는 스트림을 1급 시민으로 보고 레코드 단위로 연속 처리하므로 저지연과 정교한 이벤트 시간·워터마크 제어에 강하다. Spark Structured Streaming은 전통적으로 스트림을 작은 배치의 연속으로 다루는 마이크로배치 모델이어서 처리량과 배치 생태계 통합에 유리한 대신 배치 주기만큼의 지연이 있다(연속 처리 모드는 지연을 낮추나 보장 수준에 제약이 있다). Kafka Streams는 별도 클러스터 없이 애플리케이션에 내장되는 라이브러리라 운영이 단순하고, 상태를 체인지로그 토픽으로 복제해 복원하며 카프카 내부 입출력에 대해 exactly-once를 제공한다. 즉 저지연·정교한 시간 제어가 핵심이면 Flink, 카프카 중심의 경량 내장 처리면 Kafka Streams, 배치·ML 생태계와의 통합이 중요하면 Spark가 합리적 선택이 된다.
나. 구체 적용 사례
실시간 이상거래탐지(FDS)를 보자. 카드 승인 이벤트를 카드번호 키로 묶어, 1분 텀블링 윈도우로 승인 건수·금액을 집계하고 임계치를 넘으면 차단 알림을 발사한다. 모바일 결제 특성상 일부 이벤트가 수초 늦게 도착하므로 워터마크 지연을 5초로 두어 지각분까지 포함하되, 그보다 늦은 이벤트는 사이드 출력으로 보내 사후 정산에 반영한다. 처리 시간이 아닌 승인 시각(이벤트 시간) 기준이어야 "같은 1분 구간" 판정이 일관된다.
산업 IoT 사례로, 수천 대 설비의 센서 스트림에서 세션 윈도우(공백 임계치 10분)로 가동 세션을 분리하면, 한 세션 내 진동·온도 추세로 이상을 조기 감지할 수 있다. 가동과 정지가 불규칙한 설비에서는 고정 윈도우보다 세션 윈도우가 현실의 운전 구간을 자연스럽게 반영한다. 글로벌 차량 호출 서비스가 실시간 수요·공급 매칭과 ETA 예측에 Flink 기반 스트림 처리를 대규모로 운용하는 것도 같은 맥락으로 알려져 있다.
전자상거래 추천에서는 슬라이딩 윈도우가 쓰인다. 사용자별로 최근 10분 클릭을 1분 간격 슬라이딩으로 집계해 관심 상품 추세를 실시간 반영하면, 세션이 길어질수록 추천이 최신 행동에 민감하게 반응한다. 이때 이벤트가 여러 윈도우에 중복 포함되어 상태가 커지므로, 상태 TTL과 체크포인트 주기(예: 10초)를 함께 설계해 메모리·복구 시간을 관리한다. 예컨대 활성 사용자 100만 명 × 평균 상태 1KB면 약 1GB의 키별 상태가 상시 유지되므로, 상태 백엔드를 메모리에서 RocksDB로 옮기고 유휴 사용자의 상태를 TTL로 만료시키는 결정이 처리량·비용에 직접적인 영향을 준다.
4. 심화: 정확히 한 번 처리의 구현 메커니즘과 최신 동향
exactly-once는 "전송을 한 번만"이 아니라 "효과가 한 번만"을 보장하는 문제다. Flink는 비동기 배리어 스냅샷(Asynchronous Barrier Snapshotting)으로 이를 구현한다. 소스가 데이터 흐름에 특수한 배리어(barrier)를 주기적으로 삽입하면, 각 연산자는 배리어가 지나갈 때 자신의 상태를 스냅샷으로 떠서 저장한다. 모든 연산자의 스냅샷이 모이면 전역적으로 일관된(globally consistent) 체크포인트가 완성되고, 장애 시 이 지점으로 상태와 소스 오프셋을 함께 되돌린다. 처리를 멈추지 않고 비동기로 스냅샷을 뜨기 때문에 성능 저하가 작다는 것이 핵심이다. 이는 분산 스냅샷 이론(Chandy-Lamport 알고리즘)을 스트림 데이터 흐름에 맞게 변형한 것이다.
상태 복원만으로는 싱크로 이미 나간 출력의 중복을 막을 수 없으므로, 출력 단에서는 멱등 쓰기 또는 트랜잭션 쓰기(2단계 커밋)를 결합한다. 예컨대 Flink의 TwoPhaseCommitSinkFunction은 체크포인트 완료 시점에 맞춰 외부 트랜잭션을 커밋하여, 체크포인트와 출력 가시성을 원자적으로 묶는다. Kafka는 트랜잭션(프로듀서 트랜잭션)과 컨슈머 오프셋 커밋을 하나의 트랜잭션으로 묶는 read-process-write 패턴으로 카프카 내부 입출력에 대한 exactly-once를 제공하며, Kafka Streams는 processing.guarantee=exactly_once_v2 설정으로 이를 단순화했다(KIP-447 계열 개선으로 다수 파티션 처리 시 효율이 향상된 것으로 알려져 있다).
한편 Spark Structured Streaming은 체크포인트에 오프셋과 상태를 기록하는 WAL(Write-Ahead Log) 기반으로 장애 복구를 보장하며, 멱등 싱크나 트랜잭션 싱크와 결합해 end-to-end 정확성을 확보한다. Kafka Streams는 상태를 체인지로그 토픽으로 복제해 두었다가 다른 인스턴스에서 복원하므로, 상태를 가진 애플리케이션을 무상태 서비스처럼 수평 확장·재배치할 수 있다. 엔진마다 메커니즘의 구현은 다르지만, "일관된 상태 스냅샷 + 출력의 멱등·트랜잭션 보장"이라는 공통 원리는 동일하다는 점을 이해하는 것이 핵심이다.
최신 동향으로는 몇 가지 흐름이 두드러진다. 첫째, Google의 Dataflow 모델에서 정립된 "무엇(What)·어디서(Where, 윈도우)·언제(When, 워터마크·트리거)·어떻게(How, 지각 데이터 정정)"의 4가지 축이 Apache Beam을 통해 엔진 중립적 추상으로 확산되어, 하나의 파이프라인 정의를 여러 실행 엔진에서 재사용하는 방향이 강화되고 있다. 둘째, 스트림을 테이블처럼 SQL로 다루는 스트리밍 SQL(Flink SQL 등)이 성숙해 비개발 직군의 접근성이 높아졌다. 셋째, 상태 저장소를 클라우드 오브젝트 스토리지로 분리해 컴퓨트와 상태를 독립적으로 확장하는 상태·컴퓨트 분리 아키텍처와, 스트림·배치를 단일 테이블 포맷(예: Apache Iceberg 등 레이크하우스) 위에서 통합하려는 시도가 활발하다. 다만 이들 세부 수치·버전은 빠르게 변하므로 설계 시 최신 공식 문서로 확인하는 것이 바람직하다.
5. 고려사항 및 시사점
첫째, 지연과 완전성의 트레이드오프를 비즈니스 요구로 환산해 설계한다. 워터마크 허용 지연과 윈도우 크기는 기술 파라미터이기 이전에 "얼마나 늦은 데이터까지 결과에 반영할 가치가 있는가"라는 업무 판단이다. 정산처럼 완전성이 중요하면 허용 지연을 넉넉히 두고 사후 정정 경로를 두며, 경보처럼 즉시성이 중요하면 짧게 잡고 지각분은 별도 보정한다. 기술사 관점에서는 SLA·정확성 요구를 먼저 정의하고 파라미터를 역산하는 접근이 바람직하다.
둘째, exactly-once의 범위와 비용을 명확히 한다. exactly-once는 보통 "엔진 상태와 특정 싱크"에 한정되며, 외부 시스템 전체로 자동 확장되지 않는다. 멱등 설계가 불가능한 외부 연계(예: 외부로의 이메일·결제)는 트랜잭션으로 묶기 어려워 at-least-once + 멱등 키 설계로 보완해야 한다. 또한 체크포인트 주기를 짧게 하면 복구 시간은 줄지만 오버헤드가 커지므로, 복구목표시간(RTO)과 처리량의 균형점을 찾아야 한다.
셋째, 상태 관리와 재처리 전략을 함께 설계한다. 키별 상태가 무한히 커지지 않도록 TTL·상태 백엔드(예: RocksDB)·상태 크기 모니터링을 두고, 로직 변경 시 과거 데이터를 안전하게 재처리할 수 있도록 브로커 보존 기간과 오프셋 되감기 전략, 결과 뷰의 멱등 재생성 구조를 미리 마련한다. 카파 아키텍처([[lambda-kappa-architecture]])의 재처리 철학과 연계해 설계하면 운영 복잡도를 낮출 수 있다.
넷째, 운영 관측성과 역압(backpressure) 대응을 내재화한다. 워터마크 지연 급증, 체크포인트 실패율, 소비 지연(consumer lag), 역압 발생 구간은 실시간 파이프라인의 건강 지표다. 리액티브 스트림의 배압 제어([[reactive-streams-backpressure]])와 메시지 큐([[message-queue-kafka]])의 파티션·오프셋 관측을 결합해, 폭주 시 처리량을 자동 조절하고 상태가 복구 가능한 범위를 벗어나지 않게 해야 한다.
다섯째, 전망과 연계 기술. 스트리밍 SQL과 Beam 기반 엔진 중립 추상, 레이크하우스 기반 스트림·배치 통합은 "실시간과 배치의 경계를 흐리는" 방향으로 수렴하고 있다. 최종 일관성([[eventual-consistency]])·CQRS와 이벤트 소싱([[cqrs-event-sourcing]]) 같은 데이터 일관성 모델과 함께 보면, 스트림 처리는 단일 기능이 아니라 실시간 데이터 플랫폼 전반을 꿰는 설계 축으로 이해해야 한다.
참고자료
- Apache Flink Documentation, "Timely Stream Processing / Windows": https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/
- Tyler Akidau et al., "The Dataflow Model", VLDB 2015: https://research.google/pubs/pub43864/
- Apache Kafka Documentation, "Exactly Once Semantics": https://kafka.apache.org/documentation/#semantics
- Apache Beam Programming Guide, "Windowing & Watermarks": https://beam.apache.org/documentation/programming-guide/
한 줄 요약: 스트림 처리는 윈도잉으로 무한 데이터를 유한 단위로 묶고, 이벤트 시간·워터마크로 집계 확정 시점을 정하며, 체크포인트와 트랜잭션 싱크로 정확히 한 번 처리를 보장해 지연·정확성·완전성의 트레이드오프를 설계하는 실시간 데이터 처리 기술이다.