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

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

一件が消える — コミットのタイミングとat-most-once・at-least-once

TT Labで続きを見る

一言でいうと

コンシューマーが位置を保存するタイミングが、セマンティクスを決めます。処理の前に保存すると、処理中に 落ちたメッセージは再び届かず(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つの値を、新しいグループで比較します。