스터디 · 2026-10-09 · 기초
실시간 데이터 파이프라인 가이드: 스트리밍·CDC·이벤트 기반 구조
실시간 데이터 파이프라인은 이벤트 로그(Kafka)에 변경을 쌓습니다. 스트림 처리 엔진(Flink·Structured Streaming·Dataflow)이 이를 계속 가공하고, CDC가 운영 DB 변경을 지연 없이 옮깁니다. 이 글은 공식 문서를 근거로 각 구성 요소의 역할, AI 서비스에서 쓰는 곳, 도입 순서와 점검 항목을 정리합니다.
- 배치는 모아 둔 데이터를 분·시간 단위로 한꺼번에 처리하고, 스트림은 들어오는 기록마다 초·밀리초 단위로 결과를 갱신한다.
- Kafka 같은 이벤트 로그는 읽은 뒤에도 이벤트를 지우지 않고 보관해, 생산자와 소비자를 떼어 놓고 여러 소비자가 다시 읽을 수 있게 한다.
- 로그 기반 CDC는 운영 DB의 트랜잭션 로그를 읽어 삭제까지 포함한 모든 변경을 원본 DB 부담 없이 실시간으로 옮긴다.
- 스트림 처리의 핵심 난제는 시간이다. 윈도·워터마크·트리거로 늦게 도착하거나 순서가 뒤바뀐 데이터를 다룬다.
- '이벤트 기반'은 알림, 상태 전달, 이벤트 소싱, CQRS처럼 서로 다른 패턴을 가리키므로 어느 패턴을 쓰는지 먼저 정해야 한다.
무엇이고 왜 필요한가: 배치와 스트림의 차이
데이터 파이프라인(데이터를 원천에서 목적지로 옮기며 가공하는 흐름)은 오랫동안 배치 방식이 기본이었다. AWS 설명 글글에 따르면 배치 처리는 일정 기간 모은 데이터나 정해진 양을 주기적으로 한꺼번에 처리한다. 전체 데이터를 여러 번 훑는 깊은 분석에 맞지만 지연이 분·시간·일 단위다. 스트림 처리는 끊임없이 들어오는 데이터를 받아 기록이 하나 도착할 때마다 분석 결과를 조금씩 갱신하며, 초나 밀리초 단위의 지연을 목표로 한다.
이 차이가 중요한 이유는 데이터의 가치가 시간이 지나면서 떨어지는 경우가 많기 때문이다. AWS 글글은 사용자의 현재 위치로 식당을 추천하는 예를 든다. 위치 데이터는 바로 쓰지 않으면 의미가 사라진다. Kafka 공식 소개글도 실시간 결제 처리, 차량·화물 추적, 공장 센서 분석, 환자 상태 모니터링을 이벤트 스트리밍의 대표 용도로 든다.
둘 중 하나만 골라야 하는 것은 아니다. AWS 글글은 많은 조직이 실시간 계층과 배치 계층을 함께 두는 혼합 모델을 쓴다고 설명한다. 스트림으로 즉시 필요한 결과를 뽑고, 같은 데이터를 저장소에 쌓아 두었다가 배치로 깊이 분석하는 식이다.
배치 처리 대 스트림 처리
배치 처리모아서 한꺼번에
- 전체 또는 묶음 데이터 대상
- 지연: 분·시간·일
- 전체를 여러 번 훑는 깊은 분석
스트림 처리도착할 때마다
- 최근 윈도나 최신 기록 대상
- 지연: 초·밀리초
- 윈도 단위로 조금씩 갱신
공통 같은 데이터를 스트림으로 즉시 쓰고, 저장해 배치로 다시 분석할 수 있음
구성 요소: 이벤트 로그, 처리 엔진, CDC
실시간 파이프라인은 크게 세 부분으로 이뤄진다. 첫째는 이벤트 로그다. Kafka 공식 소개글에 따르면 Kafka는 이벤트를 쓰고 읽는 기능, 원하는 기간 동안 이벤트를 안전하게 저장하는 기능, 발생 즉시 또는 나중에 처리하는 기능을 한데 묶은 플랫폼이다. 이벤트는 키·값·타임스탬프로 이뤄지고 토픽(폴더와 비슷한 묶음)에 쌓인다. 기존 메시징 시스템과 달리 소비한 뒤에도 지워지지 않으며, 토픽마다 정한 보관 기간이 지나야 버려진다. 토픽은 여러 파티션으로 나뉘고, 같은 키(예: 고객 ID)의 이벤트는 같은 파티션에 들어가 쓴 순서대로 읽힌다. 운영 환경에서는 복제 계수 3, 즉 데이터 사본 세 개를 두는 설정이 흔하다고 한다.
둘째는 스트림 처리 엔진이다. Flink 문서글는 상태를 유지하면서(stateful) 시간을 고려하는(timely) 스트림 처리를 핵심 개념으로 둔다. 가장 낮은 수준의 Process Function부터 DataStream API, 표 중심의 Table API, 가장 높은 수준의 SQL까지 여러 추상화 단계를 제공한다. Databricks 글글에 따르면 Spark의 Structured Streaming은 배치에서 쓰던 구조화 API를 그대로 스트리밍에 쓸 수 있게 해, 배치 작업을 거의 고치지 않고 스트리밍 작업으로 바꿀 수 있다. Google Cloud Dataflow 문서글는 Apache Beam SDK로 윈도·워터마크·트리거를 다루는 방식을 설명한다. Kafka 자체에도 집계·조인·윈도를 지원하는 Kafka Streams API가 있다글.
셋째는 CDC(Change Data Capture, 변경 데이터 캡처)다. Confluent 글글은 CDC를 데이터 원천의 모든 변경을 추적해 목적지 시스템에 반영하는 과정으로 정의한다. Debezium 문서글에 따르면 Debezium은 Kafka Connect용 소스 커넥터 모음이고, 데이터베이스별 기능을 이용해 변경을 읽어 Kafka로 보낸다. Kafka Connect에는 커뮤니티가 만든 커넥터가 수백 개 있어 대개 직접 만들 필요가 없다글.
설계 원칙: 시간과 이벤트를 어떻게 다룰까
스트림 설계에서 가장 먼저 부딪히는 문제는 시간이다. Dataflow 문서글는 데이터가 시간 순서대로, 예측 가능한 간격으로 도착한다는 보장이 없다고 지적한다. 그래서 끝없이 들어오는 데이터를 윈도(시간 구간)로 나눠 집계한다. 겹치지 않는 고정 구간인 텀블링 윈도, 일정 주기로 시작해 서로 겹치는 호핑 윈도(이동 평균에 적합), 활동이 끊긴 간격으로 나누는 세션 윈도가 있다. 워터마크는 한 윈도의 데이터가 다 도착했다고 보는 기준선이다. 이 선을 넘은 뒤 도착한 데이터는 '늦은 데이터'로 따로 처리한다. 트리거는 집계 결과를 언제 내보낼지 정하며, 이벤트 발생 시각(event time)과 처리 시각(processing time)을 구분해서 쓴다.
둘째 원칙은 생산자와 소비자를 떼어 놓는 것이다. Kafka 공식 소개글는 생산자가 소비자를 기다리지 않도록 둘을 완전히 분리한 것이 높은 확장성의 핵심 설계라고 밝힌다. AWS 글글도 생산자·브로커·소비자의 3단 구조를 설명하고, 소비자가 가공한 결과를 다시 브로커에 써서 새 스트림을 만들 수 있다고 덧붙인다.
셋째는 '이벤트 기반'이라는 말이 무엇을 뜻하는지 분명히 하는 것이다. Martin Fowler의 글글은 이를 네 가지로 나눈다. 이벤트 알림은 변경 사실만 알리는 방식으로, 결합도가 낮지만 여러 알림에 걸친 흐름이 코드에 드러나지 않는다. 이벤트 기반 상태 전달은 변경된 데이터를 이벤트에 실어 받는 쪽이 자기 사본을 유지하게 한다. 원본 시스템이 멈춰도 동작하지만 사본이 많아진다. 이벤트 소싱은 모든 상태 변경을 이벤트로 기록해 언제든 다시 재생해 상태를 복원한다. Git의 커밋 기록이 대표적인 예다. CQRS는 읽기와 쓰기에 서로 다른 데이터 구조를 둔다. Fowler글는 이 패턴들을 섞어 생각하면 실패 원인을 잘못 짚게 된다고 경고한다.
AI 서비스에서 쓰는 곳
실시간 피처: 추천이나 이상 탐지 모델은 '방금 일어난 일'을 입력으로 쓸 때 정확해진다. AWS 글글은 스트림 소비자가 필터링·집계·패턴 매칭·머신러닝을 수행한다고 설명한다. 클릭·주문 이벤트를 윈도로 집계해(예: 최근 구간의 이동 평균, w7) 피처 스토어에 계속 써 넣으면, 모델 서빙이 최신 피처로 예측할 수 있다.
RAG 색인 갱신: Confluent 글글은 운영 DB의 INSERT·UPDATE·DELETE를 캐시·검색 인덱스·데이터 웨어하우스·데이터 레이크로 전파하는 것을 CDC의 기본 쓰임으로 든다. RAG가 쓰는 벡터 DB도 검색 인덱스의 한 종류이므로 같은 구조를 적용할 수 있다. 문서 테이블의 변경을 CDC로 받아 바뀐 문서만 다시 청킹·임베딩하면, 전체를 주기적으로 다시 색인하지 않아도 된다. 특히 로그 기반 CDC는 삭제도 잡아내므로글 지워진 문서가 검색 결과에 남는 문제를 막는 데 도움이 된다.
모니터링: AWS 글글은 서버 로그가 시간 순서를 유지하므로 시간에 따른 변화를 보며 이상을 찾는 데 유용하다고 설명한다. 모델 요청·응답 로그를 이벤트로 흘려보내 윈도별로 오류율이나 지연을 집계하면, 배치 보고서보다 훨씬 빨리 데이터 드리프트나 장애 징후를 볼 수 있다.
단계별 적용 방법
처음 도입할 때는 다음 순서를 권한다. 첫째, 실시간이 정말 필요한 결과를 하나 고른다. 몇 분·몇 시간 늦어도 되는 일이면 배치로 충분하다글. AWS 글글도 기업들이 보통 시스템 로그 수집과 이동 최솟값·최댓값 계산 같은 단순한 용도부터 시작해 점점 정교한 실시간 처리로 넓혀 간다고 설명한다.
둘째, 이벤트 로그를 세우고 토픽과 키를 설계한다. 순서가 중요한 단위(고객·차량 ID)를 키로 정해야 같은 파티션 안에서 순서가 보장된다글. 보관 기간과 복제 설정도 이때 정한다.
셋째, 운영 DB 변경을 CDC로 연결한다. Debezium글은 커넥터를 처음 시작할 때 오래된 트랜잭션 로그가 이미 지워졌으면 현재 상태의 초기 스냅숏을 뜰 수 있고, 운영 중에도 증분 스냅숏을 실행할 수 있다. 필요한 스키마·테이블·열만 골라 받도록 필터를 걸고, 민감한 열은 마스킹한다.
넷째, 처리 엔진을 고른다. 이미 Spark 배치를 쓴다면 Structured Streaming으로 배치 작업을 프로토타입으로 만든 뒤 스트리밍으로 바꾸는 길이 짧다글. 상태 관리와 시간 처리를 세밀하게 다뤄야 하면 Flink의 DataStream API·Process Function을, 선언적으로 빨리 만들고 싶으면 Table API나 SQL을 쓴다글. 다섯째, 윈도·워터마크·늦은 데이터 정책을 정하고 결과를 목적지(피처 스토어, 벡터 DB, 웨어하우스)에 연결한다. 마지막으로 지연과 오류를 모니터링한다. Debezium 커넥터는 대부분 JMX로 모니터링할 수 있다글.
실시간 파이프라인 도입 순서
- 1용도 하나 고르기실시간이 꼭 필요한 결과부터
- 2이벤트 로그 세우기토픽·키·보관 기간·복제
- 3CDC 연결초기 스냅숏·필터·마스킹
- 4처리 엔진 선택Structured Streaming·Flink 등
- 5시간 정책 정하기윈도·워터마크·늦은 데이터
- 6목적지 연결·모니터링피처 스토어·벡터 DB·지연 감시
점검 목록
구성 단계에서는 다음을 확인한다. 순서가 필요한 단위가 이벤트 키로 정해졌는가글. 토픽별 보관 기간이 재처리·장애 복구에 충분한가글. 복제를 설정해 브로커 하나가 멈춰도 데이터가 남는가글. CDC가 로그 기반인지, 그리고 삭제와 이전 값까지 잡는지 확인한다글글.
처리 단계에서는 이벤트 시각과 처리 시각 중 무엇을 기준으로 집계할지 정했는지 본다글. 워터마크와 늦은 데이터 정책이 있는지, 처리 엔진의 상태가 장애 뒤에도 일관되게 복구되는지도 확인한다글. 데이터가 형식이 섞이거나 일부 빠진 채로 올 수 있으므로 유효성 검사 로직이 있는지 본다글.
운영·보안 단계에서는 개인정보가 담긴 열을 CDC 단계에서 마스킹하거나 빼는지글, 커넥터와 소비자의 지연을 지켜보고 있는지글를 점검한다. 롤백된 트랜잭션의 변경이 목적지에 잘못 반영되지 않는지글도 확인한다.
실시간 파이프라인 점검 항목
구성
- 순서가 필요한 단위가 키로 정해졌는가
- 보관 기간이 재처리에 충분한가
- 복제로 브로커 장애에 대비했는가
- CDC가 삭제·이전 값까지 잡는가
처리
- 이벤트 시각·처리 시각 기준을 정했는가
- 워터마크·늦은 데이터 정책이 있는가
- 장애 뒤 상태가 일관되게 복구되는가
- 형식 오류·누락을 검사하는가
운영·보안
- 민감한 열을 마스킹하는가
- 커넥터·소비자 지연을 감시하는가
- 롤백된 변경을 걸러 내는가
흔한 실수와 한계
첫째, 이벤트를 명령처럼 쓰는 실수다. Fowler글는 받는 쪽이 특정 행동을 해 주길 기대하면서 메시지를 이벤트처럼 꾸미는 것을 '수동 공격적 명령'이라 부른다. 이벤트 알림을 여러 단계로 엮으면 전체 흐름이 어떤 코드에도 드러나지 않아, 실행 중인 시스템을 관찰해야만 파악할 수 있게 된다.
둘째, 패턴을 뒤섞는 것이다. Fowler글는 '이벤트 소싱이 재앙이었다'는 프로젝트 사례가 실제로는 CQRS나 과도한 비동기 통신 문제였을 수 있다고 지적한다. CQRS는 읽기·쓰기 모델을 따로 두는 만큼 복잡해지므로, 동료 상당수가 이를 경계한다고도 적었다. 이벤트 소싱은 외부 시스템과 얽힌 이벤트를 재생하기 어렵고, 이벤트 스키마가 시간에 따라 바뀌는 문제를 다뤄야 한다.
셋째, CDC 방식을 잘못 고르는 것이다. Confluent 글글에 따르면 타임스탬프 열 방식은 단순하지만 실제 DELETE를 잡지 못하고 원본 스키마를 바꿔야 한다. 트리거 방식은 모든 변경을 잡지만 쓰기 때마다 추가 쓰기가 생겨 원본 DB 성능을 떨어뜨린다. 로그 기반은 원본 부담이 적지만, 트랜잭션 로그 형식이 벤더마다 달라 이후 버전에서 바뀔 수 있다. 또 롤백된 변경을 걸러 내야 한다.
넷째, 확장을 나중에 생각하는 것이다. Dataflow 문서글는 Streaming Engine을 쓰지 않는 스트리밍 작업은 처음 지정한 최대 워커 수를 넘어 확장할 수 없다고 안내한다. 처리 엔진마다 이런 운영상 제약이 있으므로 도입 전에 문서를 확인해야 한다.
더 공부하려면
처음이라면 Kafka 공식 소개글부터 읽는다. 이벤트·토픽·파티션·복제 개념이 짧게 정리돼 있어 이후 글을 이해하는 바탕이 된다. 그다음 AWS 설명 글글에서 배치와 스트림의 차이 표와 생산자·브로커·소비자 구조를 확인한다.
운영 DB를 연결할 계획이라면 Confluent의 CDC 글글로 세 가지 CDC 방식의 장단점을 비교하고, Debezium 기능 문서글로 스냅숏·필터·마스킹 옵션을 살펴본다. 시간 처리는 Dataflow 문서글의 윈도·워터마크·트리거 설명이 그림과 함께 이해하기 쉽다.
처리 엔진을 고를 때는 Flink 개념 문서글와 Databricks의 Structured Streaming 소개글를 비교해 읽는다. 시스템 설계 단계에 들어가면 Martin Fowler의 글글로 우리 팀이 말하는 '이벤트 기반'이 네 패턴 중 무엇인지 정리해 둔다.
관련 용어
데이터 파이프라인 피처 스토어 벡터 DB 임베딩 파이프라인 RAG 이상 탐지 데이터 품질 데이터 레이크하우스
참고 문헌
참고 자료(공식 문서·엔지니어링 글)
- Introduction | Apache Kafka (kafka.apache.org)
- Overview | Apache Flink (nightlies.apache.org)
- Debezium Features :: Debezium Documentation (debezium.io)
- What Is Change Data Capture (CDC)? | Confluent (confluent.io)
- What Is Streaming Data? - Streaming Data Explained - AWS (aws.amazon.com)
- What do you mean by “Event-Driven”? (martinfowler.com)
- Streaming pipelines | Cloud Dataflow | Google Cloud Documentation (cloud.google.com)
- What is Structured Streaming? | Databricks (databricks.com)
참고 영상·강의
영상은 본문의 근거로 쓰지 않았고, 더 알아보려는 분을 위한 링크입니다.
AI가 참고 문헌을 바탕으로 작성하고 검수를 거친 해설입니다. 정확한 내용은 원문을 확인해 주세요.


