提交时机导致消息丢失——回退与幂等处理
한국어 원문으로 표시합니다.
이 실습은 VM 에서 돕니다
우분투 VM 에 Apache Kafka 4.3.1 이 KRaft 단일 노드로 떠 있습니다
(localhost:9092). 도우미 ship-process 는 배송 처리기를 흉내 냅니다 — 표준
입력의 주문을 전부 읽은 뒤 하나씩 /root/kafka/processed.txt 에 적다가,
SHIP_FIXED=1 이 아니면 order-2 에서 죽습니다. 처음 뜨는 데 4분쯤 걸립니다.
목표
컨슈머가 처리보다 먼저 오프셋을 커밋하면 크래시 뒤에 메시지가 "사라지는" 것을
재현하고, 그 메시지가 Kafka 에는 그대로 있음을 다른 그룹으로 확인한 뒤, 그룹
오프셋을 되감아 다시 처리합니다. 되감기가 만드는 중복은 멱등 컨슈머로 막고,
auto.offset.reset 이 새 그룹의 첫 위치를 어떻게 정하는지 봅니다.
왜 중요한가
설계 문서의 '메시지 전달 의미론' 절이 이 실습의 대본입니다. 읽고 → 위치를
저장하고 → 처리하면, 처리 중 죽었을 때 그 메시지는 다시 오지 않습니다
(at-most-once). 읽고 → 처리하고 → 저장하면, 저장 전에 죽었을 때 다시
옵니다(at-least-once). 콘솔 컨슈머처럼 자동 커밋(enable.auto.commit 기본
true, 5초 간격)에 기대는 처리기는 앞의 모양에 가깝습니다. "한 번 사라졌다" 는
대개 이것이고, 고치면 "두 번 도착" 이 됩니다 — 그래서 처리기는 멱등해야 합니다.
문서는 그것을 "메시지에 기본 키가 있어 갱신이 멱등한 경우" 라고 적습니다.
단계
- 토픽
shipments를 파티션 1개로 만들고order-1부터order-5까지 다섯 줄을 넣으세요. - 그룹
ship-svc로 처음부터 세 건을 읽어ship-process에 파이프하세요(죽습니다). 그 뒤ship-svc를 describe 해/root/kafka/ship-crash.txt에 저장하세요. CURRENT-OFFSET 은 3 인데/root/kafka/processed.txt에는order-1하나뿐이어야 합니다 —order-2,order-3이 "사라진" 것입니다. - 새 그룹
audit-svc로 처음부터 다섯 건을 읽어/root/kafka/ship-audit.txt에 저장하세요. Kafka 에는 전부 남아 있습니다. ship-svc의 오프셋을 가장 앞으로 되감고(--reset-offsets --to-earliest --execute) 그 출력을/root/kafka/ship-reset.txt에 저장하세요.SHIP_FIXED=1로ship-svc가 다섯 건을 다시 읽어ship-process에 파이프하세요.processed.txt는 여섯 줄이 되고order-1이 두 번 있어야 합니다 — 되감기의 대가인 중복입니다./root/kafka/dedup.sh를 만드세요. 표준 입력의 주문을 읽어/root/kafka/seen.txt에 없는 것만/root/kafka/processed-dedup.txt에 적고 seen 에 기록합니다(환경변수DEDUP_SEEN·DEDUP_OUT이 있으면 그 경로).shipments를 처음부터 두 번 읽어 파이프해도 다섯 줄만 남아야 합니다.--from-beginning없이 새 그룹late-svc로 5초 읽고(0건),auto.offset.reset=earliest로 새 그룹early-svc로 다섯 건을 읽으세요./root/kafka/offset-reset.txt에late_count=0,early_count=5두 줄을 쓰세요./root/kafka/consumer-report.txt에lost_after_crash=<2단계에서 사라진 건수>,duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>,unique_orders=<processed-dedup.txt 의 줄 수>세 줄을 쓰세요.
참고
- 그룹으로 읽기:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic shipments --group ship-svc --from-beginning --max-messages 3 --timeout-ms 8000 | ship-process - 되감기:
kafka-consumer-groups.sh ... --reset-offsets --group ship-svc --topic shipments --to-earliest --execute. 운영 문서대로 컨슈머가 멈춰 있어야 합니다.--execute없이 돌리면 계획만 보여 줍니다. - 새 그룹의 첫 위치: 컨슈머 설정 문서의
auto.offset.reset— 기본latest라 그룹에 오프셋이 없으면 지금 이후만 읽습니다.--command-property auto.offset.reset=earliest로 바꿉니다.--from-beginning은 콘솔 도구가 같은 일을 해 주는 단축키입니다. - 흔한 실수 1: 2단계에서
--max-messages를 빼는 것. 다섯 건을 다 읽어 버리면 '사라진 두 건' 이 재현되지 않습니다. - 흔한 실수 2: dedup 의 상태(seen)를 메모리에만 두는 것. 프로세스가 죽으면 상태도 죽어 다음 실행이 다시 중복을 냅니다. 파일이든 DB 든 처리 결과와 함께 남아야 합니다.
배송 주문 다섯 건
토픽 shipments 를 파티션 1개로 만들고 order-1 부터 order-5 까지 다섯 줄을 넣으세요.
printf 'order-1\norder-2\norder-3\norder-4\norder-5\n' | kafka-console-producer.sh .... 파티션이 하나라 순서가 전부 지켜집니다.
처리 전에 커밋하면 사라진다
그룹 ship-svc 로 처음부터 세 건을 읽어 ship-process 에 파이프하세요(죽습니다). 그 뒤 ship-svc 를 describe 해 /root/kafka/ship-crash.txt 에 저장하세요. CURRENT-OFFSET 은 3 인데 /root/kafka/processed.txt 에는 order-1 하나뿐이어야 합니다 — order-2, order-3 이 "사라진" 것입니다.
--group ship-svc --from-beginning --max-messages 3 --timeout-ms 8000 | ship-process. 콘솔 컨슈머는 세 건을 넘긴 뒤 오프셋 3 을 커밋하고 끝나는데, 처리기는 두 번째에서 죽었습니다. 다음에 이 그룹으로 읽으면 3 부터 시작합니다.
Kafka 에는 그대로 있다
새 그룹 audit-svc 로 처음부터 다섯 건을 읽어 /root/kafka/ship-audit.txt 에 저장하세요. Kafka 에는 전부 남아 있습니다.
--group audit-svc --from-beginning --max-messages 5 --timeout-ms 8000 > /root/kafka/ship-audit.txt. 소비는 삭제가 아닙니다 — 사라진 것은 메시지가 아니라 ship-svc 의 위치입니다.
그룹 오프셋을 되감는다
ship-svc 의 오프셋을 가장 앞으로 되감고(--reset-offsets --to-earliest --execute) 그 출력을 /root/kafka/ship-reset.txt 에 저장하세요.
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --reset-offsets --group ship-svc --topic shipments --to-earliest --execute. 컨슈머 위치가 정수 하나라서 가능한 일입니다 — 설계 문서는 이것을 큐의 계약을 어기지만 꼭 필요한 기능이라고 부릅니다.
다시 처리하면 두 번 도착한다
SHIP_FIXED=1 로 ship-svc 가 다섯 건을 다시 읽어 ship-process 에 파이프하세요. processed.txt 는 여섯 줄이 되고 order-1 이 두 번 있어야 합니다 — 되감기의 대가인 중복입니다.
... --group ship-svc --max-messages 5 --timeout-ms 8000 | SHIP_FIXED=1 ship-process. 그룹에 오프셋이 있으면 --from-beginning 은 무시됩니다 — 되감은 위치(0)에서 읽습니다. 사라졌던 두 건은 돌아오지만 이미 처리한 order-1 도 다시 옵니다.
멱등 컨슈머
/root/kafka/dedup.sh 를 만드세요. 표준 입력의 주문을 읽어 /root/kafka/seen.txt 에 없는 것만 /root/kafka/processed-dedup.txt 에 적고 seen 에 기록합니다(환경변수 DEDUP_SEEN·DEDUP_OUT 이 있으면 그 경로). shipments 를 처음부터 두 번 읽어 파이프해도 다섯 줄만 남아야 합니다.
grep -qxF "$line" "$SEEN" 으로 본 적이 있는지 확인하고, 없을 때만 두 파일에 적습니다. 그룹 없이 --from-beginning --max-messages 5 로 두 번 읽어 파이프하세요. 채점기는 임시 경로를 환경변수로 주고 중복이 섞인 입력을 넣어 봅니다.
새 그룹은 어디서 시작하나
--from-beginning 없이 새 그룹 late-svc 로 5초 읽고(0건), auto.offset.reset=earliest 로 새 그룹 early-svc 로 다섯 건을 읽으세요. /root/kafka/offset-reset.txt 에 late_count=0, early_count=5 두 줄을 쓰세요.
--group late-svc --timeout-ms 5000 은 아무것도 못 읽고 끝납니다(기본 latest). --group early-svc --command-property auto.offset.reset=earliest --max-messages 5 --timeout-ms 8000 은 다섯 건을 읽습니다. wc -l 로 세어 파일에 적으세요.
사라진 것과 두 번 온 것을 센다
/root/kafka/consumer-report.txt 에 lost_after_crash=<2단계에서 사라진 건수>, duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>, unique_orders=<processed-dedup.txt 의 줄 수> 세 줄을 쓰세요.
사라진 건수는 ship-crash.txt 의 CURRENT-OFFSET 에서 그때 처리된 건수(1)를 뺀 것, 중복 건수는 processed.txt 의 줄 수에서 서로 다른 주문 수를 뺀 것입니다. 채점기는 같은 파일들에서 다시 셉니다.