コミットのタイミングで一件が消える — 巻き戻しと冪等処理
このラボはVM上で動きます
UbuntuのVMに、Apache Kafka 4.3.1がKRaftの単一ノードとして起動しています
(localhost:9092)。ヘルパーship-processは、配送ハンドラーを真似ます。標準
入力の注文をすべて読んだ後、1つずつ/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秒間隔)に頼るハンドラーは、前の形に近いものです。「1回消えた」は
たいていこれで、直すと「2回届く」になります。そのため、ハンドラーは冪等でなければなりません。
ドキュメントは、それを「メッセージに主キーがあって更新が冪等になる場合」と書いています。
ステップ
shipmentsトピックを、パーティション1つで作成し、order-1からorder-5までの5行を入れてください。- グループ(
ship-svc)で、最初から3件を読んで、ship-processにパイプしてください(落ちます)。その後、ship-svcをdescribeして、/root/kafka/ship-crash.txtファイルに保存してください。CURRENT-OFFSETは3なのに、/root/kafka/processed.txtにはorder-1の1件だけがなければなりません。order-2とorder-3が「消えた」のです。 - 新しいグループ(
audit-svc)で、最初から5件を読んで、/root/kafka/ship-audit.txtファイルに保存してください。Kafkaには、すべて残っています。 ship-svcのオフセットを一番前まで巻き戻して(--reset-offsets --to-earliest --execute)、その出力を、/root/kafka/ship-reset.txtファイルに保存してください。SHIP_FIXED=1で、ship-svcが5件を再び読んで、ship-processにパイプしてください。processed.txtは6行になり、order-1が2回なければなりません。巻き戻しの代償である、重複です。/root/kafka/dedup.shファイルを作成してください。標準入力の注文を読んで、/root/kafka/seen.txtファイルにないものだけを、/root/kafka/processed-dedup.txtファイルに書き込み、seenに記録します(環境変数DEDUP_SEEN・DEDUP_OUTがあれば、そのパス)。shipmentsを最初から2回読んでパイプしても、5行だけが残らなければなりません。--from-beginningなしで、新しいグループ(late-svc)で5秒読み(0件)、auto.offset.reset=earliestで、新しいグループ(early-svc)で5件を読んでください。/root/kafka/offset-reset.txtファイルに、late_count=0、early_count=5の2行を書いてください。/root/kafka/consumer-report.txtファイルに、lost_after_crash=<2단계에서 사라진 건수>、duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>、unique_orders=<processed-dedup.txt 의 줄 수>の3行を書いてください(プレースホルダーは、ステップ2で消えた件数、ステップ5の後のprocessed.txtの重複件数、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を外すことです。5件をすべて読んでしまうと、「消えた2件」が再現されません。 - よくあるミス2: dedupの状態(seen)を、メモリだけに置くことです。プロセスが落ちると状態も消えて、次の実行がまた重複を出します。ファイルでもDBでも、処理結果と一緒に残らなければなりません。
配送注文5件
shipmentsトピックを、パーティション1つで作成し、order-1からorder-5までの5行を入れてください。
printf 'order-1\norder-2\norder-3\norder-4\norder-5\n' | kafka-console-producer.sh ...。パーティションが1つなので、順序がすべて守られます。
処理の前にコミットすると消える
グループ(ship-svc)で、最初から3件を読んで、ship-processにパイプしてください(落ちます)。その後、ship-svcをdescribeして、/root/kafka/ship-crash.txtファイルに保存してください。CURRENT-OFFSETは3なのに、/root/kafka/processed.txtにはorder-1の1件だけがなければなりません。order-2とorder-3が「消えた」のです。
--group ship-svc --from-beginning --max-messages 3 --timeout-ms 8000 | ship-process。コンソールコンシューマーは、3件を渡した後、オフセット3をコミットして終了しますが、ハンドラーは2件目で落ちました。次にこのグループで読むと、3から始まります。
Kafkaにはそのまま残っている
新しいグループ(audit-svc)で、最初から5件を読んで、/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。コンシューマーの位置が整数1つだからできることです。設計ドキュメントは、これをキューの契約には反するが、どうしても必要な機能と呼んでいます。
再び処理すると2回届く
SHIP_FIXED=1で、ship-svcが5件を再び読んで、ship-processにパイプしてください。processed.txtは6行になり、order-1が2回なければなりません。巻き戻しの代償である、重複です。
... --group ship-svc --max-messages 5 --timeout-ms 8000 | SHIP_FIXED=1 ship-process。グループにオフセットがあれば、--from-beginningは無視されます。巻き戻した位置(0)から読みます。消えていた2件は戻ってきますが、すでに処理したorder-1も再び届きます。
冪等なコンシューマー
/root/kafka/dedup.shファイルを作成してください。標準入力の注文を読んで、/root/kafka/seen.txtファイルにないものだけを、/root/kafka/processed-dedup.txtファイルに書き込み、seenに記録します(環境変数DEDUP_SEEN・DEDUP_OUTがあれば、そのパス)。shipmentsを最初から2回読んでパイプしても、5行だけが残らなければなりません。
grep -qxF "$line" "$SEEN"で見たことがあるかを確認し、ないときだけ2つのファイルに書き込みます。グループなしで--from-beginning --max-messages 5で2回読んでパイプしてください。採点ツールは、一時的なパスを環境変数で渡して、重複が混ざった入力を入れてみます。
新しいグループはどこから始めるのか
--from-beginningなしで、新しいグループ(late-svc)で5秒読み(0件)、auto.offset.reset=earliestで、新しいグループ(early-svc)で5件を読んでください。/root/kafka/offset-reset.txtファイルに、late_count=0、early_count=5の2行を書いてください。
--group late-svc --timeout-ms 5000は、何も読めずに終わります(デフォルトはlatest)。--group early-svc --command-property auto.offset.reset=earliest --max-messages 5 --timeout-ms 8000は、5件を読みます。wc -lで数えて、ファイルに書いてください。
消えたものと、2回届いたものを数える
/root/kafka/consumer-report.txtファイルに、lost_after_crash=<2단계에서 사라진 건수>、duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>、unique_orders=<processed-dedup.txt 의 줄 수>の3行を書いてください(プレースホルダーは、ステップ2で消えた件数、ステップ5の後のprocessed.txtの重複件数、processed-dedup.txtの行数です)。
消えた件数は、ship-crash.txtのCURRENT-OFFSETから、そのとき処理された件数(1)を引いたもので、重複の件数は、processed.txtの行数から、異なる注文の数を引いたものです。採点ツールは、同じファイルから数え直します。