代表的な面接トピック

データエンジニアリング面接:KafkaのShare GroupはConsumer Groupとどう違うか?

データ難しい
Offer.cc 編集チーム公開日 更新日

質問

マルチテナントのデータプラットフォームにおいて、一部のワークロードをKafkaのConsumer GroupからShare Groupに移行することを検討しています。その相違点、適正、移行ゲート、およびロールバック計画について説明してください。

質問とそれが適用される場面

あるプラットフォームでは、イベントストリームとタスクキューにKafkaを使用しています。従来のConsumer Groupは各パーティションを1つのメンバーに割り当てますが、チームはよりキューに近い並行性を得るためにKafka 4.1のShare Groupを使用したいと考えています。どのワークロードが適しているかを判断し、プレビュー機能のリスクをどのように制御するかを決定してください。

面接官が評価するポイント

  • パーティション分割されたストリーム処理と共有キューの配信セマンティクスの区別。
  • 取得制限(acquisition limits)、確認応答(acknowledgment)、障害時の再配信(redelivery)、および順序性の説明。
  • マルチテナントの公平性、ラグ、冪等性、および可観測性を1つの設計に統合する能力。
  • プレビュー機能に対する互換性、カナリア(canary)、リプレイ、およびロールバックのゲートの設定。

回答前に確認すべき質問

  1. その処理はキーの順序性やパーティションローカルな状態を必要としますか、それとも各レコードの結果整合的な処理だけで十分ですか?
  2. 失敗したレコードは、即時再試行、遅延再試行、またはデッドレターキュー(dead-letter queue)のどれで処理すべきですか?
  3. テナントはトピックとコンシューマキャパシティを共有していますか?また、安定したテナント識別子はありますか?
  4. ダウンストリームの副作用は冪等ですか?また、並行処理や重複処理は許容されますか?
  5. Kafkaのバージョン、クライアント、管理ツール、およびマネージドプラットフォームはShare Groupをサポートしていますか?

30秒の回答フレームワーク

Consumer Groupはパーティションを並行性と順序性の境界として使用するため、ストリーミング集約やキー付き状態に適しています。Share Groupは共有キューに近く、複数のコンシューマが1つのトピックパーティションから異なるレコードを取得できますが、クラスタによってパーティションあたりに取得できるレコード数が制限されます。私ならビジネスセマンティクスによってワークロードを分類し、確認応答、再配信、冪等性、公平性を検証した上で、Consumer Groupへのロールバックパスを維持しながらプレビュー機能をカナリア検証します。

ステップバイステップの詳細解説

ステップ1: API名ではなく処理セマンティクスから始める

処理がパーティションの順序、ウィンドウ状態、またはキーによる集約に依存している場合は、Consumer Groupの割り当ての方が論理的に把握しやすいです。タスクが独立しており、より高い並行性が必要で、キュースタイルの確認応答を受け入れられる場合は、Share Groupが候補になります。ベンチマークで高いスループットが約束されているという理由だけで移行してはいけません。

ステップ2: 並行性と取得境界を比較する

従来のグループの並行性は主にパーティション数によって制限され、1つのメンバーが一度に1つのパーティションを処理します。Share Groupを使用すると、複数のコンシューマが1つのトピックパーティションからレコードを取得できますが、クラスタはパーティションごとに取得できる数を引き続き制限します。バッチサイズ、処理時間、およびダウンストリームのキャパシティの積を測定してください。

ステップ3: 確認応答、障害、および再配信を定義する

移行前に、レコードが成功する条件、障害によって他のコンシューマがレコードを取得可能になるかどうか、および再配信が古い処理と重複する可能性があるかどうかを定義します。外部書き込みには、冪等キー、重複排除テーブル、または反復可能なトランザクションを使用します。回復不能な障害は、理由、テナント、試行回数を付与してデッドレターキューに送信します。

ステップ4: 順序性と状態を処理する

共有キューのセマンティクスは、パーティション順序に関する前提を無効にする可能性があります。キーごとのシリアル処理が必要なワークロードはConsumer Groupに残すか、アプリケーション内でキーレベルのシリアライズとバージョンチェックを実装します。状態ストレージにはイベントバージョン、プロセッサ、および再試行ステータスを記録し、並行更新によって暗黙的な上書きが発生しないようにします。

ステップ5: テナントの公平性とバックプレッシャーを構築する

共有トピック上での特定テナントのバーストは、ノイジーネイバー(noisy neighbor)問題を引き起こす可能性があります。安定したテナント識別子を保持し、テナントごとの待機時間、スループット、およびエラー率を監視します。必要に応じて、アプリケーションクォータ、バッチごとの取得上限、トピックの分離、またはダウンストリームのバルクヘッドを追加します。データベースや外部APIには独自の並行性制限が必要です。

