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

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

再試行が二度目の到着を生む — 冪等プロデューサーで防ぐ

TT Labで続きを見る

このラボはVM上で動きます

UbuntuのVMに、Apache Kafka 4.3.1がKRaftの単一ノードとして起動しています (localhost:9092)。ヘルパーkafka-lab-slow on|offは、ブローカーへ向かう4KB以上の リクエストだけを600ms遅らせて、プロデューサーのタイムアウトとリトライを引き起こします。最初の起動には 4分ほどかかります。

目標

プロデューサーがレスポンスを受け取れずにリトライするとき、ログに同じレコードが複数残ることを 再現し、冪等プロデューサーがそれを1つにすることを、kafka-dump-log.shで 確認します。そして、冪等が対象としない場面(プロデューサーのセッションをまたぐ再送)も 確認します。

なぜ重要なのか

設計ドキュメントが述べるように、プロデューサーはネットワークエラーに遭遇すると、そのエラーがメッセージが コミットされる前に起きたのか後に起きたのかがわかりません。そのため再送し、 そうすると、元のリクエストが成功していた場合、ログに2回残ります。at-least-onceです。 冪等プロデューサーは、ブローカーがプロデューサーごとにIDを渡し、レコードごとにシーケンス番号を付けて、 同じシーケンス番号をふるい落とす方式で、これを防ぎます。4.xではデフォルトでオンに なっていますが(enable.idempotenceのデフォルト値はtrue)、設定を1つ間違えて触ると静かに オフになり、セッションをまたぐ再送は、そもそも対象外です。この3つを自分の目で確認しておけば、 「2回届いた」という報告を受けたときに、どこを見るべきかがわかります。

ステップ

  1. paymentsトピックを、パーティション1つで作成し、デフォルトの設定で2行(pay-1、pay-2)を入れてください。
  2. kafka-dump-log.shでpayments-0のセグメントをダンプし、/root/kafka/dump-default.txtファイルに保存してください。バッチのproducerIdが-1ではなく、baseSequenceが0でなければなりません(デフォルトのプロデューサーは冪等です)。
  3. enable.idempotence=falseで1行(pay-3)をさらに入れて、もう一度ダンプし、/root/kafka/dump-noidem.txtファイルに保存してください。新しいバッチはproducerId: -1でなければなりません。
  4. kafka-lab-slow onをオンにして、payments-dupトピックを作成した後、5,000文字の1行(/root/kafka/big.txt)を、冪等をオフにして、request.timeout.ms=300、delivery.timeout.ms=2000、retries=2、max.block.ms=10000で入れてください。プロデューサーのstderrを、/root/kafka/dup-producer.logファイルに残し、終わったらkafka-lab-slow offを実行してください。ログには、同じレコードが2つ以上残っていなければなりません。
  5. 同じ条件で、payments-idemトピックに、冪等をオンにして(enable.idempotence=true)入れ、stderrを、/root/kafka/idem-producer.logファイルに残してください。リトライは起きますが、ログには1つだけ残っていなければなりません。
  6. slowをオフにした状態で、payments-appトピックにpay-77を、プロデューサーを2回別々に実行して入れてください(冪等はオンのまま)。ダンプを、/root/kafka/two-sessions.txtファイルに保存してください。2つのバッチのproducerIdが、互いに異なっていなければなりません。
  7. /root/kafka/producer-report.txtファイルに、dup_copies=<payments-dup 의 레코드 수>、idem_copies=<payments-idem 의 레코드 수>、two_sessions_copies=<payments-app 의 레코드 수>、acks_default=all、idempotence_default=trueの5行を書いてください(プレースホルダーは、それぞれのトピックのレコード数です)。

参考

デフォルトのプロデューサーで2件

paymentsトピックを、パーティション1つで作成し、デフォルトの設定で2行(pay-1、pay-2)を入れてください。

