代表的な面接トピック

Kafka バックエンド面接:Share Group はどのような場合にキューセマンティクスを提供すべきか?

バックエンド難しい
Offer.cc 編集チーム公開日 更新日

質問

ある注文処理サービスにおいて、複数のコンシューマーが1つの Kafka トピックから独立したジョブを取得し、レコード単位の確認応答と再試行を行いたいと考えています。Share Group と Consumer Group を比較し、それぞれがどのような場合に安全であるかを説明した上で、順序性、重複、移行への対処方法を述べてください。

問題とスコープ

注文通知、画像変換、請求計算などは、多くの場合、独立した作業項目です。チームはすでにイベントを Kafka に書き込んでいますが、従来の Consumer Group ではパーティションが一度に1つのメンバーにのみ割り当てられるため、パーティション数を超えてワーカーを追加しても並列処理能力は直接向上しません。パーティションごとの順序性を依然として必要とするビジネスフローを特定しつつ、複数のコンシューマーがレコード単位の確認応答、再試行、監視可能な配信試行を伴って協調して処理できる仕組みを設計してください。

この設問では、Apache Kafka KIP-932 で説明されている Share Group モデルを使用します。この KIP は、通常のトピック上で協調的なコンシュームを行うための新しいグループタイプを定義しており、Kafka を RabbitMQ と完全に同じにするものではありません。優れた回答では、本番環境での利用を推奨する前に、デプロイされているブローカーのバージョン、クライアントのサポート状況、および API の可用性を検証します。

面接官が見ているポイント

  • 排他的なパーティション割り当てと、協調的なレコード取得の違いを対比できるか?
  • acknowledge、release、reject、およびロックの期限切れを処理状態に正しくマッピングできるか?
  • Share Group では、通常のキー順序の直感を維持することなく、パーティション数を超えるコンシューマーを持てることを理解しているか?
  • 再試行、ポイズンレコード、処理期限、並行性の上限を1つの障害モデルとして整理して回答できているか?
  • 単に KIP の名前を挙げるだけでなく、クライアント、ブローカー、ACL、監視、ロールバックの詳細まで検証しているか?

不十分な回答は「Kafka もキューになれる」とだけ述べるものです。優れた回答は、キューのようなメリット、変化する保証、そしてロールアウト前に必要なゲート(検証基準)を具体的に挙げます。

最初に確認すべき明確化のための質問

  1. ジョブは本当に独立していますか?1つの注文に対するイベントをキーの順序通りに適用する必要がある場合、Share Group は不適切なプリミティブである可能性があります。
  2. 障害は一時的なものですか、手動で復旧可能なものですか、それとも恒久的に無効なものですか?それによって release、reject、隔離(クランティン)の挙動が決まります。
  3. p99 の処理時間、最大並行数、許容される重複の副作用はどの程度ですか?これらが取得ロックおよび冪等性の設計を決定します。
  4. ビジネス要件としてエンドツーエンドの Kafka トランザクションが必要ですか?Share Group が既存の Consumer Group のトランザクション設計をそのまま引き継ぐと仮定してはいけません。
  5. トピックは現在もブロードキャストやリプレイを行うコンシューマーに利用されていますか?あるグループを変更しても、別のグループの読み取り規約を変更してはなりません。

30秒での回答

「まず、ジョブが順不同で完了してもよいか、またデプロイされている Kafka とクライアントが KIP-932 をサポートしているかを確認します。Share Group を使用すると、メンバーはトピックから協調してレコードを取得でき、メンバー数がパーティション数を超えることが可能になり、レコード単位の確認応答(acknowledge)、解放(release)、拒絶(reject)がサポートされます。一方、パーティション内の順序性やオフセットに基づく推論には、Consumer Group の方が適しています。私なら各副作用を冪等にし、処理レイテンシに基づいてロックと再試行のポリシーを設定し、acquire、acknowledge、release、reject、タイムアウトの状態を監視して、ポイズンレコードを隔離します。順序性、トランザクション境界、クライアントサポートが未解決の場合は、Consumer Group を維持し、移行前に小さな規模で別の作業用トピックを検証します。」

ステップごとの論理的思考

1. 2つの割り当てモデルを図解する

Consumer Group は通常、パーティションをメンバーに割り当てます。そのグループ内では1つのメンバーが特定のパーティションを読み取るため、並列度はパーティション数によって制限されます。Share Group では、メンバーがサブスクライブされたトピックからレコードを協調して取得できます。複数のメンバーが1つのパーティションから異なるレコードを処理でき、メンバー数がパーティション数を超えることも可能です。これは独立したジョブには有用ですが、グローバルな順序性を意味するものではありません。

