본문으로 건너뛰기
DDIA 11장: 스트림 처리

DDIA 11장: 스트림 처리

2026년 7월 1일

이 책은 논리적인 흐름으로 스트림 처리를 어떻게 하는지 전개한다.

개요

우선 개요에서는 이전 장인 10장에서 다루던 일괄 처리와의 특징인 입력 크기를 사전에 한정할 수 있다는 특징을 들어 설명한다. ( 력 크기를 한정 -> 입력을 읽는 작업이 끝나는 시점을 알 수있다.) 이는 입력이 모두 끝나야, 출력을 할 수 있지, 조기에 출력을 할 수 없다는 것이다.

하지만 실제 데이터는 끊임없이 들어오고(생산되고) 이를 처리해야만 한다. 그래서 일괄 처리 프로세서는 스케줄링을 걸어서, 일정기간 데이터 청크를 나눈다.

그러나 이 배치처리의 문제점은 입력의 변화가 하루가 지나야 반영된다는 것 -> 비즈니스 요구사항에 맞는가?

아니다 바로 반영이 된다면 좋겠다! 이러려면 그러면 배치를 줄이는 식으로 하면 되겠다! 하루를 1시간으로, 분으로, 초로… 이렇게 배치를 줄이는 방식(= 고정된 시간조각)에서 사고방식을 벗어나서

대안은!

단순히 이벤트가 발생할 때마다 처리해야한다! 이벤트는 스트림의 최소단위이고 특정 시각에 일어난 불변의 사실

이게 스트림 처리이다

스트림은 시간 흐름에 따라 점진적으로 생산된 데이터

데이터 관리 메커니즘으로 이벤트 스트림을 설명

  1. 스트림을 표현하는 방법
  2. 저장하는 방법
  3. 네트워크 상에서 전송하는 방법
  4. 스트림과 데이터베이스 사이의 관계
  5. 스트림 처리

이벤트 스트림 전송(배치랑은 무엇이 다른가?)

이벤트는 스트림의 최소단위이고 특정 시각에 일어난 불변의 사실 추가적으로 이벤트는 일기준 시계를 따르는 이벤트 발생 타임스탬프를 포함 이벤트는 text,json,바이너리로 부호화될 수 있음 -> 파일저장 + 전송의 용도 -> 스키마 메모리 -> 파일 -> 네트워크 -> 파일 -> 메모리

Producer가 이벤트를 생산하면 해당 이벤트를 복수의 Consumer가 처리할 수 있다

이론 상으로는 파일이나 데이터베이스가 있으면 생산자와 소비자를 연결하기 충분 -> 폴링해서 변경분을 알아내면 되기 때문 -> 다만 OLTP에 부하가 걸림

오히려 소비자쪽에서 체크하기보다는 변경에 대한 knowledge가 있는 생산자가 이벤트가 발생했음을 소비자에게 알리는 게 낫다. + 책에서는 안나왔지만 A->B->C에서 최종상태인 C이벤트만 집계될 수 있음

기존 트리거 기능이 있지만, 이는 한계가 있는듯

이벤트 스트림 전송

메시징 시스템(PUB/SUB)

메시징 시스템 생산자는 이벤트를 포함한 메세지를 전송 메시지는 소비자에게 전송

가장 간단한 방법 unix pipe, TCP -> Pub 1—1 SUB 반면에 메시징 시스템은 다수의 생산자 노드가 같은 토픽으로 메시지를 전송할 수 있음, 다수의 소비자 노드가 토픽 하나에서 메시지를 받아 갈 수 있음

PUB — Topic —SUB

근데 이제 2개의 핵심 질문이 나옴

  1. 생산자가 소비자가 메시지를 처리하는 속도보다 빠르게 메시지를 전송한다면 어떻게 될까?
    1. 메시지를 버리거나
    2. 큐에 메시지를 버퍼링
    3. 백프레셔
    4. (사실 queue를 일단 쓰고, 그 다음에 꽉찼을 때 drop이냐 backpressure로 분기, 또 카프카는 또 다른 흐름임)
  2. 노드가 죽거나 일시적으로 오프라인이 된다면 어떻게 될가? 손실되는 메시지가 있을까?
    1. 디스크 기록
      1. 방지하는 장애는 프로세스 재시작
    2. 복제본 생성
      1. 노드 완전 소실. 하드웨어 자체가 맛이감
    3. 아니면 둘다

