「厳密に一度」は嘘だ
一言でいうと
「正確に1回」は、配信層が作ってくれるものではなく、少なくとも1回の配信+受信側の重複排除で作るものであり、後半の半分はこちらの役割です。
なぜキューを使うのか
同期連携は、相手が生きていて初めて成立します。非同期は、その前提をなくします。
- 時間の分離: 相手が今いなくても、あとで処理されます
- 負荷の吸収: 毎秒1万件が集中しても、キューに溜めておき、毎秒1000件ずつ処理します
- 障害の分離: 受信側の障害が、送信側に広がりません
- 再処理: 失敗したメッセージを再投入できます
その代わり、支払う代償があります。複雑さと重複です。「応答をすぐに受け取れない」ということは、結果の通知チャネルを別に設計する必要があるということで、「再試行する」ということは、同じメッセージが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種類に分けます。
- 永続エラー(形式エラー、必須値の欠落、存在しないコード) → ただちにDLQ
- 一時エラー(DBコネクションの失敗、相手システムの5xx、タイムアウト) → 再試行
この区別がないとキューが詰まり、詰まったキューは、そのまま全業務の停止です。
キューの滞留(lag)は最も重要なメトリクス
非同期連携の健全性を1つの数字で見るなら、コンシューマーラグ(consumer lag)です。「溜まったメッセージ数」または「最も古い未処理メッセージの経過時間」です。
- 普段は0付近 → 正常
- 増え続ける → コンシューマーが追いつけていない(性能の問題、またはコンシューマーのダウン)
- 突然急増 → 生産側の急増、またはコンシューマーの障害
- 減らない一定の値 → ポイズンメッセージで詰まっている可能性
アラートは、絶対値ではなくトレンドと持続時間で設定します。「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つあります。
- 別のファイルシステムの間の
mvは、コピー+削除なので、アトミックではありません。同じマウントの中で動かす必要があります。 - 送信側がファイルを書き終える前に、コンシューマーが持っていってしまうことがあります。そのため、一時的な名前で書き、完了後に名前を変更するか、完了フラグファイル(
.ok)も一緒に作る規約を使います。
再処理の手順は事前に作っておく
障害が起きたら、必ず再処理をすることになります。そのときに必要なもの。
- 何が失敗したか: error/のメッセージと理由
- なぜ失敗したか: 理由別の分類(形式/業務/システム)
- 直せるか: データを修正して再投入するか、ソースに再送を依頼するか
- 再処理しても安全か: 冪等性が保証されているか
- 誰が承認するか: 金額が動くインターフェースは、承認が必要なことがある
この手順を、サービスイン前に文書化してリハーサルしておく必要があります。障害当日に作ると、その夜は徹夜になります。
現場での姿
キューを導入したプロジェクトで実際に問題になるのは、キュー自体ではなく、コンシューマー側の前提です。
最もよくあるのは重複です。ネットワークが一度切れて復旧すると、ブローカーは、まだ確認応答を受け取っていないメッセージを再送します。これは故障ではなく、規格どおりに動作したのですが、コンシューマーがそれを知らなければ、同じ注文が2件ロードされます。そして、この事故は、たいてい月末の精算で金額が合わなくて発見されます。事故が起きてから3週間後です。
2つ目は順序です。キュー全体が順序を守ってくれると信じて作ると、パーティションやコンシューマーを増やした瞬間に崩れます。順序が保証される範囲は、たいてい1つのパーティションの内側なので、同じ注文のメッセージが同じパーティションに行くように、キーを設定する必要があります。
3つ目は滞留です。キューは負荷を吸収してくれるので、コンシューマーが遅くても、送信側は何の異常も感じません。そのため、滞留(lag)を監視していないと、数時間遅れていることを、「今日のデータが見えないのですが」という問い合わせで知ることになります。