プログラミングとソフトウェア開発

Goを使用してセッションの順序を維持するKafka処理パイプラインの構築方法

Joshua Oluikpeが、各セッション内でメッセージを厳密な順序で処理しつつ、独立したセッションを並行して処理できるKafkaの実践的な設計を解説する。この設計は、コンシステントハッシュ、セッションごとのGoワーカー、同一パス内での再試行、障害後のメッセージ損失を防ぐ隣接コミットポイントの管理に基づいている。

2026-10-07
1 分で読めます
1 閲覧数
certi.news Editorial Team
Goを使用してセッションの順序を維持するKafka処理パイプラインの構築方法

数千のセッションが同じパーティションを共有する場合、パーティション内のKafkaの順序だけに依存することは十分ではない。チャットシステムや連続的なタスクでは、「ロンドンへのフライトを予約して」というメッセージの後に「目的地をパリにして」という訂正メッセージが到着することがあるが、後者を前者より先に実行すると文脈が壊れる。InfoQでJoshua Oluikpeが説明している解決策は、順序付けの責任の一部をGoで記述されたアプリケーション層に移す。

セッション内の順序とセッション間の並列性

この設計は2層のワーカーに基づいている。分配層はKafkaレコードを受け取り、コンシステントハッシュを使用してセッションIDに基づいて振り分ける。これにより、同じセッションのメッセージは同じワーカーに届く一方、異なるセッションは複数のワーカーに分散できる。その後、アクティブな各セッションに専用のgoroutineが割り当てられ、到着順に一度に1件のメッセージを処理する。

これにより、フラットなワーカープールで1つのセッションの再試行がワーカー全体を占有してしまう場合とは異なり、処理の遅いセッションが他のセッションの妨げになることはない。goroutineは必要に応じて作成され、一定期間非アクティブになると削除されるため、リソース消費は想定されるセッション総数ではなく、アクティブなセッション数に比例する。

メッセージを追い越さない再試行

再試行は同じセッションのgoroutine内で行われる。一時的なエラー、例えば依存サービスのタイムアウトが発生した場合、ワーカーは一定のランダム性を加えた指数バックオフを使用して待機し、その後続のメッセージへ進む前に同じメッセージを再処理する。後続のメッセージは失敗したメッセージの後ろでチャネルに留まるため、4番目のメッセージが3番目のメッセージを追い越すことはない。

一方、不正なデータのような修復不能なエラーは、直接デッドレターキュー(DLQ)へ送られる。一時的なエラーと最終的なエラーをこのように分けることで、再試行を管理する別個の状態機械は不要になる。ただし、バックオフ中に特定のセッションのメッセージが蓄積する可能性があることも意味する。

安全なコミットと障害後の復旧

複数のセッションがメッセージを並行して処理する場合、完了した中で最も高いオフセットだけをコミットしてはならない。102より先に104の処理が完了し、104をコミットすると、コンシューマーの障害によって102が恒久的にスキップされる可能性がある。そのためシステムは、各パーティションについて処理中のオフセットと完了したオフセットを追跡し、隣接して連続的に完了した最高のオフセットまでしかコミットポイントを進めない。

その結果、再起動はより古いオフセットから開始され、一部のメッセージが再処理される可能性がある。設計では重複排除を支援するため固定イベントIDを使用するが、「完全な1回限り」の実行を実現すると主張してはいない。外部への副作用は、依存サービスが重複を冪等に処理できるかどうかに依然として左右される。まれなケースでは、パーティションの進行を回復するため、所定のタイムアウト後に停滞したギャップをスキップし、運用上の障害として記録したうえでメッセージをDLQに送ることもできる。これは継続性と処理の完全性の間にある明示的なトレードオフである。

実際の本番運用に必要なもの

順序の保証は基本ロジックだけでは完成しない。パーティションの再割り当て中、コンシューマーは新しいレコードの受信を停止し、実行中の処理が完了するのを待ったうえで、所有権を移転する前に最後の隣接watermarkをコミットする。ワーカーのチャネルが満杯になった場合は、メッセージを破棄するのではなく、取得を一時的に停止する。また、すでに取得したレコードはパーティションごとの専用バッファに保存し、より新しいメッセージがそれらを追い越さないようにする。

さらにシステムには、停滞したオフセットを個別に検出する機能、非同期なKafka境界をまたぐトレーシングの関連付け、取得を停止している間の自動スケーリングに必要な遅延シグナルの維持、終了前の処理の排出とDLQへの書き込みが必要になる。製品側では、非同期なプロデュース時にもメッセージの順序を保証するため、送信側にもコンシステントハッシュを適用している。

結果と制約

50,000件のメッセージと10個のセッションを対象に、メッセージの10%へ一時的なエラーを注入した合成テストでは、システムは順序違反なしに毎秒14,027件のメッセージを処理し、中央値のレイテンシーは11.7ミリ秒だった。1,000セッションでの負荷テストでは、送信エラーなしに、目標の50,000件に対して実効スループットは毎秒48,805件に達した。ただし、テールレイテンシーは200ミリ秒の予算を超えた。これはセッション順序付けの仕組みではなく、API Gatewayにおけるメッセージレートの割り当てが原因だった。

著者によれば、このシステムは本番環境で4,000万件を超えるメッセージを処理し、目立った順序違反はなく、DLQへ送られたのはわずか460件、すなわち約0.002%だった。ただし、この数値は運用監視に基づくものであり、すべてのメッセージの順序を包括的に検査したものではない。また、より高い負荷テストでは、各セッションの順序が独立に検証されたわけでもない。そのため、この設計の実務的な価値は、順序をKafkaレベルの前提からアプリケーションが強制する特性へと変換することにある。同時に、明示的なシーケンスのテストと、依存サービスにおける冪等性の保証は依然として必要である。

ニュースの出典
InfoQ - Architecture Articles
原文を開く ↗
c
著者

certi.news Editorial Team

同じカテゴリー

おすすめ記事

すべてのニュースを見る