代表的な面接トピック

システムデザイン面接:Kafkaグループのコオペラティブリバランスへの移行

システム設計難しい
Offer.cc 編集チーム公開日 更新日

質問

高スループットのKafkaコンシューマーグループにおいて、スケール変更のたびに全パーティションが一時停止してしまいます。メッセージ損失を防ぎ、ロールアウト中の重複を制限し、メンバーのクラッシュや設定ロールバックから復旧できる、eagerリバランスからCooperativeStickyAssignorへの移行を設計してください。

課題と背景

この設問では、Kafkaの設定を単に暗記しているかではなく、分散所有権の移譲を理解しているかが試されます。Eagerリバランスは最初に全パーティションを取り消す(revoke)のに対し、コオペラティブリバランスはメンバーが移動の不要なパーティションを保持し続け、移行対象のセットのみを取り消します。完全な回答には、メンバーのバージョン、アサイナープロトコル、オフセットコミット、障害ウィンドウ、テレメトリが含まれます。

面接官の評価ポイント

  • eagerプロトコルとコオペラティブプロトコルの間にあるrevokeとハンドオフの違いを説明できるか。
  • 単一メンバーが単独で切り替えられると思い込まず、互換性のあるローリング順序を計画できるか。
  • 所有権、オフセット、処理中の作業(in-flight work)、コミットのタイミングが整合しているか。
  • クラッシュ、タイムアウト、重複、ロールバック、キャパシティ制限が処理されているか。

明確化のための質問

クライアントのバージョン、現在のアサイナー一覧、静的メンバーシップの有無、メッセージごとの処理時間、許容される重複ウィンドウ、リバランスのレイテンシバジェットを確認します。処理が冪等であるか、シンクが重複排除をサポートしているか、ロールアウトを段階的に実施およびロールバック可能かを尋ねます。ピーク時のスケーリング、パーティション数、アラート閾値を設定します。

30秒の回答アウトライン

まず、互換性のあるアサイナー一覧を維持しながら、全メンバーをコオペラティブプロトコルをサポートするバージョンにアップグレードします。プロトコルサポートが統一されたら、優先戦略としてCooperativeStickyAssignorをロールアウトし、古いものを削除します。各リバランスでは移動するパーティションのみがrevokeされます。コンシューマーはそれらのフェッチを停止して完了済みオフセットをコミットし、新しい所有者はコミット済みオフセットから再開します。リバランス回数、revokeされたパーティション、処理レイテンシ、重複を監視します。クラッシュ時はセッションタイムアウトとオフセットリカバリに依存し、ロールバックは別のローリング変更を通じて古い互換設定を復元します。

ステップごとの解決策

1. プロトコル互換性マトリクスの定義

アサイナーはグループレベルで交渉(ネゴシエーション)されるため、1つのインスタンスを変更するだけでは不十分です。コオペラティブリバランスに依存する前に、それを理解できるバージョンへすべてのクライアントをアップグレードします。第1フェーズでは互換性エントリを維持し、グループの準備が整った後にコオペラティブを優先して古い戦略を削除します。

properties
partition.assignment.strategy=\
org.apache.kafka.clients.consumer.CooperativeStickyAssignor,\
org.apache.kafka.clients.consumer.RangeAssignor

各フェーズの後に、グループのネゴシエーションされたプロトコルと割り当てを検証します。設定ファイルを確認するだけでは、実行中のグループが切り替わった証拠にはなりません。

2. パーティションのrevokeとハンドオフの設計

コオペラティブのrevokeセットには、移動が必要なパーティションのみが含まれます。revoke時には、それらのパーティションのフェッチを停止し、現在の安全なバッチを完了または中止して、完了済みオフセットをコミットします。revokeされていないパーティションの処理は継続します。新しい所有者はコミット済みオフセットから開始するため、ビジネスの冪等性キーまたはシンクの重複排除によって重複に対処します。

3. オフセットと処理中作業の整合

ビジネス上の副作用の実行後にコミットし、決してその前にはコミットしません。バッチ処理中にrevokeが発生した場合は、停止フラグを設定して安全なポイントで終了します。期限が迫った場合はフェッチを停止し、未完了のバッチを記録します。非同期処理では、連続して完了したプレフィックスのみがコミットされるように、パーティションごとのシーケンス追跡が必要です。

4. ローリングロールアウトの計画

