데이터 메시·패브릭과 스트림 처리: 주문 이벤트 중복·지연·재처리 설계
1장. 스트림을 다시 돌렸더니 매출이 두 배가 됐다#
주문 O1의 결제 금액이 25,000원이라고 하겠습니다.
결제 이벤트가 스트림에 들어옵니다.
event_id = e100
order_id = O1
amount = +25,000처리기가 이 이벤트를 읽고 오늘 매출에 25,000원을 더했습니다.
오늘 매출
25,000원그 직후 장애가 발생했습니다.
문제는 결과는 저장됐지만 이 이벤트를 처리했다는 위치 정보가 아직 확정되지 않은 경우입니다.
처리기가 재시작됩니다.
같은 이벤트 e100을 다시 읽습니다.
그리고 또:
+25,000을 적용합니다.
결과:
오늘 매출
50,000원이 됩니다.
원래 주문은 한 건입니다.
이벤트를 다시 읽었을 뿐인데 업무 사건이 두 번 발생한 것처럼 집계됐습니다.
스트림 처리에서 중요한 원칙은 이것입니다.
이벤트를 다시 읽는 것과 업무 사건을 다시 발생시키는 것은 같은 일이 아니다.
2장. 실시간 시스템에서는 재전송이 예외가 아니라 정상 상황이다#
메시지 시스템에서 같은 이벤트가 다시 전달되는 이유는 다양합니다.
처리기 장애
네트워크 타임아웃
확인 응답 유실
체크포인트 실패
소비자 재시작
수동 재처리따라서 다음과 같은 가정을 하면 위험합니다.
이벤트는 항상 정확히 한 번만 도착한다.보다 현실적인 가정은:
같은 이벤트가
0번, 1번 또는 여러 번 관찰될 수 있다.입니다.
그래서 이벤트 소비자는 중복이 들어와도 같은 최종 상태에 수렴하도록 설계해야 합니다.
3장. 데이터 메시는 조직 구조와 데이터 책임의 문제다#
데이터 메시는 단순한 데이터 저장 기술 이름이 아닙니다.
핵심은 데이터를 만드는 도메인이 자신이 제공하는 데이터를 소비 가능한 데이터 제품으로 책임진다는 데 있습니다.
예를 들어 주문 도메인이 있다고 하겠습니다.
주문팀은 다음 의미를 가장 잘 알고 있습니다.
주문 생성이 무엇인가?
결제 완료는 언제인가?
취소와 환불은 어떻게 다른가?
부분 취소는 어떻게 표현하는가?
같은 주문의 이벤트 순서는 어떻게 해석하는가?이런 의미를 중앙 데이터팀이 추측하게 두지 않고 도메인이 명확한 계약으로 제공하는 것이 데이터 메시의 중요한 관점입니다.
4장. 주문 데이터 제품에는 데이터 파일만 있으면 되는 것이 아니다#
주문팀이 다음 파일을 올렸다고 하겠습니다.
orders.json이것만으로는 다른 팀이 신뢰하고 사용할 수 있는 데이터 제품이라고 보기 어렵습니다.
최소한 다음 정보가 필요합니다.
소유자
스키마
필드 정의
이벤트 ID 의미
갱신 주기
품질 기준
접근 권한
변경 정책
문의 경로예를 들어:
status = CANCELLED이 어떤 상황을 의미하는지 문서가 없다면 소비자는 매출에서 빼야 하는지 판단하기 어렵습니다.
5장. 주문 이벤트의 계약을 예로 보면#
| 항목 | 예 |
|---|---|
| 제품 이름 | order_item_event |
| 소유자 | 주문 도메인 |
| 업무 단위 | 주문 품목 사건 |
| 식별자 | event_id |
| 주문 식별자 | order_id |
| 이벤트 시각 | occurred_at |
| 수집 시각 | ingested_at |
| 스키마 버전 | schema_version |
| 중복 기준 | 동일 event_id |
| 문의 책임 | 주문 데이터 담당자 |
이 정도 정보가 있어야 이벤트를 다른 시스템에서 안전하게 사용할 수 있습니다.
6장. 데이터 패브릭은 데이터를 찾고 연결하는 문제에 초점을 둔다#
조직 안에 다음 데이터가 있다고 하겠습니다.
주문
상품
고객
물류
결제
광고각 시스템이 서로 다른 장소에 존재합니다.
사용자는 묻습니다.
주문 데이터는 어디에 있습니까?
고객ID 정의는 무엇입니까?
이 매출 보고서는 어느 원천에서 만들어졌습니까?
데이터 패브릭은 이런 환경에서:
검색
카탈로그
메타데이터
계보
접근 권한
연결
통합을 제공하는 접근으로 이해할 수 있습니다.
7장. 데이터 메시와 데이터 패브릭은 경쟁 관계일 필요가 없다#
다음처럼 역할을 나눌 수 있습니다.
주문 도메인#
취소 이벤트의 의미를 정의
스키마 변경 책임
품질 기준 관리공통 데이터 플랫폼#
카탈로그
접근 제어
계보
수집
모니터링즉:
flowchart LR
D["도메인<br/>데이터 의미·품질 책임"] --> P["데이터 제품"]
F["공통 플랫폼<br/>카탈로그·계보·접근"] --> P
P --> C["소비자"]처럼 함께 사용할 수 있습니다.
8장. 메시라고 각 팀이 마음대로 표준을 정하는 것은 아니다#
도메인 자율성을 다음처럼 이해하면 곤란합니다.
주문팀은 고객ID를 숫자로 사용
배송팀은 이메일을 고객ID로 사용
결제팀은 전화번호를 고객ID로 사용이런 구조에서는 데이터를 연결하기 어렵습니다.
그래서 데이터 메시에서도 공통 규칙이 필요합니다.
예:
고객 식별자 형식
날짜·시간 기준
통화 코드
개인정보 규칙
스키마 버전 규칙이런 공통 기준을 연방 거버넌스 관점으로 설계할 수 있습니다.
9장. 패브릭도 단순 데이터 가상화만을 뜻하지 않는다#
데이터 패브릭을:
데이터를 움직이지 않고
가상으로 연결하는 기술이라고만 설명하면 범위가 좁아집니다.
실제 통합 환경에서는 필요에 따라:
가상화
수집
복제
물리적 이동
메타데이터 관리
계보
정책 적용등이 함께 포함될 수 있습니다.
핵심은 데이터가 어디에 있든 발견하고 이해하고 신뢰하며 접근할 수 있는 연결 구조를 만드는 데 있습니다.
10장. 조직 구조가 좋아도 중복 이벤트는 따로 해결해야 한다#
데이터 메시를 도입했습니다.
데이터 패브릭도 구축했습니다.
그래도 주문팀이 같은 결제 이벤트를 두 번 발행하면:
e100
+25,000
e100
+25,000매출 계산기는 두 번 받을 수 있습니다.
패브릭이 이 두 이벤트의 계보를 추적할 수는 있습니다.
하지만 매출 집계에서 한 번만 반영하는 책임은 별도의 처리 규칙이 필요합니다.
조직 아키텍처와 이벤트 처리 정확성은 다른 문제입니다.
11장. 실시간 매출을 계산하려면 먼저 지표를 정의해야 한다#
오늘의 매출을 다음처럼 정의한다고 하겠습니다.
순매출
=
결제액
-
취소·환불액그러면 이벤트를 최소한 다음처럼 구분할 수 있습니다.
PAYMENT
→ 양수
REFUND
→ 음수예:
e100
PAYMENT
+25,000
e101
REFUND
-5,000최종 순매출:
20,000입니다.
12장. 이벤트 시간과 처리 시간은 서로 다르다#
환불이 실제로 발생한 시각:
10:05시스템에 도착한 시각:
10:20이라고 하겠습니다.
두 개의 시간이 있습니다.
Event Time
→ 실제 업무 사건 발생 시각
Processing Time
→ 처리 시스템이 사건을 다룬 시각둘을 구분하지 않으면 시간대별 매출이 달라질 수 있습니다.
13장. 처리 시간으로 계산하면 10시 5분 환불이 10시 20분 매출에 들어갈 수 있다#
10분 단위 창을 사용한다고 하겠습니다.
10:00~10:10
10:10~10:20
10:20~10:30환불 이벤트:
발생
10:05
도착
10:20입니다.
Processing Time 기준으로 집계하면:
10:20~10:30에 들어갈 수 있습니다.
하지만 실제 업무 발생 시점은:
10:00~10:10입니다.
14장. 이벤트 시간 기준 집계는 실제 업무 시간을 복원하려는 것이다#
이벤트에 다음 정보가 있다고 하겠습니다.
event_id = e101
occurred_at = 10:05
received_at = 10:20
amount = -5,000업무 발생 기준 보고서라면:
occurred_at을 사용합니다.
따라서 이 환불은:
[10:00, 10:10)창의 결과를 수정합니다.
15장. 워터마크는 이벤트 시간 처리가 어디까지 진행됐는지를 나타내는 신호다#
이벤트는 발생 순서대로 정확하게 도착하지 않습니다.
예:
10:01 이벤트
10:04 이벤트
10:03 이벤트
10:09 이벤트따라서 시스템은 어느 시점의 이벤트까지 대체로 도착했다고 판단할 기준이 필요합니다.
워터마크는 이벤트 시간 처리 진행 정도를 표현하는 데 사용됩니다.
개념적으로:
Watermark = 10:10이라면 시스템은 10:10보다 이전 이벤트가 대부분 처리됐다고 판단할 수 있습니다.
하지만 늦은 이벤트가 절대로 더 오지 않는다는 증명은 아닙니다.
16장. 워터마크와 허용 지연은 같은 개념이 아니다#
워터마크:
현재 이벤트 시간 진행 추정허용 지연:
창이 닫힌 뒤
얼마나 늦은 이벤트까지 받을 것인가입니다.
예를 들어:
Window
10:00~10:10
Allowed lateness
15분이라면 10:05 이벤트가 10:20에 도착해도 정책에 따라 10시 창을 수정할 수 있습니다.
17장. 늦은 이벤트를 어떻게 처리할지는 업무 정책이다#
창이 이미 닫혔습니다.
그 뒤 환불이 도착합니다.
선택지는 여러 가지입니다.
무시
기존 결과 수정
별도 보정 이벤트 기록
다음 보고서에서 수정
오류 큐로 이동어느 것이 정답인지는 업무에 따라 달라집니다.
실시간 대시보드와 회계 확정 보고서는 같은 정책을 사용할 필요가 없습니다.
18장. Tumbling Window는 겹치지 않는 시간 창이다#
1시간 Tumbling Window라면:
12:00~13:00
13:00~14:00처럼 서로 겹치지 않습니다.
12:40 이벤트는:
12:00~13:00창 하나에만 포함됩니다.
19장. Sliding Window는 하나의 이벤트가 여러 창에 들어갈 수 있다#
창 길이:
1시간이동 간격:
30분이라고 하겠습니다.
창:
12:00~13:00
12:30~13:3012:40 이벤트는 두 창에 모두 들어갑니다.
따라서:
한 이벤트
=
한 창이라고 가정하면 안 됩니다.
20장. Session Window는 사용자 활동 간격으로 경계를 만든다#
고객 A의 이벤트:
10:00
10:05
10:12그리고 다음 이벤트:
11:00이 있다고 하겠습니다.
비활동 기준이 20분이라면 앞의 세 이벤트는 하나의 세션으로 묶이고 11시 이벤트는 새로운 세션이 될 수 있습니다.
세션 창은 고정된 시계 기준이 아니라 활동 간격을 기준으로 만들어집니다.
21장. Lambda 아키텍처는 빠른 결과와 재계산 경로를 분리한다#
Lambda 구조를 단순화하면 다음과 같습니다.
flowchart TD
E["원본 이벤트"] --> B["Batch Layer"]
E --> S["Speed Layer"]
B --> V["배치 결과"]
S --> R["실시간 결과"]
V --> Q["Serving"]
R --> Q실시간 스트림은 빠른 결과를 제공합니다.
배치는 원본을 다시 읽어 정확한 결과를 재계산할 수 있습니다.
22장. Lambda에서 가장 어려운 문제는 두 계산 경로의 규칙을 맞추는 것이다#
실시간 처리에서는 환불을:
-5,000으로 처리합니다.
그런데 배치 코드에서는:
주문 상태가 취소면
원매출 제외방식으로 계산한다고 하겠습니다.
두 방식이 업무적으로 같은 결과를 만들지 검증해야 합니다.
코드가 두 벌이면 시간이 지나면서 규칙이 달라질 수 있습니다.
23장. 배치 결과와 실시간 결과의 경계도 명확해야 한다#
배치가:
어제 23:59:59까지계산했습니다.
실시간 계층은:
오늘 00:00 이후만 제공하면 경계가 명확합니다.
하지만 배치가 오늘 오전까지 다시 계산했는데 기존 스트림 결과를 그대로 더하면 같은 주문을 두 번 셀 수 있습니다.
따라서:
배치 확정 범위
스트림 유효 범위를 명확하게 관리해야 합니다.
24장. Kappa 아키텍처는 동일한 스트림 처리 경로를 재사용한다#
Kappa 구조를 단순화하면:
flowchart LR
L["보존된 이벤트 로그"] --> P["Stream Processor"]
P --> O["결과"]새 이벤트도 같은 처리 경로를 사용합니다.
과거 데이터를 고칠 때도 보존된 이벤트를 다시 재생합니다.
처리 코드가 하나라는 점은 장점이 될 수 있습니다.
25장. Kappa에서도 재처리는 공짜가 아니다#
지난달 데이터를 다시 계산하려면 지난달 이벤트가 남아 있어야 합니다.
재처리 범위
30일인데 로그 보존 기간이:
7일이라면 지난달 이벤트를 재생할 수 없습니다.
따라서 Kappa 구조에서는 로그 보존 기간이 재처리 요구와 연결됩니다.
26장. 코드 버전도 재처리 결과를 바꾼다#
9월에 사용한 매출 규칙:
배송완료 주문만 매출10월에 규칙이 변경되었습니다.
결제완료 주문부터 매출9월 이벤트를 현재 코드로 다시 재생하면 9월 당시 보고서와 다른 결과가 나올 수 있습니다.
따라서 재처리에서는:
원본 이벤트
처리 코드 버전
참조 데이터 버전
지표 정의를 함께 고려해야 합니다.
27장. 재처리 결과를 기존 집계에 더하면 가장 쉽게 두 배가 된다#
기존 9월 매출:
100,000,000원9월 이벤트를 다시 재처리했습니다.
새 결과:
100,000,000원이 결과를 기존 값에 더하면:
200,000,000원이 됩니다.
재처리는 새 업무 사건을 추가하는 것이 아니라 같은 업무 사실을 다시 계산하는 작업일 수 있습니다.
28장. 재처리 결과는 별도 버전으로 만들어 교체하는 방식이 안전할 수 있다#
예:
sales_2026_09_v1
→ 기존 결과
sales_2026_09_v2
→ 재처리 결과새 결과를 검증한 뒤 조회 대상을 v2로 바꿉니다.
flowchart LR
V1["기존 결과 v1"] --> C["검산"]
V2["재처리 결과 v2"] --> C
C --> S["v2로 조회 전환"]이 방식은 기존 값에 증분으로 더해서 중복되는 문제를 줄일 수 있습니다.
29장. Exactly Once라는 표현은 범위를 명확히 해야 한다#
다음 설명은 위험합니다.
Kafka를 사용하면
Exactly Once다.무엇이 정확히 한 번인지 범위를 확인해야 합니다.
입력 읽기
상태 변경
출력 토픽 쓰기
DB UPDATE
이메일 발송
결제 API 호출모두 같은 장애 경계를 가지지 않습니다.
30장. 이벤트가 네트워크에서 한 번만 전송된다는 뜻도 아니다#
Exactly Once를:
같은 메시지는 물리적으로
절대 두 번 전달되지 않는다.로 이해하면 안 됩니다.
일반적으로 중요한 것은 여러 번 읽거나 재시도되더라도 논리적 결과를 한 번 반영한 것과 같은 상태를 만드는 것입니다.
이를 위해:
체크포인트
트랜잭션
이벤트 ID
멱등 처리등을 조합할 수 있습니다.
31장. 외부 이메일이나 결제 API는 별도의 멱등성 설계가 필요하다#
스트림 처리기가 이벤트 e100을 읽고 이메일을 전송했습니다.
그 뒤 체크포인트 전에 장애가 났습니다.
재시작 후 e100을 다시 읽습니다.
이메일을 다시 보내면 고객에게 두 통이 갑니다.
메시지 시스템 내부의 처리 보장이 외부 이메일 시스템까지 자동으로 확장되는 것은 아닙니다.
32장. 이벤트 ID가 중복 제거의 핵심 기준이 될 수 있다#
다음 이벤트를 사용하겠습니다.
event_id = e100
order_id = O1
type = PAYMENT
amount = 25,000같은 이벤트가 다시 도착했습니다.
event_id = e100이벤트 ID가 안정적이라면 이미 처리했는지 확인할 수 있습니다.
33장. 재시도마다 새로운 이벤트 ID를 만들면 중복 제거가 어렵다#
첫 전송:
event_id = e100
refund = -5,000재전송:
event_id = e999
refund = -5,000업무상 같은 환불인데 ID가 바뀌었습니다.
소비자 입장에서는 서로 다른 두 사건처럼 보일 수 있습니다.
따라서 이벤트 ID는 전송 시도가 아니라 업무 사건을 안정적으로 식별하는 것이 좋습니다.
34장. 같은 이벤트 ID에 다른 값이 들어오면 단순 중복으로 숨기면 안 된다#
첫 이벤트:
e101
-5,000두 번째:
e101
-8,000입니다.
ID는 같지만 내용은 다릅니다.
이를:
이미 처리한 이벤트
→ 무시라고 하면 실제 데이터 오류를 숨길 수 있습니다.
이 경우:
원본 해시 비교
충돌 기록
정정 이벤트 생성
담당자 알림같은 별도 처리 정책이 필요합니다.
35장. 처리 원장을 만들어 이벤트를 한 번만 인정할 수 있다#
PostgreSQL에서 다음과 같은 원장을 만들 수 있습니다.
CREATE TABLE revenue_event_ledger (
event_id text PRIMARY KEY,
order_id text NOT NULL,
occurred_at timestamptz NOT NULL,
delta_amount numeric(12, 2) NOT NULL
);이벤트 ID를 기본키로 사용합니다.
36장. 결제와 환불을 원장에 넣어 보자#
INSERT INTO revenue_event_ledger VALUES
(
'e100',
'O1',
TIMESTAMPTZ '2026-09-01 10:02:00+09',
25000
),
(
'e101',
'O1',
TIMESTAMPTZ '2026-09-01 10:05:00+09',
-5000
)
ON CONFLICT (event_id) DO NOTHING;현재 데이터:
e100
+25,000
e101
-5,000입니다.
37장. 같은 환불이 다시 들어와도 새로운 행을 만들지 않는다#
INSERT INTO revenue_event_ledger VALUES
(
'e101',
'O1',
TIMESTAMPTZ '2026-09-01 10:05:00+09',
-5000
)
ON CONFLICT (event_id) DO NOTHING;e101은 이미 존재합니다.
새 행이 만들어지지 않습니다.
38장. 원장에서 직접 합계를 계산하면 중복 이벤트가 매출을 두 번 올리지 않는다#
SELECT
order_id,
COUNT(*) AS accepted_events,
SUM(delta_amount) AS net_amount
FROM revenue_event_ledger
GROUP BY order_id
ORDER BY order_id;결과:
| order_id | accepted_events | net_amount |
|---|---|---|
| O1 | 2 | 20,000 |
재전송된 e101은 새로운 행이 아니므로 결과에 다시 반영되지 않습니다.
39장. 처리 원장과 별도 집계를 둘 때는 두 변경의 원자성을 고민해야 한다#
이벤트 원장:
e101 저장 완료하지만 대시보드 집계를 수정하기 전에 장애가 났다고 하겠습니다.
상태:
원장
→ 환불 처리됨
대시보드
→ 아직 25,000불일치입니다.
반대 상황도 가능합니다.
대시보드
→ 20,000으로 갱신
원장
→ 저장 실패재전송되면 다시 -5,000을 적용할 수도 있습니다.
40장. 같은 데이터베이스라면 원장과 집계를 하나의 트랜잭션으로 묶을 수 있다#
개념적으로:
BEGIN
e101 원장 INSERT
대시보드 매출 -5,000
COMMIT처럼 처리할 수 있습니다.
하지만 데이터베이스가 서로 다르거나 외부 시스템까지 포함되면 문제는 복잡해집니다.
그때는 이벤트 기반 재구축, 아웃박스, 멱등 소비 같은 다른 전략이 필요할 수 있습니다.
41장. 늦게 도착한 환불을 실제 시간 창에 적용해 보자#
결제 이벤트:
e100
발생 10:02
도착 10:02
+25,000환불 이벤트:
e101
발생 10:05
도착 10:20
-5,000시간 창:
[10:00, 10:10)입니다.
업무 발생 시간 기준으로 두 이벤트는 같은 창에 속합니다.
42장. 환불이 반영되기 전의 대시보드#
10:10에 창이 처음 계산됐습니다.
당시 도착한 이벤트는 e100뿐입니다.
순매출
25,000원입니다.
43장. 10시 20분에 환불이 늦게 도착했다#
허용 지연이:
10:25까지라고 하겠습니다.
e101은 10:20에 도착했습니다.
따라서 정책상 아직 반영할 수 있습니다.
새 결과:
25,000 - 5,000
=
20,000원입니다.
대시보드 값은 수정됩니다.
44장. 같은 환불이 10시 21분에 다시 도착했다#
두 번째 이벤트도:
event_id = e101입니다.
원장에 이미 있으므로 다시 차감하지 않습니다.
최종 순매출은 계속:
20,000원입니다.
45장. 늦은 이벤트와 중복 이벤트는 서로 다른 문제다#
e101의 첫 도착:
발생
10:05
도착
10:20은 늦은 이벤트 문제입니다.
e101이 10:21에 다시 오는 것은 중복 이벤트 문제입니다.
둘을 같은 개념으로 처리하면 안 됩니다.
Late
→ 어느 시간 결과를 고칠 것인가?
Duplicate
→ 이미 반영한 사건인가?질문이 다릅니다.
46장. 취소가 먼저 도착하고 주문 생성이 나중에 올 수도 있다#
분산 시스템에서는 순서도 바뀔 수 있습니다.
원래 업무 순서:
e200
ORDER_CREATED
e201
ORDER_CANCELLED그런데 네트워크 지연으로:
e201
먼저 도착
e200
나중 도착할 수 있습니다.
단순히 도착 순서대로 상태를 적용하면 문제가 생길 수 있습니다.
47장. 상태 버전을 함께 보내면 오래된 이벤트를 구분할 수 있다#
예:
ORDER_CREATED
version = 1
ORDER_CANCELLED
version = 2취소 version 2가 먼저 도착했습니다.
현재 주문 상태:
version = 2
cancelled나중에 생성 version 1이 도착합니다.
1 < 2이므로 오래된 상태로 판단해 무시할 수 있습니다.
48장. 델타 이벤트는 상태 이벤트보다 중복 처리에 더 민감할 수 있다#
상태 이벤트:
balance = 80
version = 5는 같은 상태를 반복 적용해도 결과가 크게 달라지지 않을 수 있습니다.
하지만 델타 이벤트:
balance -= 20를 두 번 적용하면:
100 → 80 → 60이 됩니다.
실제 의도가 한 번 차감이라면 잘못된 결과입니다.
따라서 델타 이벤트에서는 안정적인 이벤트 ID가 특히 중요합니다.
49장. 주문 건수도 지표 정의에 따라 취소 처리 방법이 달라진다#
지표 A:
주문 생성 건수라면 주문이 나중에 취소돼도 생성 사건은 사라지지 않습니다.
ORDER_CREATED
→ +1
ORDER_CANCELLED
→ 생성 수에는 영향 없음지표 B:
유효 주문 수라면:
ORDER_CREATED
→ 1
ORDER_CANCELLED
→ 0이 될 수 있습니다.
같은 취소 이벤트라도 지표마다 역할이 다릅니다.
50장. 취소 이벤트를 무조건 -1로 더하면 음수가 될 수도 있다#
취소 이벤트가 생성 이벤트보다 먼저 도착합니다.
단순 집계:
현재 0
취소
-1이면:
유효 주문 수
-1이 됩니다.
업무적으로 말이 되지 않을 수 있습니다.
따라서 상태 기반 지표에서는 주문별 현재 상태를 확인하는 것이 더 안전할 수 있습니다.
51장. 생성 뒤 취소·취소 뒤 생성·중복 취소를 모두 시험해야 한다#
동일한 최종 업무 결과를 기대하는 세 입력 순서입니다.
정상 순서#
생성
↓
취소역순 도착#
취소
↓
생성중복#
생성
↓
취소
↓
취소최종 상태는 모두:
cancelled로 수렴해야 할 수 있습니다.
이런 테스트가 스트림 처리 안정성을 검증합니다.
52장. 중복 제거 기록의 보존 기간은 재처리 범위보다 짧으면 위험하다#
중복 처리 기록을:
7일동안 보존한다고 하겠습니다.
그런데 운영자가:
지난 30일 이벤트 전체 재처리를 실행합니다.
8일 이전 이벤트의 처리 기록은 이미 삭제되어 있습니다.
동일한 event_id라도 새 이벤트처럼 처리될 수 있습니다.
53장. 재처리 요구와 원장 보존 기간을 맞춰야 한다#
재처리 가능한 범위:
90일이라면 중복 방지 정보도 최소한 그 범위의 재처리 정책과 맞아야 합니다.
또는 재처리 방식 자체를:
기존 결과를 비우고
전체를 다시 계산하는 방식으로 설계할 수 있습니다.
중요한 것은:
기존 결과에 덧붙이는 재처리와:
전체 결과를 교체하는 재계산을 섞지 않는 것입니다.
54장. 전체 재계산과 증분 보정은 다른 전략이다#
전체 재계산#
9월 기존 결과 제거
9월 원본 전체 재처리
새 9월 결과 생성증분 보정#
오류 이벤트만 식별
정정 이벤트 추가
기존 결과 조정각각 장단점이 있습니다.
전체 재계산은 단순하지만 비용이 클 수 있습니다.
증분 보정은 효율적이지만 정정 로직이 복잡해질 수 있습니다.
55장. 원본 이벤트 보존은 재처리의 보험이다#
집계 테이블만 있고 원본 이벤트가 없다면 계산 규칙이 잘못됐을 때 다시 만들기 어렵습니다.
예:
지난달 매출 계산식 오류 발견원본 이벤트가 남아 있으면 새 코드로 다시 계산할 수 있습니다.
flowchart LR
E["원본 이벤트"] --> V1["기존 코드"]
E --> V2["수정 코드"]
V1 --> O1["기존 결과"]
V2 --> O2["새 결과"]두 결과를 비교할 수 있습니다.
56장. 하지만 원본만 보존하면 충분한 것도 아니다#
이벤트를 다시 재생했는데 결과가 과거와 다릅니다.
원인은 외부 참조 데이터일 수 있습니다.
예:
상품 분류
환율
고객 등급
세금 규칙이 값이 현재 기준으로 바뀌었다면 같은 이벤트도 다른 결과를 만들 수 있습니다.
재현성을 위해서는 참조 데이터의 버전도 필요할 수 있습니다.
57장. 데이터 제품에는 품질 기준도 포함되어야 한다#
예를 들어 주문 이벤트 제품의 품질 기준을 다음처럼 정의할 수 있습니다.
event_id NULL 0%
order_id NULL 0%
중복 event_id 0%
24시간 내 도착률 99.9%
금액 음수는 환불 타입에서만 허용이렇게 해야 소비자가 어느 수준까지 데이터를 신뢰할 수 있는지 알 수 있습니다.
58장. 스키마 변경도 소비자 계약 문제다#
기존 이벤트:
{
"order_id": "O1",
"amount": 25000
}새 이벤트:
{
"order_id": "O1",
"gross_amount": 25000,
"currency": "KRW"
}amount를 갑자기 삭제하면 기존 소비자가 깨질 수 있습니다.
따라서:
스키마 버전
호환 기간
폐기 일정같은 변경 정책이 필요합니다.
59장. 데이터 계보는 숫자가 어디서 왔는지 보여준다#
대시보드:
10시 순매출
20,000원이 숫자를 다음으로 추적할 수 있어야 합니다.
flowchart TD
D["Dashboard<br/>20,000"] --> A["집계 테이블"]
A --> P["Stream Processor v3"]
P --> E1["e100 +25,000"]
P --> E2["e101 -5,000"]이 구조가 있으면 숫자 변경 이유를 설명하기 쉽습니다.
60장. 데이터 메시와 스트림 처리의 책임이 만나는 지점#
주문팀은 다음을 정의해야 합니다.
event_id의 의미
취소 이벤트 의미
이벤트 순서
정정 방법
스키마 변화플랫폼팀은 다음을 제공할 수 있습니다.
이벤트 전달
카탈로그
권한
계보
모니터링
재처리 도구소비팀은:
어떤 버전을 사용할지
어떤 지표를 계산할지
재처리 후 결과를 어떻게 검증할지책임집니다.
61장. 데이터 제품의 계약에 재처리 조건까지 넣는 것이 좋다#
예:
| 항목 | 정책 |
|---|---|
| 원본 보존 | 180일 |
| event_id 보장 | 동일 업무 사건에 동일 ID |
| 순서 | order_id 안에서 version 증가 |
| 중복 가능성 | 있음 |
| 최대 정상 지연 | 30분 |
| 늦은 이벤트 | 별도 보정 허용 |
| 정정 이벤트 | 새 ID + corrected_event_id |
| 스키마 호환 | 최소 90일 |
| 재처리 문의 | 주문 데이터 팀 |
이 정보가 있으면 소비자가 안전한 재처리 전략을 만들기 쉽습니다.
62장. 실시간과 정확성은 반대말이 아니다#
다음 식의 설명은 지나치게 단순합니다.
스트림
→ 빠르지만 부정확
배치
→ 느리지만 정확스트림도 올바른 이벤트 시간·중복 제거·상태 관리가 있으면 정확한 결과를 만들 수 있습니다.
배치도 잘못된 원본이나 잘못된 집계 규칙을 사용하면 틀립니다.
정확성은 처리 방식 이름이 아니라 업무 의미와 실패 처리 설계에 달려 있습니다.
63장. Lambda와 Kappa를 다시 비교하면#
| 항목 | Lambda | Kappa |
|---|---|---|
| 실시간 경로 | 별도 Speed Layer | 동일 스트림 경로 |
| 과거 재계산 | 배치 경로 | 로그 재생 |
| 코드 경로 | 둘 이상일 수 있음 | 상대적으로 단일화 가능 |
| 주요 부담 | 배치·스트림 규칙 일치 | 로그 보존과 재생 비용 |
| 결과 교체 | 배치와 실시간 경계 관리 | 재처리 결과 버전 관리 |
어느 구조가 절대적으로 더 정확한 것은 아닙니다.
64장. 스트림 처리 체크리스트#
- 이벤트 ID는 업무 사건을 안정적으로 식별하는가?
- 같은 이벤트가 여러 번 와도 안전한가?
- 같은 ID에 다른 내용이 오면 감지하는가?
- Event Time과 Processing Time을 구분하는가?
- 시간 창은 어떤 기준으로 닫는가?
- 허용 지연은 얼마인가?
- 늦은 이벤트가 오면 결과를 수정하는가?
- 이벤트 순서가 뒤집힐 수 있는가?
- 상태 버전을 가지고 있는가?
- 체크포인트와 외부 결과 저장 사이 장애를 고려했는가?
- 외부 API 호출은 멱등한가?
- 재처리 범위보다 로그 보존 기간이 짧지 않은가?
- 중복 제거 기록의 보존 기간은 충분한가?
- 재처리 결과를 기존 값에 더하는가 교체하는가?
- 처리 코드와 참조 데이터 버전을 재현할 수 있는가?
65장. 데이터 제품 운영 체크리스트#
소유권#
이 데이터의 의미를
최종적으로 누가 설명하는가?품질#
어떤 오류율과 지연을 허용하는가?변경#
스키마 변경을 소비자에게
어떻게 알리는가?접근#
누가 어떤 필드까지 볼 수 있는가?계보#
보고서 숫자에서 원본 이벤트까지
추적할 수 있는가?재처리#
지난달 결과를 다시 만들 수 있는가?66장. 장애를 재현하는 테스트 시나리오를 만들어야 한다#
정상 이벤트만 보내서 테스트하면 부족합니다.
다음 입력을 준비하는 것이 좋습니다.
정상 결제
동일 결제 2회 전달
결제 후 환불
환불 2회 전달
환불이 결제보다 먼저 도착
30분 늦게 온 이벤트
재처리 후 같은 이벤트 재등장
처리 완료 직후 체크포인트 실패이 모든 경우에서 최종 업무 결과가 올바른지 검증합니다.
67장. 가장 중요한 검증은 “같은 업무 사실로 수렴하는가”다#
예를 들어 최종 업무 사실이:
주문 O1
결제 25,000
환불 5,000
최종 순매출 20,000이라고 하겠습니다.
입력 순서가 달라도:
결제 → 환불
환불 → 결제
결제 → 결제 중복 → 환불
결제 → 환불 → 환불 중복최종 결과는 정책에 따라 동일한 20,000원으로 수렴해야 합니다.
스트림 시스템의 신뢰성은 여기에서 드러납니다.
68장. 핵심 정리#
데이터 메시·데이터 패브릭·스트림 처리는 서로 다른 문제를 다룹니다.
데이터 메시는:
누가 데이터 의미와 품질을 책임하는가?를 다룹니다.
데이터 패브릭은:
데이터를 어떻게 찾고,
연결하고,
접근하고,
계보를 추적할 것인가?에 초점을 둡니다.
스트림 처리는:
이벤트가 늦고,
중복되고,
순서가 바뀌고,
다시 재생될 때
어떻게 같은 업무 결과를 만들 것인가?를 다룹니다.
실시간 이벤트 처리에서 가장 먼저 구분해야 하는 시간은:
Event Time
→ 실제 업무 사건 시각
Processing Time
→ 시스템 처리 시각입니다.
워터마크는 이벤트 시간 진행 상태를 판단하는 신호이고, 허용 지연은 이미 닫힌 시간 창에 늦은 사건을 얼마나 더 받아들일지 정하는 정책입니다.
중복 이벤트 문제에서는 안정적인 event_id가 중요합니다.
같은 업무 사건
→ 같은 event_id가 되어야 재전송을 식별하기 쉽습니다.
그리고:
Late Event와:
Duplicate Event는 서로 다른 문제입니다.
늦은 이벤트는:
어느 과거 결과를 수정할 것인가?
의 문제입니다.
중복 이벤트는:
이미 반영한 사건인가?
의 문제입니다.
Lambda 아키텍처에서는 실시간 경로와 배치 재계산 경로의 규칙과 경계를 맞춰야 합니다.
Kappa 아키텍처에서는 같은 처리 경로를 재사용할 수 있지만 과거 이벤트를 다시 읽을 수 있을 만큼 로그를 오래 보존해야 합니다.
어느 구조를 사용하든 재처리 결과를 기존 집계에 무조건 더하면 같은 사건을 두 번 계산할 수 있습니다.
따라서 재처리는:
기존 결과 교체
또는
명확한 멱등 증분 보정중 어떤 방식인지 분명해야 합니다.
또한 Exactly Once라는 표현도 범위를 명확하게 해야 합니다.
메시지 시스템에서 한 번의 논리적 결과를 보장하는 기능이 있다고 해도:
외부 DB
이메일
결제 API
검색 인덱스까지 자동으로 한 번만 처리되는 것은 아닙니다.
최종적으로 좋은 스트림 데이터 제품은 다음을 설명할 수 있어야 합니다.
같은 사건이 두 번 와도 왜 한 번만 반영되는가?
늦게 온 사건은 어느 시간의 결과를 어떻게 수정하는가?
지난달 데이터를 다시 계산해도 왜 같은 업무 사실로 수렴하는가?
이 세 질문에 답할 수 있어야 데이터 메시의 책임 구조, 데이터 패브릭의 연결 구조, 스트림 처리의 재처리 구조가 실제 운영 가능한 데이터 플랫폼으로 이어집니다.