메시지의 유실을 허용할지 말지

생산자에서 소비자로 메시지를 직접 전달하기(직접 메시징 시스템)

직적 메시징은 중간 노드 없이 생산자와 소비자를 네트워크로 직접 통신한다

직접 메시징 방식은 간단하지만 장애에 취약하다. 소비자가 오프라인일 때의 처리가 힘듬

메세지 브로커

그래서 메시지 브로커(Message Queue)를 사용하는 거고, 이는 메세지 스트림을 처리하는데 최적화된 일종의 데이터베이스

메시지 브로커는 서버 생산자와 소비자는 서버의 클라이언트

브로커에 데이터가 모이기 때문에 클라이언트의 상태변경(접속, 접속해제, 장애)에 쉽게 대처할 수 있다

큐 대기를 하면 소비자는 비동기로 동작. 생산자는 큐에 잘 넣었는지만 확인하고 소비자는 고려하지 않음 메시지가 실제 소비자에게 전달되는 시점은 미래 시점이지만(아니면 순식간에 일수도) 때로는 큐에 백로그가 있다면 상당히 늦을 수 있다.

메세지 브로커와 데이터베이스의 비교

그럼 데이터베이스랑 뭐가 같고 다른가? 여기서 다루는 메시지 브로커는 AMQP,JMS스타일 브로커를 기준

  1. 데이터 보관
    1. 데이터베이스는 영구적으로 보관
    2. 메시지큐는 소비자에게 전달 성공하면 삭제
  2. 작업 집합 크기
    1. DB는 대용량 저장이 정상
    2. MQ는 큐가 짧게 유지된다고 가정
  3. 데이터 검색
    1. DB는 인덱스로 임의 조회
    2. MQ는 토픽을 구독하지 과거 메시지를 임의 조회는 불가
  4. 결과 관찰 방식
    1. DB는 클라이언트가 질의하고 서버가 대답
    2. MQ는 브로커가 소비자에게 푸시

  1. 전통적 메시지 브로커에서 “3분 전에 지나간 메시지를 다시 보고 싶다"가 안 되는 이유는? 그리고 이 한계를 어떤 방식이 해결하나요?
    1. 로그베이스 MQ가 해결을 해주긴 하는데, 사실 KTable같은거로 질의를 해야하는거 아닌가? 이걸 해결이라고 볼 수가 있나 모르겠네..
  2. “큐는 짧게 유지된다"는 가정이 깨지면(백로그 폭증) 어떤 문제가 생기나요?
    1. 전통 브로커 -> 디스크로 스필, 처리량 감소와 지연 증가, 큐가 한계면 drop or backpressure
    2. 로그베이스는 컨슈머 랙이 증가, 디스크 retention이 지나면 오래된 메세지 삭제되어 유실
  3. Push 모델이 앞서 정리한 “폴링(polling)의 문제"를 어떻게 해결하는지 한 문장으로?
    1. 생산 주체와 소비주체를 분리하고, 소비자는 비동기로 이벤트 처리만

소비자 측에서 하는 일. 여기서 KTable은 Kafka Streams에서 만드는 거

KTable, GlobalKTable은 추상개념이고, 실제로는 RocksDB에 저장됨. 디스크 기반이지만, LSM-Tree라는 자료구조쓰고, 읽기는 캐시써서 빠르다. 대안으로 InMemoryKeyValueStore 메모리기반.

다만 Kafka Streams에서의 KTable은 반쪽짜리 DB인 이유는 상태가 파티션별로 쪼개져있다. 그리고 SQL 부재, 세컨더리 인덱스 부재, 최종일관성 보장 x. 마치 샤딩, 근데 그 샤딩된 곳에서 각각 테이블을 띄우니 전역 테이블이 필요하다. 그래서 GlobalKTable으로 추상화한다

