一件が消える — コミットのタイミングとat-most-once・at-least-once
一言でいうと
コンシューマーが位置を保存するタイミングが、セマンティクスを決めます。処理の前に保存すると、処理中に
落ちたメッセージは再び届かず(at-most-once、「消失」)、処理の後に保存すると、
保存の前に落ちたメッセージが再び届きます(at-least-once、「2回届く」)。後者を選んで、
ハンドラーを冪等にするのが答えで、巻き戻し(--reset-offsets)は、消えたものを
取り戻すための道具です。
なぜ必要なのか
設計ドキュメントの「メッセージ配信のセマンティクス」 の節は、コンシューマー側を2つの段落で終えます。コンシューマーがメッセージを読んで位置を保存してから 処理すると、保存の後で処理の前に落ちたとき、引き継いだプロセスは保存された位置から 始めるので、その前のメッセージは処理されません。at-most-onceです。読んで処理してから 保存すると、処理の後で保存の前に落ちたとき、引き継いだプロセスがすでに処理したメッセージを 再び受け取ります。at-least-onceです。そして、ドキュメントはこう付け加えます。多くの場合、メッセージに 主キーがあって、更新が冪等になります(同じメッセージを2回受け取っても、同じレコードを上書き するだけです)。
この2つの段落が、コースのタイトルの2つの事故です。「消えた1件」は、処理の前に保存した結果で、 それを直すと「2回届く」になります。どちらか一方を選ぶのではなく、2回届くほうを選んで、2回処理 しても問題ないようにするのが設計です。
どう動くのか
自動コミットは、処理と無関係に動きます: コンシューマー設定のドキュメント
のenable.auto.commitは、デフォルトがtrueで、auto.commit.interval.ms(デフォルト5000)
ごとに、オフセットがバックグラウンドでコミットされます。コンソールコンシューマーのように、自動コミットに頼る
ハンドラーは、「読んだもの」をコミットするのであって、「処理したもの」をコミットするのではありません。コンシューマーが3
件を渡してオフセット3をコミットしたのに、ハンドラーが2件目で落ちたなら、次のコンシューマーは
3から読みます。2件目と3件目は、Kafkaにそのまま残っていますが、このグループは永遠に
通り過ぎてしまいます。
shipments P0: order-1 order-2 order-3 order-4 order-5
ship-svc 읽음 3건 → 커밋 3 → 처리기 order-2 에서 크래시
처리됨: order-1 사라짐: order-2, order-3
消えたのはメッセージではなく、位置です: 消費は削除ではないので、別のグループで
最初から読めば、5件すべてあります。そして、運用ドキュメント
の--reset-offsetsで、グループの位置を巻き戻せます。--to-earliest、
--to-latest、--to-offset、--shift-by、--to-datetimeなどのシナリオがあり、
--executeを付けて初めて実際に変わります(なければ計画だけを表示します)。また、コンシューマー
インスタンスが停止していなければなりません。巻き戻すと、消えていた2件が戻ってきますが、すでに
処理したorder-1も再び届きます。巻き戻しの代償が、重複です。
冪等なコンシューマー: ドキュメントが述べた「主キーがあって更新が冪等になる場合」を、ハンドラー 側で作るのです。処理した注文番号を、結果と同じ場所に残し、入ってきた 注文がすでにあれば、スキップします。状態をメモリだけに置くと、プロセスが落ちるときに一緒に 消えて、次の実行がまた重複を出します。ファイルでもDBでも、処理結果と一緒に生き残る必要が あります。
新しいグループはどこから始めるのか: auto.offset.resetの項目は、グループにコミットされた
オフセットがないか、そのオフセットがもう存在しないとき(データが削除された場合)に、
何をするかです。earliestは一番前へ、latest(デフォルト)は一番後ろへ、
by_duration:<ISO8601>は今からその期間だけ前へ、noneは例外を投げます。
デフォルトがlatestなので、新しいグループを作って何も設定せずに接続すると、以前のメッセージは
すべてスキップされます。「新しいサービスが古い注文を見られなかった」という報告の原因です。コンソール
ツールの--from-beginningは、これをearliestに切り替えてくれるショートカットで、グループに
すでにオフセットがあれば無視されます。
コミットは消えます: ブローカー設定のドキュメント
のoffsets.retention.minutes(デフォルト10080、7日)によると、グループが空であるか、トピックの
購読を止めたままこの時間が過ぎると、コミットされたオフセットが破棄されます。その後コンシューマーが
戻ってくると、オフセットがないのでauto.offset.resetが適用されます。1週間止まっていた
バッチのコンシューマーが戻ってきて、latestで始めると、その間のデータをすべてスキップします。
現場での姿
「注文が消えました」という報告を受けたら、手順はこうです。グループをdescribeしてコミットされた
オフセットを見て、別のグループ(またはグループなし)でその区間を読み、メッセージがあるかを
確認し、あれば位置の問題です。コンシューマーが停止した状態で、--reset-offsets --to-offsetで巻き戻して再処理します。再処理が重複を生まないかどうかは、ハンドラーが
冪等かどうかにかかっていて、そうでなければ、巻き戻しの前にそれを直さなければなりません。
逆に、「同じ注文が2回処理されました」という報告は、たいていリバランスや再起動の後の at-least-onceの正常な動作です。それをバグと見てコミットを早めると、次の報告は 「消えました」になります。2つの報告は、同じつまみの両端です。
次のラボですること
配送ハンドラーが2件目の注文で落ちて、2件が消えることを再現し、別の
グループで、それがKafkaに残っていることを確認した後、グループを巻き戻して再処理し、
重複が生じることを確認し、冪等なハンドラーで防ぎます。最後に、auto.offset.reset
の2つの値を、新しいグループで比較します。