ステップ6: プレビュー環境と運用パスを検証する

Kafkaのドキュメントでは、Share Groupはプレビュー版であり、デフォルトでは有効になっていないとされています。まず、ブローカー、クライアント、管理ツール、メトリクス、障害復旧、およびアップグレードの互換性を検証します。擬似イベントを使用して、再起動、コンシューマの減少、重複確認応答、ブローカーの変更、ラグ、およびデッドレターをテストします。

ステップ7: 移行とロールバックを設計する

まず、小規模でクリティカルではないワークロードをShare Groupにミラーリングします。スループット、p99待機時間、重複、再配信、テナントの公平性、およびダウンストリームのエラーを比較します。元のトピックまたはリプレイ可能なオフセット境界を保持します。順序性、重複による副作用、またはプレビューコンポーネントにリグレッションが発生した場合は、新しいトラフィックを一時停止し、Consumer Groupに戻します。

質の高い模範回答

セマンティクスに基づいて切り分けます。パーティション順序、ウィンドウ状態、またはキー付き集約を必要とするイベントはConsumer Groupに残し、独立して並行処理可能で冪等なタスクはShare Groupの候補とします。従来のグループはパーティションを並行性の境界として使用し、1つのメンバーがそれを処理しますが、Share Groupは共有キューのように動作し、パーティションごとの取得制限を維持しながら、複数のコンシューマが1つのトピックパーティションから異なるレコードを取得できるようにします。移行前に、確認応答と再配信、冪等キー、デッドレター、テナントの公平性メトリクスをダウンストリームのバルクヘッドとともに定義します。Kafka 4.1のドキュメントでShare Groupはプレビュー機能とされているため、バージョンと運用の互換性を検証し、クリティカルでないテナントでカナリア検証を行い、待機時間、重複、ラグ、再配信、テナントごとのスループット、およびダウンストリームのエラーを比較します。リプレイ可能なデータとConsumer Groupへのロールバックパスを保持し、順序性や副作用のリグレッションが発生した場合は処理を一時停止して元の構成に戻します。

よくある間違い

  • Share Groupを単に「コンシューマの増加」として扱い、配信および確認応答のセマンティクスを無視すること。
  • パーティション順序に依存するステートフルなストリームを直接共有キューに移動すること。
  • 冪等性を確保せずに再配信や並行処理を許容すること。
  • 全体のスループットばかりに注目し、テナントの待機時間、重複、ダウンストリームの飽和を見落とすこと。
  • クライアント、ツール、アップグレード全体におけるプレビュー機能の互換性を無視すること。
  • 移行後に元のデータを削除し、リプレイやロールバックができなくなること。

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

フォローアップ1: Share Groupはパーティションを不要にするものですか?

いいえ。トピックパーティションは依然としてストレージとレプリケーションの境界です。変わった点は、クラスタの取得制限に従いながら、1つのパーティションから複数のshare-groupメンバーが異なるレコードを取得できるようになったことです。

フォローアップ2: 同じキーに対する順序性は維持されますか?

維持されると想定してはいけません。順序に依存する処理はConsumer Groupに残すか、アプリケーション内でキーごとにシリアライズし、バージョンチェックを使用して並行更新が互いに上書きしないことを保証してください。

フォローアップ3: 失敗したレコードはどうなりますか?

実装と構成を確認し、再度取得可能になるのか、遅延されるのか、並行して再配信される可能性があるのかを把握してください。冪等キー、試行回数、デッドレターの理由を用いて副作用を保護します。

フォローアップ4: テナントのバーストによるキャパシティ圧迫をどう防ぎますか?

テナント識別子を保持し、テナントごとの待機時間とスループットを監視した上で、アプリケーションクォータ、取得上限、トピック分離、またはダウンストリームのバルクヘッドを組み合わせます。公平性をアラート可能な目標値として設定してください。

フォローアップ5: なぜすべてを移行しないのですか?

ストリームワークロードとキューワークロードでは、必要な順序性、状態、再試行、運用の保証が異なります。プレビュー機能にはバージョンや障害のリスクも伴うため、1つのモデルを強制するのではなくワークロードを適切に分類すべきです。

フォローアップ6: 処理済みのレコードをどのようにロールバックしますか?

リプレイ可能なイベントと処理バージョンを保持し、新しいShare Groupへのトラフィックを停止して、Consumer Groupの境界から再開します。外部への副作用は冪等に補償または調整し、盲目的に二重書き込みをしないようにします。

公開情報ソース

関連する質問