二度目の到着 — 再試行による重複と冪等プロデューサー
一言でいうと
プロデューサーがレスポンスを受け取れないと、メッセージがコミットされたかどうかを知る術がないので再送し、 元のリクエストが成功していたなら、ログに2つ残ります。冪等プロデューサーは、ブローカーが渡した プロデューサーIDとレコードのシーケンス番号で再送をふるい落とし、1つだけ残します。4.xでは デフォルトでオンになっていますが、セッションをまたぐ再送は対象外です。
なぜ必要なのか
設計ドキュメントの「メッセージ配信のセマンティクス」 の節は、問題を正確に書いています。プロデューサーが発行中にネットワークエラーに遭遇すると、そのエラーが メッセージがコミットされる前に起きたのか後に起きたのかがわかりません。 自動生成のキーを持つテーブルにINSERTしている最中に、接続が切れたのと同じです。0.11より前の プロデューサーは、再送する以外に選択肢がなく、そのため、at-least-onceでした。元の リクエストが実は成功していたなら、再送がログに同じメッセージをもう1回書き込みます。
0.11から、冪等な配信のオプションが加わりました。ブローカーがプロデューサーごとにIDを渡し、プロデューサーは メッセージごとにシーケンス番号(sequence)を付けて送り、ブローカーは同じIDとシーケンス番号のメッセージを ふるい落とします。同じ時期にトランザクションも導入され、複数のパーティションに原子的に書き込めるように なりました。このコースは、冪等までを扱います。トランザクションは、その上の層です。
どう動くのか
acksは何を待つかを決めます: プロデューサー設定のドキュメント
のacksの項目によると、0はサーバーの確認をまったく待たないので、受け取ったという保証がなく、
リトライも起きません(オフセットは常に-1)。1はリーダーが自分のログに書き込むと
応答するので、フォロワーが複製する前にリーダーが落ちると失います。all(= -1)はISR
全体が確認するまで待ち、ISRのうち1つでも生きていれば失わない、最も
強い保証です。デフォルト値はallです。冪等をオンにするには、allでなければなりません。
リトライは、デフォルトではほぼ無限です: retriesのデフォルト値は2147483647で、ドキュメントは、
この値には手を触れず、delivery.timeout.ms(デフォルト120000)でリトライの合計時間を
管理するよう勧めています。delivery.timeout.msは、request.timeout.ms(デフォルト30000) +
linger.ms以上でなければなりません。request.timeout.msの項目には、こんな一文があります。
この値は、ブローカーのreplica.lag.time.max.ms(デフォルト30000)より大きくすると、不要な
リトライによるメッセージ重複の可能性を減らせます。ドキュメントが重複をリトライの
結果として明示している箇所です。
冪等には条件があります: enable.idempotenceの項目によると、オンにすれば、各メッセージがちょうど
1つだけストリームに書き込まれ、オフにすると、ブローカーの障害などによるリトライが、重複を書き込む
ことがあります。オンにするには、max.in.flight.requests.per.connectionが5以下、retriesが0
より大きく、acksがallでなければなりません。競合する設定があり、冪等を明示的に
オンにしていなければ、冪等は静かにオフになります。明示的にオンにしたまま競合すると、
ConfigExceptionです。そのため、acks=1を「性能のために」入れた瞬間に、重複の防御が
なくなるのに、何の警告もありません。
멱등 프로듀서의 배치 헤더 (kafka-dump-log.sh)
producerId: 1 producerEpoch: 0 baseSequence: 0 lastSequence: 1 ← ID 와 순번
멱등을 끈 배치
producerId: -1 producerEpoch: -1 baseSequence: -1 lastSequence: -1 ← 걸러 낼 재료가 없다
このコードブロックの韓国語コメントは、上の2行が冪等プロデューサーのバッチヘッダーで、IDとシーケンス番号が入っていること、下の2行が冪等をオフにしたバッチで、重複をふるい落とす材料がないことを述べています。
ブローカーは、このヘッダーで再送を見分けます。実測(4.3.1、ローカルのブローカーに600msの遅延を
かけ、request.timeout.ms=300、retries=2)では、冪等をオフにしたプロデューサーは、同じ
レコードを3つ残し、冪等をオンにしたプロデューサーは1つだけ残しました。どちらも、プロデューサーは
最終的に「失敗」と報告しました。レスポンスを受け取れなかっただけで、ブローカーにはあったのです。
冪等が対象としないもの: IDとシーケンス番号は、プロデューサーのセッションのものです。プロセスが
新しく起動すると、新しいIDを受け取り、シーケンス番号は0からになるので、アプリケーションが失敗を報告されて
再びsendするのは、ブローカーから見れば新しいレコードです。transactional.idの項目が、これを
「複数のプロデューサーセッションにまたがる信頼性」と呼び、トランザクションの領域にしています。
その層がないなら、答えは、業務キー(注文番号)でコンシューマー側でふるい落とすことです。
現場での姿
決済サービスで「同じ決済が2回記録された」という報告が来たら、最初に確認するのは、プロデューサーの
設定です。acks=1やmax.in.flight=10のような値があれば、冪等が静かにオフになって
います。2番目の確認は、ログのバッチヘッダーです。producerIdが異なる2つのバッチに
同じ決済があれば、アプリケーション層の再送で、これはプロデューサーの設定では
防げません。
逆に、「送ったと言うのにない」は、acks=0の姿です。プロデューサーは、ソケットバッファに入れた
瞬間に成功と報告し、その後に何が起きても知りません。遅延に敏感なログ収集なら
受け入れられる取引ですが、決済では受け入れられません。
次のラボですること
デフォルトのプロデューサーと、冪等をオフにしたプロデューサーのバッチヘッダーを、ダンプで比較し、遅い ネットワークをオンにして、リトライが3つを残すことと、冪等が1つにすることを再現した 後、プロデューサーを2回別々に実行して、セッションをまたぐ再送は防げないことを確認します。