이거를 모두 다해준 게 ksqlDB인데…

이거의 쓸모가 무엇일까?

왜냐하면 대부분 그냥 소비자측에서 이벤트 받아서 DB에 저장하고 조회용은 DB를 OLTP로 쓰면 되는건데 말야.

  1. “스트림 처리 중 매 이벤트를 로컬에서 enrich하는 조인 입력”
  2. 단일 키 + 고볼륨 + 이벤트소싱/감사 중심 + saga 수용 가능.
    1. 예: SKU별 재고예약, 계정별 홀드, 주문 상태머신, 결제 dedup, fraud velocity

진실원은 RocksDB가 아니라 Kafka 로그. Streams는 단일 키 결정의 SoT는 되지만, 다중 SKU(cross-key→saga) + 외부 부작용(EOS 밖→idempotent executor) + 조회(파생 DB) 때문에 커머스 전체 SoT는 혼자 못 됨. → 대부분은 Postgres 단일 + outbox가 saner default.

복수 소비자(road balancing, fan out)

한 토픽에 소비자가 여러 개 붙을 때, 메시지를 어떻게 나눠줄 것인가?

  1. 로드밸런싱
    1. 여러 소비자 중 하나에게 전송
  2. 팬아웃
    1. 각 메시지를 모든 소비자에게 전송

Kafka에서는

  1. 그룹들 사이 = fan-out (결제 그룹도 전부 받고, 분석 그룹도 전부 받음)
  2. 그룹 내부 = load balancing (그룹 안에서 파티션을 워커들이 나눔)

클러스터 1 — n 토픽 1 —n 파티션 1—n레코드

클러스터 1—n 브로커 파티션 n—1 브로커(리더) 파티션 1—n replica 브로커 n—n 파티션

확인 응답과 재전송(ACK, retry)

소비자는 언제라도 장애가 날 수 있음 메시지를 받고, 처리를 다했고, ACK직전에 죽을 수도 있는거고 메시지를 받기도 전에 죽을 수도 있고 어느 순간에 다 죽을 수 있음

그래서 메시지를 잃어버리지 않기 위해서 메시지 브로커는 확인 응답(ACK)을 사용. 클라이언트는 메시지 처리가 끝났을 때, 브로커가 메시지를 큐에서 제거할 수 있게 브로커에게 명시적으로 알림

브로커는 ACK못받은 메세지(클라 연결 끊김, 타임아웃)가 미처리되었다고 가정하고 다른 소비자에게 재전송함 -> 이래서 멱등 처리가 필수적임!

부하 균형 분산과 결합할 때, 이런 재전송 행위는 메시지 순서에 영향을 미친다.

근데 왜 하나의 토픽 내에서 같은 키로 묶이는 이벤트들에서 메시지 순서를 보장하는 게 중요할까?

  1. 메시지 순서가 중요하지 않은 경우
    1. 서로 다른 엔티티
    2. 집계용 이벤트
    3. 멱등한 절대값 이벤트
    4. 읽기 전용 조회
  2. 메시지 순서가 중요한 경우
    1. 한 엔티티의 상태전이
      1. Created -> Paid -> Shipped -> Delivered인데, Paid가 먼저 컨슈머에 들어가면?
      2. 답: kafka 순서보장은 “같은 토픽 + 파티션 + 단일 프로듀서"일 때만
      3. 해결방법
        1. single writer: order-service만 발행, ?
          1. writer가 상태전이 검증 + 순서있게 발행
          2. outbox+cdc로 DB상태와 이벤트가 원자적
          3. 한 토픽+ order_id 키로 -> 파티션 순서 보장
          4. 컨슈마: status로 if/case 분기 타고, versioned_upsert로 멱등처리
        2. full-state+version(ECST): 전이 대신 현재 전체상태 발행 -> 컨슈머는 LWW upsert만, 전이 로직은 불필요
    2. CDC
      1. INSERT -> UPDATE -> DELETE
      2. DB에도 이 순서대로 진행해야함
    3. 증분, 누적 연산 (비멱등이고 순서를 의존)
      1. 증분은 고립상태로는 멱등화가 불가
        1. dedup-then-increment
        2. event-as-face
        3. 프레임워크 EOS
    4. 이벤트 간 선행 조건
      1. 한 이벤트가 발생한 이후에 다른 이벤트가 발생하는 구조
    5. 정정 취소 이벤트 AMQP는 왜 순서가 깨지나 공유 큐 + 경쟁 소비자 -> 여러 소비자가 한 큐를 나눠서 소비함 카프카는 파티션 1개= 소비자 1개+offset순차 멱등의 핵심은 연산을 재실행 안전하게 설계

