再試行が二度目の到着を生む — 冪等プロデューサーで防ぐ
このラボは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回届いた」という報告を受けたときに、どこを見るべきかがわかります。
ステップ
paymentsトピックを、パーティション1つで作成し、デフォルトの設定で2行(pay-1、pay-2)を入れてください。kafka-dump-log.shでpayments-0のセグメントをダンプし、/root/kafka/dump-default.txtファイルに保存してください。バッチのproducerIdが-1ではなく、baseSequenceが0でなければなりません(デフォルトのプロデューサーは冪等です)。enable.idempotence=falseで1行(pay-3)をさらに入れて、もう一度ダンプし、/root/kafka/dump-noidem.txtファイルに保存してください。新しいバッチはproducerId: -1でなければなりません。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つ以上残っていなければなりません。- 同じ条件で、
payments-idemトピックに、冪等をオンにして(enable.idempotence=true)入れ、stderrを、/root/kafka/idem-producer.logファイルに残してください。リトライは起きますが、ログには1つだけ残っていなければなりません。 - slowをオフにした状態で、
payments-appトピックにpay-77を、プロデューサーを2回別々に実行して入れてください(冪等はオンのまま)。ダンプを、/root/kafka/two-sessions.txtファイルに保存してください。2つのバッチのproducerIdが、互いに異なっていなければなりません。 /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行を書いてください(プレースホルダーは、それぞれのトピックのレコード数です)。
参考
- ダンプ:
kafka-dump-log.sh --files /var/lib/kafka/<토픽>-0/00000000000000000000.log --print-data-log(プレースホルダーはトピック名です)。バッチの行にproducerId・baseSequenceが、レコードの行にpayloadが表示されます。 - レコード数:
kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic <토픽>の最後の数字(ログの末尾のオフセット)です(プレースホルダーはトピック名です)。 - クライアントの設定は、
--command-property 키=값で指定します(4.3では--producer-propertyは非推奨。プレースホルダーはキーと値です)。5,000文字の1行は、head -c 5000 /dev/zero | tr '\0' x > big.txt; echo >> big.txtで作ります。 - プロデューサー設定のドキュメント:
delivery.timeout.msはrequest.timeout.ms + linger.ms以上でなければならず、冪等をオンにするには、acks=all、retries>0、max.in.flight.requests.per.connection<=5でなければなりません。冪等を明示的にオンにしたまま、競合する値を指定すると、ConfigExceptionです。 - よくあるミス1: slowをオンにしたままステップ4を終えて、オフにしないことです。小さなリクエストは影響がなく気づきにくいですが、大きなレコードを扱う次のステップが遅くなります。
- よくあるミス2: ステップ6で、1つのプロデューサーに2行を入れることです。それは同じセッションなので、シーケンス番号が続き、1つのバッチになります。「2回別々に実行」して初めて、セッションが2つになります。
デフォルトのプロデューサーで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つの数字を今のブローカーで数え直して比較します。