ロールアウトコントローラーはメンバーを小さなバッチ単位で再起動し、各バッチの後にグループの安定とラグの回復を待ちます。ベースラインを記録し、クライアントをアップグレードし、プロトコルを観察し、優先アサイナーを切り替え、スケール変更をリハーサルした上で、バッチサイズを増やします。リバランスストームやレイテンシ違反が発生した場合は一時停止し、セッションタイムアウトとmax poll intervalを同時に変更することは避けます。

5. クラッシュおよびロールバック動作の設計

メンバーがクラッシュした場合、セッションタイムアウトが切れた時点でコーディネーターがそのパーティションを再割り当てします。代替メンバーは最後にコミットされたオフセットから再開するため、クラッシュ前の副作用が繰り返される可能性があります。ロールバックでは互換性リスト内の古いアサイナーを復元し、同じローリング順序を使用します。稼働中の新しいメンバーを強制的に削除してはいけません。競合(レース)の診断のために、ジェネレーション、メンバーID、revokeセット、コミット失敗を記録します。

6. キャパシティとテレメトリのガードレール追加

リバランスの頻度と所要時間、revokeされたパーティション数、コンシューマーラグ、poll間隔、コミットレイテンシ、重複率、未割り当てメンバーを追跡します。同時再起動、ホットパーティション、max.poll.intervalを超える処理、ネットワークジッター、メンバー数に近いパーティション数をテストします。キャパシティが不足している場合は、ロールアウトのバッチサイズを縮小するか、コンシューマーを追加してから続行します。

質の高い模範解答

すべてのクライアントがコオペラティブ割り当てをサポートしていることを確認した上で、2フェーズのローリング設定を使用します。バージョンをアップグレードする間は互換アサイナーを維持し、その後にのみCooperativeStickyAssignorを優先して古い戦略を削除します。Revokeコールバックは移動対象のパーティションのみを停止し、安全なポイントを完了させて連続したオフセットをコミットします。保持されたパーティションは継続します。新しい所有者はコミット済みオフセットから再開し、冪等性によって重複を処理します。スモールバッチコントローラーがリバランス、ラグ、poll間クス、コミット失敗、重複率を監視します。クラッシュ時はセッションタイムアウトとオフセットによって回復し、ロールバックは稼働中のメンバーを強制削除するのではなく、同じ互換ローリング順序に従います。

よくある間違い

  • 1つのコンシューマーのみを変更し、グループレベルのアサイナーネゴシエーションを無視する。
  • 移動するパーティションには依然としてハンドオフが発生するにもかかわらず、コオペラティブリバランスを停止時間ゼロとして扱う。
  • ビジネスの副作用を実行する前にオフセットをコミットする。
  • revoke後にフェッチを行ったり、非連続な非同期結果をコミットしたりする。
  • 複数のタイムアウト設定を同時に変更し、因果関係の証拠を見失う。
  • リバランス頻度、revokeセット、重複を無視してラグのみを監視する。

フォローアップの質問

コオペラティブリバランスは重複ゼロを保証しますか?

いいえ。クラッシュ、コミットの再試行、revokeの境界によって処理が繰り返される可能性があります。目的はグループ全体の停止を減らし、重複ウィンドウを制限することです。シンクには依然として冪等性または重複排除が必要です。

なぜすべてのメンバーがコオペラティブ割り当てをサポートしなければならないのですか?

アサイナープロトコルはグループ全体でネゴシエーションされます。コオペラティブセマンティクスを解釈または実行できないメンバーが存在すると、ネゴシエーションが失敗するか、eager動作が強制される可能性があるため、まず互換性のあるバージョンをロールアウトする必要があります。

処理がmax.poll.intervalを超えた場合はどうなりますか?

poll呼び出しを適切なタイミングで維持しながら、バッチを小さくするか、制御された非同期プールを使用するか、慎重に検討されたパラメータ変更を行います。単にタイムアウトを増やすだけでは、障害検出が遅れ、パーティションの占有が長引く可能性があります。

ロールバックの安全性をどのように検証しますか?

ステージンググループでメンバーのクラッシュ、ネットワークジッター、コミット失敗を注入します。ジェネレーション、オフセット、revokeセット、重複排除された副作用を記録します。未コミットのオフセットをスキップすることなく、古いアサイナーが互換性マトリクス内で安定することを検証します。

公開情報ソース

関連する質問

関連面接ツール

システム設計の回答には「回答する」を使用

まず要件を明確にし、スケール、アーキテクチャ、コンポーネント選定、トレードオフの順に進めます。

ツールを見る