자연멱등, versioned upsert, dedup

파티셔닝된 로그 (메세지를 저장하면서 이벤트 처리를 할 수는 없을까?)

Log-based message broker

  1. 데이터베이스의 지속성 있는 저장 방법
  2. 메시징 시스템의 지연 시간이 짧은 알림 기능

로그를 사용한 메시지 저장소(append-only and make event, partioning,offset)

  1. append only
  2. 디스크 하나를 쓰는 방법에서 처리량을 늘리려면 -> 파티션을 나눈다.
  3. 브로커는 모든 메시지에 offset이라는 단조 증가하는 순번을 부여 동일 파티션 내에서 보장하는 순서

로그 방식과 전통적인 메시지 방식의 비교 (메시지 순서가 중요한가?)

lbm(logbased)는 팬아웃 메시징 방식을 제공

리마인드

  1. 팬아웃 -> kafka는 컨슈머 그룹에서는 팬아웃
  2. 라운드로빈 -> kafka는 라운드로빈은 아니고, 프로듀서 레벨에서는 해쉬 키 기반 파티션 배정하고 , 컨슈머레벨에서는 컨슈머 그룹이 파티션을 소비자에게 배정, 결정

각 클라이언트는 할당된 파티션의 메시지를 모두 소비 소비자에게 파티션이 할당되면 소비자는 단일 스레드로 파티션에서 순차적으로 메시지를 읽음

불리한점

  1. 토픽 하나를 소비하는 작업을 공유하는 노드 수는 토픽의 파티션 수로 제한
  2. 특정 메시지 처리가 느리면 파티션 내 후속 메시지 처리가 지연됨

메시지 처리 비용이 비싸고(오래걸림) 메시지 단위로 병렬화 처리하고 싶고, 순서 안중요 -> AMQP 처리량 많고, 메시지 처리 속도가 빠름, 순서가 중요함 -> 로그 기반

소비자 오프셋 (log-based-message-broker는 오프셋만 변경한다)

파티션을 순차처리하면 메시지를 어디까지 처리했는지 알기 쉽다. 브로커는 소비자 오프셋만 기록하면 된다. -> 추적 오버헤드 감소 메시지 오프셋은 로그 순차번호인 lsn과 유사 장애가 나도, 다시 마지막 offset 이후부터 처리하면 돼. 장애전에 처리가 되고, offset commit이 안된 경우는 중복 처리되었을 가능성이 있지만 이는 멱등처리로 해결

디스크 공간 사용 (로그를 쌓다보면 디스크 공간을 사용하는데 어떻게 할것인가? ring buffer)

디스크 공간을 전부 사용하면 어떻게 할것인가? 원형버퍼, 링버퍼

  1. retention.ms 시간기준
  2. retention.bytes 크기기준

정리방식 cleanup.policy

  1. delete 시간/크기로 오래된것 삭제
  2. compact 키별 최신값만 남김

드물지만, 메시지 처리속도(소비자가 consume하는 속도)가 생산되는 속도를 따라잡지 못하면, 소비자가 뒤쳐져서 소비자 오프셋이 이미 삭제한 조각을 가리킬 수 있음 -> 메시지 유실

소비자가 생산자를 따라갈 수 없을 때 (consumer lag)

소비자가 메시지를 전송하는 생산자를 따라갈 수 없을때 선택지

  1. 메시지 drop
  2. buffering
  3. backpressure