text
Consumer Group:  partition-0 -> worker-A
                 partition-1 -> worker-B
                 extra workers wait for another partition

Share Group:     partition-0 records -> worker-A, worker-B, worker-C
                 each acquired record is locked for one consumer

Share Group を選択する理由は、単に「パーティションが少なすぎるから」ではなく、柔軟な作業取得とレコード単位の完了処理であるべきです。特定の顧客のイベントを順序通りに適用する必要がある場合は、Consumer Group を維持するか、アプリケーションレベルでシリアライズされたステートマシンを追加してください。

2. レコードのライフサイクルをモデル化する

KIP-932 では、時間制限付きの取得ロックについて説明されています。レコードを取得した後、コンシューマーは成功を確認(acknowledge)したり、別の配信のために解放(release)したり、処理不能として拒絶(reject)したり、ロックが切れるまで何もしないでおくことができます。KIP ではデフォルト値として30秒が記載されていますが、本番環境の動作はデプロイされたブローカーの設定に従う必要があります。デフォルト値は SLA ではありません。

text
available -> acquired -> acknowledged
                    -> released -> available
                    -> rejected  -> terminal or quarantine
                    -> lock timeout -> available

ハンドラーは、外部への副作用を実行する前に冪等性キーを登録する必要があります。そうしないと、クライアントのクラッシュやロック切れにより、二重請求、二重発送、または二重通知が発生する可能性があります。確認応答(acknowledge)はこの取得が完了したことを示すものであり、別のシステムによってすでにコミットされた副作用をロールバックすることはできません。

3. 再試行、ポイズンレコード、並行性を制限する

配信試行回数をカウントすることで、一時的な障害と恒久的に無効なレコードを区別しやすくなります。ネットワーク障害の場合はバックオフして release します。決定論的なスキーマエラーやバリデーションエラーの場合は、隔離トピックや手動キューに reject します。無制限に release し続けてはいけません。1つのポイズンレコードが無期限にロックと下流のキャパシティを消費し続ける可能性があります。

ロック時間は、説明可能なジッターマージンを含め、通常の p99 処理時間より長く設定します。短すぎると重複した再配信が発生し、長すぎると復旧が遅れます。また、パーティションごとに取得するレコード数を制限し、ワーカーのセマフォ、データベースプール、外部 API のクォータを調整します。アクティブなロック数、ロックタイムアウト、試行回数の分布、reject 数、エンドツーエンドの完了レイテンシを合わせて監視してください。

4. 順序性と重複に関する保証を再確認する

Kafka に関する説明では、「パーティション内で順序保証される」ことが「ビジネス処理が順序通りに行われる」と混同されがちです。Share Group のメンバーはレコードを並行して取得できるため、あるキーの完了順序は書き込み順序と異なる場合があり、release や再配信によってその差はさらに拡大します。順序が重要な場合は、キーレベルのシリアライズ、バージョンチェック、またはステートマシンをアプリケーションに実装してください。単に「Kafka は順序が保証されている」とだけ回答してはなりません。

また、Exactly-once(正確に1回)もグループタイプによって自動的に得られるわけではありません。レコードの取得、ビジネスの書き込み、確認応答の間の境界を追跡してください。外部データベースや決済サービスには、依然として冪等性キー、重複排除の制約、またはトランザクショナルアウトボックスが必要です。組み合わせがサポートされていない場合は、Exactly-once と呼ぶのではなく、冪等性を備えた At-least-once(少なくとも1回)の設計であると明記してください。

5. 移行とロールバックを計画する

ブローカーのバージョン、クライアント API、グループ設定、ACL、メトリクス、運用コマンドを検証します。その後、別のトピックまたは小さなワークロードで負荷テストを実施します。取得後のクラッシュ、ロック時間を超える処理、繰り返される reject、ブローカーの再起動、コーディネーターの移動など、障害を注入してテストします。すべてのレコードについて、ビジネスキー、試行回数、状態、タイムスタンプを記録します。

古いコンシューマーが順序性やトランザクションに依存している場合は、同じグループをインプレースで変更してはいけません。作業を専用のトピックにコピーし、新しいグループに段階的にトラフィックを移行させます。エラー率、重複による副作用、レイテンシが基準を満たすまで、古いパスをリプレイ可能な状態に保ちます。ロールバックでは新規の取得を停止し、未移行のレコードを古いパスで消費させます。明示的な重複排除境界なしに、2つのアクティブなパスが同じ副作用を実行してはなりません。

