이벤트 스트림
고친 사람 github-actions[bot]
이벤트 스트림은 일어난 일을 한 건씩 순서대로 이어 붙여 여러 곳이 함께 읽게 합니다. 읽었다고 지우지 않아서 받는 쪽마다 제 속도로 따라 읽습니다. 서비스 사이에 변경을 건네거나 흩어진 데이터를 한곳으로 모을 때 이 흐름을 씁니다. 화면에 클릭이 잇달아 들어오는 흐름도 같은 이름으로 부릅니다. 이 편은 서비스 사이에 남겨 두는 기록을 다룹니다.
쉽고 빠른 이해
이벤트 스트림은 일어난 일을 순서대로 쌓아 둡니다. 여러 곳이 그 기록을 각자 읽어 갑니다. 주문 서비스가 「주문됨」「결제됨」「배송 시작됨」을 차례로 쌓으면 검색도 알림도 같은 기록을 따라 읽습니다.
이게 없으면 보내는 쪽이 받는 쪽마다 따로 보내야 합니다. 받는 쪽이 하나 늘 때마다 보내는 코드를 고칩니다. 읽으면 사라지는 방식이면 나중에 붙은 쪽은 지난 일을 볼 수 없습니다.
- 보내는 쪽이 일어난 일을 스트림 끝에 한 건 붙입니다
- 스트림은 붙은 순서대로 번호를 매겨 정해진 기간 동안 남겨 둡니다
- 받는 쪽은 어디까지 읽었는지 기억해 두고 그다음 건부터 이어 읽습니다
대가는 늦음과 운영 부담입니다. 받는 쪽은 늘 조금씩 뒤처집니다. 쓰자마자 받는 쪽에서 같은 값을 읽을 수는 없습니다. 기록을 남겨 두는 서버를 따로 띄우고 돌봐야 합니다.
상세
회사 메신저의 팀 채널을 떠올려 봅시다. 말은 올라온 순서대로 맨 아래에 붙습니다. 누가 읽었다고 사라지지 않습니다.
팀원은 저마다 읽다 만 데부터 이어 읽습니다. 나중에 채널에 들어온 사람은 위로 올려 지난 말을 봅니다. 이벤트 스트림은 이 채널처럼 움직입니다.
이벤트는 「주문 8871 결제됨」처럼 이미 일어난 일 하나를 데이터 한 건으로 적은 것입니다. 이벤트 스트림은 이런 이벤트를 일어난 순서대로 끝에 이어 붙인 흐름입니다. 영어로는 event stream 이라고 씁니다.
주문 서비스 하나를 예로 들어 봅니다. 손님이 주문하면 「주문됨」이 붙습니다. 결제가 끝나면 「결제됨」이, 상자가 나가면 「배송 시작됨」이 그 뒤에 붙습니다. 이 셋이 한 줄로 이어진 것이 주문 서비스의 이벤트 스트림입니다.
스트림에 이벤트를 붙이는 쪽을 생산자(producer)라고 부릅니다. 스트림에서 이벤트를 읽어 가는 쪽은 소비자(consumer)입니다. 생산자와 소비자는 서로를 모릅니다. 둘 다 스트림만 봅니다.
스트림을 보관하면서 쓰기와 읽기를 받아 주는 서버가 따로 있습니다. 이 서버를 브로커라고 부릅니다. 생산자와 소비자는 브로커에 붙어 쓰고 읽습니다.
이 절은 먼저 스트림이 이벤트를 어떻게 쌓고 소비자가 어떻게 읽는지 봅니다. 이어서 메시지 큐와 갈리는 대목, 지난 기록을 다시 읽는 일, 스트림을 여러 갈래로 나누는 법을 짚습니다. 끝으로 데이터 파이프라인에서 맡는 일과 치르는 대가를 봅니다.
끝에만 붙이는 기록
이벤트 스트림은 끝에만 붙입니다. 가운데에 끼워 넣지 않습니다. 이미 붙은 건을 고치지도 않습니다. 이렇게 끝에 붙이는 쓰기만 받는 기록을 추가 전용 로그라고 부릅니다.
고치지 않는 까닭은 이벤트가 이미 일어난 일이기 때문입니다. 결제가 취소됐다면 「결제됨」을 지우지 않습니다. 「결제 취소됨」을 하나 더 붙입니다.
붙는 건마다 차례대로 번호가 매겨집니다. 이 번호가 오프셋입니다. 첫 건이 0이면 다음 건은 1, 그다음은 2입니다.
| 오프셋 | 이벤트 | 주문 |
|---|---|---|
| 0 | 주문됨 | 주문 8871 |
| 1 | 결제됨 | 주문 8871 |
| 2 | 주문됨 | 주문 5020 |
| 3 | 배송 시작됨 | 주문 8871 |
표를 보면 서로 다른 주문의 이벤트가 한 스트림에 섞여 있습니다. 번호는 주문별로 매기지 않고 스트림에 붙은 차례로 매깁니다. 오프셋 3을 보면 0부터 2까지가 그보다 먼저 붙었다는 것을 압니다.
소비자마다 따로 기억하는 읽을 위치
소비자는 스트림을 한 번에 끝까지 읽지 않습니다. 한 건이나 몇 건씩 읽어 처리합니다. 처리가 끝나면 다음에 읽을 오프셋을 적어 둡니다.
이 위치는 소비자마다 따로 둡니다. 같은 주문 스트림을 소비자 둘이 읽는다고 해 봅시다. 알림 소비자는 주문 이벤트를 읽어 손님에게 알림을 보냅니다.
다른 하나는 검색 색인 소비자입니다. 검색 색인은 검색을 빠르게 하려고 따로 만들어 두는 저장소입니다. 이 소비자는 같은 주문 이벤트를 읽어 색인을 고칩니다.
flowchart TD
subgraph S["주문 이벤트 스트림"]
E0["0 · 주문됨"] --> E1["1 · 결제됨"]
E1 --> E2["2 · 주문됨"]
E2 --> E3["3 · 배송 시작됨"]
end
A["검색 색인 소비자"] -.->|다음에 읽을 건| E2
B["알림 소비자"] -.->|다음에 읽을 건| E3
P["생산자"] -->|새 건은 3 뒤에 붙는다| E3
그림에서 검색 색인 소비자는 다음에 2를 읽을 차례입니다. 알림 소비자는 3을 읽을 차례입니다. 한 소비자가 느려도 다른 소비자는 기다리지 않습니다. 생산자는 소비자가 어디쯤인지 모른 채 끝에 붙이기만 합니다.
소비자가 멈췄다 다시 뜨면 적어 둔 위치부터 이어 읽습니다. 멈춰 있던 동안 붙은 이벤트는 스트림에 남아 있으므로 놓치지 않습니다.
소비자가 읽은 위치와 스트림 끝 사이의 거리를 소비자 지연(consumer lag)이라고 부릅니다. 이 값이 계속 커지면 소비자가 붙는 속도를 못 따라가고 있다는 뜻입니다.
메시지 큐와 갈리는 대목
메시지 큐도 보내는 쪽과 받는 쪽 사이에서 건을 받아 둡니다. 겉모습이 비슷해서 둘을 자주 헷갈립니다. 갈리는 곳은 읽은 뒤에 그 건이 어떻게 되느냐입니다.
전형적인 메시지 큐는 한 건을 받는 쪽 하나에게만 건넵니다. 받는 쪽이 처리를 마쳤다고 알리면 큐는 그 건을 지웁니다. 일감을 여러 일꾼에게 나눠 맡기는 데 맞춘 방식입니다.
이벤트 스트림은 읽어도 지우지 않습니다. 같은 건을 소비자 여럿이 저마다 읽습니다. 지우는 시점은 누가 읽었느냐가 아니라 얼마나 오래 남겼느냐로 정합니다.
둘을 나란히 놓으면 이렇게 갈립니다.
| 메시지 큐 | 이벤트 스트림 | |
|---|---|---|
| 한 건을 받는 쪽 | 하나 | 스트림을 읽는 소비자마다 한 번씩 |
| 읽은 뒤 | 지운다 | 남겨 둔다 |
| 나중에 붙은 받는 쪽 | 지난 건을 못 본다 | 남아 있는 건을 처음부터 읽는다 |
| 어디까지 읽었나 | 큐가 건마다 기억한다 | 소비자가 오프셋 하나로 기억한다 |
할 일을 나눠 맡길 때는 큐를 고릅니다. 일어난 일을 여러 곳에 알리면서 기록으로도 남길 때는 스트림을 고릅니다. 두 방식을 함께 내놓는 제품도 있어서 경계가 늘 또렷하지는 않습니다.
지난 기록을 다시 읽기
스트림은 이벤트를 끝없이 들고 있지 않습니다. 디스크가 한없이 늘 수 없어서 오래된 건부터 지웁니다. 얼마나 남길지는 며칠이나 몇 주 같은 기간으로 정하거나 전체 크기로 정합니다. 이렇게 정한 기간을 보존 기간이라고 부릅니다.
보존 기간 안의 건은 몇 번이든 다시 읽을 수 있습니다. 소비자가 읽을 위치를 앞으로 되돌리기만 하면 됩니다. 이것을 다시 읽기, 영어로 replay 라고 부릅니다.
다시 읽기는 세 경우에 씁니다. 소비자 코드의 버그로 잘못 처리한 구간을 고쳐 다시 돌릴 때 씁니다. 새 소비자를 붙여 지난 이벤트로 제 저장소를 처음부터 채울 때도 씁니다. 받는 쪽 저장소가 망가졌을 때 스트림으로 다시 쌓기도 합니다.
기간으로 지우는 대신 다른 방식으로 줄이기도 합니다. 같은 주문 번호의 이벤트는 가장 최근 것 하나만 남기고 옛것을 지웁니다. 이 방식이 로그 컴팩션입니다.
파티션으로 나눠 쓰기
주문이 많아지면 스트림 하나를 서버 한 대가 다 받기 어렵습니다. 소비자 하나가 다 읽기에도 벅찹니다. 그래서 스트림 하나를 여러 갈래로 나눕니다. 이 갈래가 파티션입니다.
생산자는 이벤트마다 키를 하나 붙입니다. 키는 이 이벤트가 무엇에 관한 것인지 가리키는 값입니다. 주문 이벤트라면 주문 번호가 키입니다. 어느 파티션에 붙일지는 이 키로 정합니다.
키를 정해진 계산에 넣으면 수가 하나 나옵니다. 그 수로 붙일 파티션을 고릅니다. 이 계산이 해시입니다. 같은 키는 늘 같은 수가 나오므로 늘 같은 파티션으로 갑니다.
파티션을 나누면 읽는 일도 나눠 맡을 수 있습니다. 같은 일을 하는 소비자 여럿을 한 무리로 묶습니다. 이 무리가 소비자 그룹(consumer group)입니다. 파티션 하나는 무리 안에서 소비자 하나만 맡습니다.
앞에서 본 알림 소비자와 검색 색인 소비자는 하는 일이 달라 무리도 따로입니다. 무리마다 스트림 전체를 저마다 읽습니다. 아래 그림은 알림 소비자 무리 안에 소비자 1·2가 있는 모습입니다.
flowchart TD
P["생산자"] --> H["키로 파티션을 고른다"]
H -->|주문 8871| P0["파티션 0"]
H -->|주문 5020| P1["파티션 1"]
H -->|주문 3310| P2["파티션 2"]
subgraph G["알림 소비자 그룹"]
C1["소비자 1"]
C2["소비자 2"]
end
P0 --> C1
P1 --> C2
P2 --> C2
그림에서 주문 8871의 이벤트는 모두 파티션 0에 붙습니다. 파티션 안에서는 붙은 순서가 지켜집니다. 그래서 8871의 「주문됨」「결제됨」「배송 시작됨」은 소비자 1에게 그 순서로 닿습니다.
순서는 파티션 안에서만 지켜집니다. 파티션 0의 건과 파티션 1의 건 가운데 무엇이 먼저 읽힐지는 정해져 있지 않습니다. 그래서 순서가 중요한 이벤트끼리는 같은 키를 붙여 한 파티션에 모읍니다.
그림의 소비자 2처럼 무리 안의 소비자 하나가 파티션 둘을 맡기도 합니다. 파티션 수보다 소비자가 적으면 이렇게 됩니다.
끝나지 않는 흐름과 스트림 처리
파일은 끝이 있습니다. 다 읽은 뒤에 합계를 내면 됩니다. 이벤트 스트림은 새 건이 계속 붙으므로 끝이 없습니다.
그래서 스트림을 읽는 계산은 끝을 기다리지 않습니다. 「지난 1분 동안 들어온 주문 수」처럼 도착하는 대로 값을 고쳐 나갑니다. 끝나지 않는 흐름을 받아 쉬지 않고 계산하는 일을 스트림 처리라고 부릅니다.
그 반대편에 배치 처리가 있습니다. 데이터를 모아 두었다가 정해진 때에 한꺼번에 돌립니다. 밤마다 하루치 주문을 모아 매출을 내는 일이 배치 처리입니다.
데이터 파이프라인에서 맡는 일
데이터 엔지니어링에서 데이터가 처음 생기는 곳을 원천이라고 부릅니다. 그 데이터를 옮겨 받아 쓰는 곳은 대상입니다. 이벤트 스트림은 원천과 대상 사이에 섭니다. 원천은 스트림에 쓰고 대상은 스트림에서 읽습니다.
원천으로 흔한 것이 변경 데이터 캡처입니다. 데이터베이스에 생긴 변경을 한 건씩 뽑아 스트림에 붙입니다. 주문 서비스처럼 애플리케이션이 직접 이벤트를 붙이기도 합니다.
대상은 여럿이 같은 스트림을 저마다 읽습니다. 앞에서 본 검색 색인이 대상입니다. 분석용으로 데이터를 모아 두는 데이터 웨어하우스도 대상입니다. 자주 읽는 값을 가까이 복사해 두는 캐시도 스트림을 읽어 원본을 따라갑니다.
flowchart TD
DB["원천 데이터베이스"] --> CDC["변경 데이터 캡처"]
CDC --> S["이벤트 스트림"]
APP["주문 서비스"] --> S
S --> T1["검색 색인"]
S --> T2["데이터 웨어하우스"]
S --> T3["캐시"]
그림에서 원천과 대상 사이에는 선이 없습니다. 대상을 하나 더 붙여도 원천은 손대지 않습니다. 새 대상은 스트림에 남은 가장 오래된 건부터 읽거나 지금 끝부터 읽기 시작합니다.
같은 건이 두 번 오는 일
소비자는 이벤트를 처리하고 나서 읽을 위치를 적습니다. 처리와 위치 적기 사이에 소비자가 죽으면 문제가 생깁니다. 다시 떴을 때 이미 처리한 건을 한 번 더 읽습니다.
순서를 바꿔 위치를 먼저 적으면 반대가 됩니다. 위치를 적고 처리하기 전에 죽으면 그 건을 건너뜁니다. 두 실패 가운데 하나를 골라야 합니다. 대개 두 번 받는 쪽을 고릅니다.
한 건도 잃지 않는 대신 몇 건은 두 번 받는 방식을 최소 한 번 전달이라고 부릅니다. 그래서 소비자는 같은 건을 두 번 처리해도 결과가 같게 만듭니다. 같은 「결제됨」을 두 번 받아도 포인트는 한 번만 쌓습니다. 이미 반영한 이벤트인지 이벤트마다 붙여 둔 고유 번호로 먼저 확인합니다.
같은 일을 여러 번 해도 결과가 한 번 한 것과 같은 성질을 멱등성이라고 부릅니다. 스트림을 읽는 소비자는 대개 이 성질을 갖추도록 만듭니다.
치르는 대가
첫째는 늦음입니다. 생산자가 붙인 순간과 소비자가 반영한 순간 사이가 벌어집니다. 그 사이에 원천과 대상을 나란히 읽으면 값이 다릅니다. 시간이 지나면 같아지는 이런 상태를 최종 일관성이라고 부릅니다.
둘째는 이벤트 모양이 바뀌는 일입니다. 이벤트에 어떤 칸이 어떤 형식으로 들어가는지 정한 것을 스키마라고 부릅니다. 생산자가 칸을 하나 더하거나 빼면 옛 모양을 기대하던 소비자가 멈춥니다.
스트림에는 옛 모양의 건과 새 모양의 건이 섞여 남습니다. 다시 읽을 때도 둘 다 읽을 수 있어야 합니다. 그래서 스키마를 한곳에 모아 두고 바꿀 때마다 옛 모양과 맞는지 따지기도 합니다. 이 저장소가 스키마 레지스트리입니다.
셋째는 운영입니다. 브로커는 기록을 디스크에 남기는 서버입니다. 따로 띄우고 용량과 장애를 돌봐야 합니다. 이벤트가 많을수록 보존 기간만큼 저장 공간이 듭니다.
넷째는 흐름이 코드에 안 보인다는 것입니다. 생산자 코드를 읽어도 누가 그 이벤트를 읽는지 알 수 없습니다. 어느 소비자가 어느 스트림을 읽는지는 따로 적어 두어야 합니다.
받는 쪽이 하나뿐이고 곧바로 결과를 돌려받아야 하는 호출이라면 스트림을 거치지 않습니다. 상대 서비스를 직접 부르면 흐름이 호출 한 줄로 끝납니다. 스트림은 받는 쪽이 여럿이거나 앞으로 늘 때, 지난 일을 다시 읽어야 할 때 씁니다.
다른 분야에서 부르는 이벤트 스트림
이 이름은 서비스 사이의 데이터 전달 밖에서도 쓰입니다. 어느 분야든 이벤트가 끝없이 이어진다는 점은 같습니다. 갈리는 것은 지나간 건을 남겨 두느냐입니다.
| 분야 | 이벤트 스트림이 가리키는 것 | 지나간 건 |
|---|---|---|
| 서비스 사이 데이터 전달 | 브로커가 보관하는 이벤트 기록 | 보존 기간 동안 남는다 |
| 화면 프로그래밍 | 클릭과 키 입력이 잇달아 들어오는 흐름 | 지나가면 사라진다 |
| 브라우저와 서버 | 서버가 연결을 열어 둔 채 이벤트를 이어 보내는 통로 | 지나가면 사라진다 |
화면 프로그래밍에서는 입력 흐름에 거르기와 합치기 같은 연산을 걸어 값처럼 다룹니다. 이런 방식을 반응형 프로그래밍이라고 부릅니다. 브라우저와 서버 사이의 통로는 서버 전송 이벤트라는 이름으로 부릅니다.
관련 항목
이벤트 스트림을 이루는 구성 요소
이벤트 · 토픽 · 파티션 · 오프셋 · 이벤트 키 · 추가 전용 로그
이벤트 스트림에 쓰고 읽는 참여자
생산자 · 소비자 · 컨슈머 그룹 · 메시지 브로커 · 구독자
이벤트 스트림을 구현한 제품
Apache Kafka · Amazon Kinesis · Apache Pulsar · Redis Streams · Redpanda · NATS
이벤트 스트림과 맞세워지는 전달 방식
메시지 큐 · 발행-구독 · 이벤트 버스 · 웹훅 · 폴링 · 배치 처리
이벤트 스트림을 읽어 계산하는 처리 방식
스트림 처리 · 윈도우 · 워터마크 · 이벤트 시간 · Kafka Streams · Apache Flink
이벤트 스트림이 지켜야 하는 전달 성질
순서 보장 · 최소 한 번 전달 · 정확히 한 번 전달 · 멱등성 · 백프레셔 · 컨슈머 랙
이벤트 스트림이 지난 기록을 남기고 지우는 규칙
보존 기간 · 로그 컴팩션 · 툼스톤 · 재처리 · 체크포인트
이벤트 스트림을 중심에 두는 설계 방식
이벤트 기반 아키텍처 · 이벤트 소싱 · CQRS · 아웃박스 패턴 · 변경 데이터 캡처 · 구체화 뷰
이벤트 스트림이 놓이는 데이터 파이프라인 단계
원천 · 대상 · 수집 · 적재 · 파이프라인 · 데이터 엔지니어링
이벤트 스트림을 읽어 채우는 저장소
검색 색인 · 데이터 웨어하우스 · 데이터 레이크 · 캐시 · 읽기 모델
이벤트 스트림이 스키마 변경을 견디게 하는 수단
스키마 · 스키마 레지스트리 · 스키마 진화 · 하위 호환성 · 직렬화 · Avro
이벤트 스트림과 이름이 겹치는 헷갈리는 이웃
다른 이름: event stream · event streams · event streaming · 이벤트 스트리밍