TT Lab
はじめる
学ぶ 学習パス コース

注文が二度届き、一度は消えた

コミットのタイミングで一件が消える — 巻き戻しと冪等処理

TT Labで続きを見る

このラボは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回届く」になります。そのため、ハンドラーは冪等でなければなりません。 ドキュメントは、それを「メッセージに主キーがあって更新が冪等になる場合」と書いています。

ステップ

  1. shipmentsトピックを、パーティション1つで作成し、order-1からorder-5までの5行を入れてください。
  2. グループ(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が「消えた」のです。
  3. 新しいグループ(audit-svc)で、最初から5件を読んで、/root/kafka/ship-audit.txtファイルに保存してください。Kafkaには、すべて残っています。
  4. ship-svcのオフセットを一番前まで巻き戻して(--reset-offsets --to-earliest --execute)、その出力を、/root/kafka/ship-reset.txtファイルに保存してください。
  5. SHIP_FIXED=1で、ship-svcが5件を再び読んで、ship-processにパイプしてください。processed.txtは6行になり、order-1が2回なければなりません。巻き戻しの代償である、重複です。
  6. /root/kafka/dedup.shファイルを作成してください。標準入力の注文を読んで、/root/kafka/seen.txtファイルにないものだけを、/root/kafka/processed-dedup.txtファイルに書き込み、seenに記録します(環境変数DEDUP_SEEN・DEDUP_OUTがあれば、そのパス)。shipmentsを最初から2回読んでパイプしても、5行だけが残らなければなりません。
  7. --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行を書いてください。
  8. /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の行数です)。

参考

配送注文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の行数から、異なる注文の数を引いたものです。採点ツールは、同じファイルから数え直します。