lbmq는 고정 크기의 링 버퍼를 사용하는 버퍼링 형태 -> 지속되면 오래된 메시지를 버림

그래서 컨슈머랙을 관측하는게 필요함

오래된 메시지 재생 (replay)

AMQP는 메시지를 처리하는 행위는 메시지를 이후 버리기때문에 파괴적 연산 로그 기반 메시지 브로커는 메시지를 소비하는게 로그를 변화시키지 않는 읽기 전용 연산이기 때문이다

소비자 출력(처리)를 제외한 처리의 유일한 부수적 효과는 소비자 오프셋 이동이다.

소비자 오프셋은 소비자의 관리하에 있기 때문에, 어제 오프셋 기반으로 다른 위치에 replay해서 파생 데이터를 만들 수 있음

처리 코드를 수정해서 재처리하는게 가능

일괄처리랑 유사한 측면 변환처리를 반복해도 입력데이터에는 영향이 없고, 파생데이터만 만든다 많은 실험, 오류와 버그를 복구하기 쉽고, 조직내에서 데이터 플로를 통합하기 좋은 도구

데이터베이스와 스트림

그동안 DB의 아이디어(영속성)을 메시징에 적용했음 반대로 메시징과 스트림에서 아이디어를 가져와서 DB에 적용

사건, 이벤트라는 개념을 사용자 활동, 측정 판독에서 DB에 발생한 이벤트(UPDATE, INSERT, DELETE)까지 확장하자.

DB 내에서 발생하는 이벤트 -> 복제로그에 append-only로 저장됨 이거를 producer로 만들면 DB기록 이벤트를 생산할 수 있음 팔로워는 이를 구독, 소비하여 복제본에 완전히 동일한 데이터 복사본을 만들 수 있다.

시스템 동기화 유지하기.

Silver bullet은 없다.

  1. 데이터 저장
  2. 데이터 질의
  3. 데이터 처리

이 모두를 만족시키는 단일 시스템은 없음

-> 그래서 어쩔 수 없이 기능에 특화된 여러 시스템을 운용하는 수밖에 없는데, 그렇다보면 동기화가 필수임.

원시적인 해결법

  1. 데이터베이스 덤프
  2. dual write

이중 쓰기의 문제점

  1. race condition
  2. 한쪽 쓰기가 성공했지만, 다른 쓰기는 실패할 수 있음

공통 문제는 쓰기 지점이 여러개라는 것. 이거를 하나로 줄일 수 없을까?

변경 데이터 캡처 (CDC)

리더(leader) 데이터 변경 로그를 읽어 변경 내용을 스트림으로 제공

  1. oplog
  2. binlog
  3. etc… 원래는 CDC용이 아니라 내부 복제용 로그

다른 파생 데이터 시스템도 변경 스트림의 소비자로 만들어버림

변경 데이터 캡쳐의 구현

Publisher(Producer) -> 변경 사항을 캡처할 데이터베이스 하나를 리더로 둠 Subscriber(Consumer) -> SOT의 상태를 따라가고 싶은 데이터시스템들 (DW, ES,Cache etc..)

이는 MQ를 사용하므로, 비동기로 전환됨. 그리고 그 MQ는 순서보장(같은 키 범위 내에서)이 반드시 필요 A->B, B->C면 반드시 A->B가 먼저 처리되어야함 장점: 느린 소비자의 영향을 받지 않음 단점: 복제 지연 -> 파생 데이터 시스템은 일관성을 보장하기 어려움

초기 스냅샷

스냅샷은 특정 시점 DB 상태, LSN과 일치시켜야 함.

멈추고 못 하니 인터리빙

청크마다 lw 찍고, 그 청크 SELECT, hw 찍고.

그다음 lw~hw 사이에 로그에서 바뀐 행은 청크에서 빼고(로그가 최신=LWW), 안 바뀐 청크 행만 내보내서 스냅샷과 로그를 일치시킴.

이게 Netflix DBLog.

로그 컴팩션