--create --topic payments --partitions 1の後に、printf 'pay-1\npay-2\n' | kafka-console-producer.sh ...。レコード数は、kafka-get-offsets.shで確認します。

デフォルトのプロデューサーは冪等

kafka-dump-log.shでpayments-0のセグメントをダンプし、/root/kafka/dump-default.txtファイルに保存してください。バッチのproducerIdが-1ではなく、baseSequenceが0でなければなりません(デフォルトのプロデューサーは冪等です)。

kafka-dump-log.sh --files /var/lib/kafka/payments-0/00000000000000000000.log --print-data-log。ブローカーがプロデューサーに渡したIDと、レコードのシーケンス番号が、バッチヘッダーにそのまま書かれています。これが、重複をふるい落とす材料です。

冪等をオフにするとIDがない

enable.idempotence=falseで1行(pay-3)をさらに入れて、もう一度ダンプし、/root/kafka/dump-noidem.txtファイルに保存してください。新しいバッチはproducerId: -1でなければなりません。

--command-property enable.idempotence=false。冪等がオフのプロデューサーはIDを受け取らないので、ブローカーは、同じレコードが再び届いても、見分ける術がありません。

リトライが2度目の到着を生む

kafka-lab-slow onをオンにして、payments-dupトピックを作成した後、5,000文字の1行(/root/kafka/big.txt)を、冪等をオフにして、request.timeout.ms=300、delivery.timeout.ms=2000、retries=2、max.block.ms=10000で入れてください。プロデューサーのstderrを、/root/kafka/dup-producer.logファイルに残し、終わったらkafka-lab-slow offを実行してください。ログには、同じレコードが2つ以上残っていなければなりません。

リクエストが300ms以内に答えを受け取れないと、プロデューサーは同じバッチを再送し(retries)、ブローカーは、遅れて届いた元のリクエストも、遅れて届いたリトライも、すべて書き込みます。プロデューサーは最終的に失敗と報告しますが、ログには3つあります。kafka-get-offsets.shで数えてみてください。

冪等プロデューサーは1つだけ残す

同じ条件で、payments-idemトピックに、冪等をオンにして(enable.idempotence=true)入れ、stderrを、/root/kafka/idem-producer.logファイルに残してください。リトライは起きますが、ログには1つだけ残っていなければなりません。

ステップ4とまったく同じようにslowをオンにして送りますが、enable.idempotence=trueだけを変えます。リトライされたバッチは、同じシーケンス番号を付けて届くので、ブローカーがふるい落とします。プロデューサーのログには、依然としてREQUEST_TIMED_OUTが出力されます。リトライはあったのに、重複だけがないのです。

セッションが違うと冪等が対象にしない

slowをオフにした状態で、payments-appトピックにpay-77を、プロデューサーを2回別々に実行して入れてください(冪等はオンのまま)。ダンプを、/root/kafka/two-sessions.txtファイルに保存してください。2つのバッチのproducerIdが、互いに異なっていなければなりません。

printf 'pay-77\n' | kafka-console-producer.sh ...を2回。プロセスが違うと、ブローカーが新しいプロデューサーIDを渡し、シーケンス番号も0からなので、ブローカーから見れば別のレコードです。アプリケーションが失敗の後に再送するのは、この形で、それを防ぐのは、冪等プロデューサーではなく、注文番号のような業務キーです。

3つのトピックのレコード数でまとめる

/root/kafka/producer-report.txtファイルに、dup_copies=<payments-dup 의 레코드 수>、idem_copies=<payments-idem 의 레코드 수>、two_sessions_copies=<payments-app 의 레코드 수>、acks_default=all、idempotence_default=trueの5行を書いてください(プレースホルダーは、それぞれのトピックのレコード数です)。

3つの数字は、kafka-get-offsets.shの最後のフィールドです。残りの2つは、プロデューサー設定のドキュメントのデフォルト値です。採点ツールは、3つの数字を今のブローカーで数え直して比較します。