質の高い模範解答

「まず、ジョブが順不同で完了してもよいか、処理が冪等であるか、そしてデプロイされているブローカーとクライアントが KIP-932 をサポートしているかを確認します。Share Group はトピック内の独立したレコードを協調的な作業として扱います。複数のメンバーが1つのパーティションから異なるレコードを取得でき、メンバー数がパーティション数を超えることが可能で、各レコードには acknowledge、release、reject、ロックタイムアウトのパスが存在します。従来のグループにおける割り当てと順序性の直感が変わるため、ビジネスキーで順序が必要な場合は Consumer Group を維持するかバージョンチェックを追加します。

私ならすべてのレコードに冪等性キーを割り当て、p99 処理時間から取得ロックの時間を設定し、アクティブなロック数と下流の並行数を制限します。一時的な障害はバックオフを伴って release し、決定論的な不正データは隔離用トピックに reject し、試行回数の閾値によって自動再試行を停止します。acquire、acknowledge、release、reject、タイムアウト、重複の副作用、完了レイテンシを監視します。移行前には、別のトピックでバージョン、ACL、クライアントの挙動、注入された障害を検証します。レコードの取得、ビジネスの書き込み、確認応答が1つの実績あるトランザクション境界を共有していない限り、この設計は Exactly-once ではなく、冪等性を伴う At-least-once であると位置付けます。」

よくある間違い

  • Share Group を RabbitMQ のクローンと呼ぶ → ストレージ、リプレイ、管理方法が異なる → KIP-932 で文書化されている協調的取得と確認応答のセマンティクスのみを保証として提示する。
  • ワーカー数をパーティション数で制限する → Share Group では複数のメンバーが1つのパーティションを処理できる → ロック、下流のキャパシティ、エンドツーエンドのレイテンシによって並行数を制限する。
  • キーの順序が維持されると思い込む → 並行取得と再配信により完了順序が変わる → 順序が要件である場合はキーのシリアライズやバージョンチェックを行う。
  • 冪等性なしで確認応答(acknowledge)を行う → クラッシュやロック切れによりレコードが再配信される可能性がある → 確認応答を行う前にビジネスキーで重複を排除する。
  • ポイズンレコードを無制限に release し続ける → 再試行によってロックと下流のリソースが枯渇する → エラータイプ、試行回数、隔離ポリシーに基づいて自動再試行を停止する。
  • 30秒を一律の保証として扱う → ブローカー設定や処理レイテンシは環境によって異なる → デプロイされたロック設定を p99 に照らし合わせてテストする。

フォローアップの質問と回答

1つの注文に対するイベントを厳密に順序付ける必要がある場合はどうしますか?

直接 Share Group に切り替えてはいけません。注文 ID でパーティショニングされた Consumer Group を維持するか、アプリケーションレベルでシリアライズされたステートマシンを使用します。共有取得が必須である場合は、バージョンチェック、前提条件の検証、障害時の順序再調整を追加し、それに伴う複雑性の増加を受け入れます。

コンシューマーが取得後に2分間ハングした場合はどうなりますか?

通常の p99 よりわずかに長いロックを設定し、ロックタイムアウト時にアラートを発報します。期限切れ後の再配信を許可しますが、ビジネス処理の冪等性を必須とします。長時間かかる作業の場合は、ロックを無期限に延長するのではなく、再開可能なステップに分割するか外部リースを使用します。

スキーマエラーが5回連続で発生した場合はどのように対処しますか?

決定論的なエラーとして扱います。閾値に達した後に reject し、ペイロード、スキーマバージョン、理由を隔離用トピックに書き込みます。コンシューマーを修正した後、制御された手順に従ってリプレイします。メインのワークフローで無制限に再試行を続けさせてはなりません。

現在のシステムが Kafka トランザクションに依存している場合、直接切り替えることはできますか?

トランザクションの読み取り、処理、書き込みの境界を洗い出し、Share Group における実際のクライアントおよびトランザクションサポート状況を検証します。外部への副作用が同一トランザクション外にある場合は、アウトボックス、冪等性キー、補償処理を使用します。サポートが実証されていない場合は Consumer Group を維持します。

移行によって顧客への二重請求が発生しなかったことをどのように証明しますか?

重複排除制約を持つ一意のビジネスキーのもとで、すべての実行を記録します。クラッシュ、ロック切れ、再試行、ロールバックを意図的に注入し、実行回数と確認応答数を比較します。副作用の発生数、レイテンシ、エラー率の不変条件が維持されている場合にのみ、トラフィックを拡大します。

公開情報ソース

関連する質問