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

システム間連携 (EAI)

「厳密に一度」は嘘だ

TT Labで続きを見る

一言でいうと

「正確に1回」は、配信層が作ってくれるものではなく、少なくとも1回の配信+受信側の重複排除で作るものであり、後半の半分はこちらの役割です。

なぜキューを使うのか

同期連携は、相手が生きていて初めて成立します。非同期は、その前提をなくします。

その代わり、支払う代償があります。複雑さと重複です。「応答をすぐに受け取れない」ということは、結果の通知チャネルを別に設計する必要があるということで、「再試行する」ということは、同じメッセージが2回処理されうるということです。

配信保証の3つの等級

等級 意味 現実
at-most-once 最大1回。消失の可能性あり ログ・メトリクスのように失ってもよいもの
at-least-once 最小1回。重複の可能性あり ほとんどの実務のメッセージング
exactly-once 正確に1回 条件付きでのみ成立

ここで必ず理解すべきこと。

「正確に1回」=「少なくとも1回の配信」+「受信側の重複排除」

つまり、exactly-onceは、配信層が魔法で作ってくれるものではなく、受信側が重複を取り除いて、結果的にそう見せているものです。ブローカーが「exactly-onceをサポート」と宣伝しても、それは特定の条件(同じクラスター、トランザクションAPIの使用)の中での話です。外部システムへの書き込みをした瞬間に崩れます。

したがって、コンシューマーは常に重複を前提に作る必要があります。例外ではなく、デフォルトです。

順序保証の範囲

もう1つ、よく誤解されること。

順序はパーティション(またはキュー)の中でだけ保証されます。トピック全体のグローバルな順序は保証されません。

そのため、「同じ注文番号のメッセージは順番どおりに処理されなければならない」という要件があれば、注文番号をパーティションキーに使う必要があります。そうすれば、同じ注文のメッセージは同じパーティションに入り、順序が維持されます。

これを知らずにラウンドロビンで分配すると、キャンセルのメッセージが注文のメッセージより先に処理されることが起きます。そして、それは1日に1件か2件しか発生しないので、原因を探すのが非常に難しいです。

コンシューマー設計の標準形

1. 메시지 수신
2. 파싱 및 형식 검증        → 실패: 즉시 error/DLQ (재시도해도 똑같다)
3. 멱등 확인 (이미 처리했나) → 이미 처리: 아무것도 안 하고 ack
4. 업무 처리 (DB 트랜잭션)
5. 처리 이력 기록           ← 4와 같은 트랜잭션 안에서
6. ack (큐에서 제거)

このコードブロックの韓国語は、コンシューマーの6つの手順を示しています。メッセージの受信、パースと形式の検証(失敗したら、ただちにerror/DLQへ。再試行しても同じ)、冪等の確認(すでに処理したか。処理済みなら何もせずack)、業務処理(DBトランザクション)、処理履歴の記録(項目4と同じトランザクションの中で)、ack(キューから削除)です。

項目5を項目4と同じトランザクションに入れることが核心です。別にすると、「業務は処理したのに、履歴は残せていない」状態が生まれ、その状態で再試行すると、重複処理になります。

そして、項目6を項目4より先に行ってはいけません。先にackして処理している途中で死ぬと、メッセージが消えます。これがat-most-onceになるポイントです。

パース失敗は再試行しない

項目2を別に置いた理由があります。JSONが壊れていたり、必須フィールドがなかったりするメッセージは、100回再試行しても100回失敗します。ところが、再試行ロジックに引っかかると、そのメッセージがキューの先頭で失敗し続け、後ろの正常なメッセージを塞ぎます。これをポイズンメッセージ(poison message)と呼びます。

そのため、エラーを2種類に分けます。

この区別がないとキューが詰まり、詰まったキューは、そのまま全業務の停止です。

キューの滞留(lag)は最も重要なメトリクス

非同期連携の健全性を1つの数字で見るなら、コンシューマーラグ(consumer lag)です。「溜まったメッセージ数」または「最も古い未処理メッセージの経過時間」です。

アラートは、絶対値ではなくトレンドと持続時間で設定します。「lag 1000超過」より、「lagが10分連続で増加」のほうがよい条件です。バッチ的な流入があるシステムでは、一時的にlagが大きいのが正常だからです。

ファイルキューもキューです

ブローカー(Kafka、RabbitMQ)がないSI現場も多いです。そのようなときは、ディレクトリベースのキューを使います。意外と堅牢です。

/data/if/inbox/       ← 도착
/data/if/processing/  ← 처리 중 (원자적 mv 로 이동 = 잠금)
/data/if/done/        ← 성공
/data/if/error/       ← 실패 (DLQ 역할)

このコードブロックの韓国語コメントは、順に、到着、処理中(アトミックなmvで移動すること=ロック)、成功、失敗(DLQの役割)を意味します。

核心は、mvが同じファイルシステムの中でアトミックであることです。inbox → processingの移動に成功したプロセスだけが、そのメッセージを持ちます。コンシューマーを複数起動しても、重複処理になりません。

注意点が2つあります。

  1. 別のファイルシステムの間のmvは、コピー+削除なので、アトミックではありません。同じマウントの中で動かす必要があります。
  2. 送信側がファイルを書き終える前に、コンシューマーが持っていってしまうことがあります。そのため、一時的な名前で書き、完了後に名前を変更するか、完了フラグファイル(.ok)も一緒に作る規約を使います。

再処理の手順は事前に作っておく

障害が起きたら、必ず再処理をすることになります。そのときに必要なもの。

この手順を、サービスイン前に文書化してリハーサルしておく必要があります。障害当日に作ると、その夜は徹夜になります。

現場での姿

キューを導入したプロジェクトで実際に問題になるのは、キュー自体ではなく、コンシューマー側の前提です。

最もよくあるのは重複です。ネットワークが一度切れて復旧すると、ブローカーは、まだ確認応答を受け取っていないメッセージを再送します。これは故障ではなく、規格どおりに動作したのですが、コンシューマーがそれを知らなければ、同じ注文が2件ロードされます。そして、この事故は、たいてい月末の精算で金額が合わなくて発見されます。事故が起きてから3週間後です。

2つ目は順序です。キュー全体が順序を守ってくれると信じて作ると、パーティションやコンシューマーを増やした瞬間に崩れます。順序が保証される範囲は、たいてい1つのパーティションの内側なので、同じ注文のメッセージが同じパーティションに行くように、キーを設定する必要があります。

3つ目は滞留です。キューは負荷を吸収してくれるので、コンシューマーが遅くても、送信側は何の異常も感じません。そのため、滞留(lag)を監視していないと、数時間遅れていることを、「今日のデータが見えないのですが」という問い合わせで知ることになります。