이 소비자에게 “중간 이력"이 필요한가, “최신 상태"만 필요한가? 로그 중에서 키 기준으로 최신상태만 두고, 그 로그는 영속하게 저장해둠

log-compaction이 유리한 경우

최신 상태는 중요한데, 중간 이력은 중요하지 않은 경우

log-compaction을 쓰면 안되는 경우

변경 이력 전체가 중요한 경우 긴 retention(혹은 영구보존)

변경 스트림용 API 지원

DB가 각자 제공하는 changefeed API에서 CDC로 수렴하게 됨 제품은 사장, 개념은 CDC로 승계.

이벤트 소싱

이벤트는 불변(immutable)+ 추가만(append-only) -> 원장(ledger)를 생각하면 됨

이벤트 로그에서 현재 상태 파생하기

명령과 이벤트

==이벤트(Event)==와 ==명령(Command)==을 구분

  1. 어플리케이션은 명령이 실행가능한지 확인
  2. 무결성 검증, 명령 승인 -> 지속성 있는 불변 이벤트 =사실(fact)

상태와 스트림 그리고 불변성

트랜잭션 로그는 데이터베이스에 적용된 모든 변경사항을 기록한다. 로그는 고속으로 덧붙여지고, 덧붙이기가 로그를 변경하는 유일한 방법이다. 이런 측면에서 데이터베이스의 내용은 로그의 최근 레코드값을 캐시하고 있는 셈이다. 즉 로그가 진실이다. 데이터베이스는 로그의 부분집합의 캐시다. 캐시한 부분 집합은 로그로부터 가져온 각 레코드와 색인의 최신 값이다.

불변 이벤트의 장점

원장, 불변 개념이 주는 장점

  1. 감사 완전성
    1. 언제 뭐가 왜 바뀌었는지 불변 이벤트 자체로 남음
  2. 버그 복구
    1. 처리 로직에 버그 있었어도 이벤트는 그대로 남음. 로직 고치고 replay하면 됨 가변상태였으면 원본이 이미 덮어써져 복구 불가
  3. 같은 이벤트 로그에서 여러 read_view를 파생할 수 있음
    1. 몇번 씩 재처리해도 안전 왜냐? 그냥 불변 이벤트 로그로부터 읽어내면 되니까
  4. 읽기와 쓰기의 분리가 가능해짐 (CQRS)
    1. 쓰기는 그냥 append-only
    2. 읽기는 파생이어서 쓰기로부터 분리가능
      1. 장: 읽기, 쓰기 각각 독립 스케일링, 각자 목적에 맞게 설계가능해짐
      2. 단: 읽기 모델을 비동기로 만들때(컨슈머) 일관성 보장이 어려워짐.

동일한 이벤트 로그로 여러가지 뷰 만들기

동시성 제어

불변성 제어

솔직히 이벤트 소싱이라는 개념이 너무 이론적인 것 같아 건너뜀

스트림 처리

스트림 처리 유스케이스

  1. DB 변경분 파생 시스템에 동기화
  2. 이벤트를 사용자에게 직접 보냄. 사람이 스트림의 최종 소비자.
  3. ==하나 이상의 입력 스트림을 처리해 하나 이상의 출력 스트림을 만든다. ==

이번 장에서는 3번에 집중

스트림을 처리하는 코드 조각을 연산자(operator)나 작업(job)이라고 부른다.

  1. 유닉스 프로세스
  2. 맵리듀스

스트림 처리의 사용

복잡한 이벤트 처리

스트림 분석

구체화 뷰 유지하기

스트림 상에서 검색하기

메시지 전달과 RPC

시간에 관한 추론

이벤트 시간 대 처리 시간

준비 여부 인식

어쨌든 어떤 시계를 사용할 것인가?

윈도우 유형

스트림 조인

스트림 스트림 조인

스트렘 테이블 조인

테이블 테이블 조인

조인의 시간 의존성

내결함성

마이크로 일괄 처리와 체크포인트

원자적 커밋 재검토

멱등성

실패 후에 상태 재구축하기

정리