プロンプトと適用場面
あるKafkaトピックには24のパーティションがあり、ピーク時には毎秒120,000レコードを受信します。ある1つのテナントがトラフィックの45%を生成しており、プロデューサーはtenant_idでパーティショニングしているため、他のほとんどのコンシューマーがアイドル状態であるにもかかわらず、1つのパーティションだけがラグを蓄積し続けています。現在のハンドラーでは、1つのコンシューマーが毎秒8,000レコードを維持できます。ビジネス要件ではorder_id内でのみ順序性が必要とされており、テナント全体のすべてのイベントにわたる順序性は不要です。また、インシデント中にパーティションを追加することはできません。
コンシューマー、ブローカー、またはダウンストリームの障害ではなく、キーの偏り(key skew)がこの事象を引き起こしたことをどのように証明するか説明してください。次に、現時点でバックログの増加を抑える方法、パーティションキーとマイグレーションを再設計する方法、および非同期処理を追加した後にオフセットを安全にコミットする方法を説明してください。最後に、修正を実証するメトリクスと障害テストで締めくくってください。
これはデータエンジニアリングおよびストリーミングプラットフォームのトラブルシューティングに関する質問です。現在公開されているKafkaの面接ガイドでも、本質的に同じ組み合わせが提示されています。すなわち、1つのホットパーティション、遅延するコンシューマー、他のアイドル状態のコンシューマー、キー指定プロデューサー、そして即時のパーティション数変更不可という状況です。この設問では、プロデューサーのパーティショニング、コンシューマーの並列処理、および順序性を結びつけて説明することが候補者に求められます。Apache Kafkaのドキュメントには、デフォルトのプロデューサーは指定されたキーのハッシュ値によってパーティションを選択すると記載されています。セマンティックパーティショニングは選択されたパーティション内での局所性と順序性を維持しますが、同時にそのキーのトラフィックをそこに集中させます。Huawei CloudのKafkaガイダンスでも同様に、1つのパーティションは同時に1つのコンシューマーによってのみ消費可能であり、一時的にパーティションを追加しても既存のパーティションのバックログを迅速に解消することはできないと述べられています。
このプロンプトにおけるスループット、トラフィック比率、コンシューマーの処理能力、および順序性の範囲は面接用の前提条件であり、特定の企業の実稼働環境の数値ではありません。
面接官が評価するポイント
第1のシグナルは、集約されたメトリクスからパーティションレベルの証拠へと掘り下げられるかどうかです。トピック全体のラグ、コンシューマーの平均CPU使用率、コンシューマー数などは、すべて1つのホットパーティションの存在を覆い隠してしまう可能性があります。優れた回答では、各パーティションのプロデュースレート、消費レート、ラグの傾き、リーダーブローカー、キー分布、およびダウンストリームの処理時間を同じタイムライン上に並べます。これにより、プロデューサーの偏りと、低速なハンドラー、頻発するリバランス、ブローカーのリソース圧迫とを切り分けます。
第2のシグナルは、インシデントを定量化できるかどうかです:
120,000 × 45% = 54,000 records/secondそのパーティションに割り当てられたコンシューマーが毎秒8,000レコードしか処理できない場合、バックログは次のペースで増加します:
54,000 - 8,000 = 46,000 records/second
46,000 × 600 = 27,600,000 records in 10 minutesこの計算により、通常のコンシューマーを追加してもこのパーティションの上限が変わらない理由や、アラートのしきい値を変更してもインシデントが緩和されない理由が説明できます。
第3のシグナルは、真の順序性の不変条件(invariant)を特定できるかどうかです。現在のキーは順序性のスコープをテナント全体に広げていますが、ビジネス上必要なのは注文内での順序のみです。キーをorder_idに変更すれば、カーディナリティが高まり新規の注文が分散されますが、それにはマイグレーションによって1つの注文が2つのパーティションやトピックにまたがらないようにすることが前提となります。
第4のシグナルは、オフセットの正確性です。1つのパーティションからのレコードをワーカープールで処理すると、完了順序がオフセットの順序と異なる場合があります。オフセット104が実行中である間にオフセット105が完了した場合、105までコミットしてしまうとクラッシュ後に104がスキップされる恐れがあります。堅牢な設計では、連続して完了した最大のオフセット(highest contiguous completed offset)を追跡し、完了済みだが未コミットのレコードが再実行される可能性があるため、ダウンストリームへの影響を冪等(idempotent)にします。
最後のシグナルは、ハードリミットを明確に述べられるかどうかです。パーティションを増やせば並列スロットは増えますが、単一のキーに拘束されている巨大なエンティティを分割することはできません。従来のコンシューマーグループにおいてコンシューマーを増やしても、2つのコンシューマーが同一パーティションを同時に所有することはできません。もし単一の注文だけで安全な単一パーティションの容量を超えており、かつ厳密な順序維持が必要な場合、残された手段はシリアルパスの最適化、スロットリング、または独立したビジネスシーケンスの再定義のみとなります。
回答前に明確にすべき質問
- 順序性はテナント単位、注文単位、あるいはさらに小さなイベントストリーム単位で必要ですか? テナント全体での順序が必須である場合、テナントの分割は無効です。注文単位の順序で十分であれば、
order_idがより正確なパーティション境界となります。 - 45%という数値はレコード数、バイト数、それとも処理コストのどれを表していますか? レコード数が均等に見えても、巨大なレコードや負荷の高いダウンストリームへの書き込みによってコストの偏りが発生する可能性があります。レコード数、バイト数、およびハンドラーの処理時間を調査します。
- プロデュースレートが上昇したのか、それとも消費能力が低下したのか? 8,000レコードの処理能力に対して毎秒54,000レコードが安定して流入している場合は、キーの偏りとキーあたりの容量不足を示しています。流入量が変わらないのに消費量が8,000から2,000に低下した場合は、まずダウンストリームシステム、ガベージコレクション、ネットワーク、ディスク、およびリバランスを調査します。
- コンシューマーはどこに書き込みますか? Kafka-to-Kafkaパイプラインであれば、Kafkaトランザクションを使用して出力オフセットと入力オフセットをアトミックにコミットできます。データベース、オブジェクトストレージ、または外部APIの場合、通常はat-least-once(少なくとも1回)処理に加えてビジネス冪等性キーまたは
topic-partition-offsetが必要になります。 - 注文はどのくらいの期間アクティブであり続けますか? 短命な注文であれば、新規注文が新しいルートを使用する一方で、既存注文が完了するまでレガシールートに残すことができます。長命な注文の場合は、明示的なバリア、シーケンス、またはルーティング状態が必要です。
- インシデント対応でトラフィックのスロットリングや縮退運転は可能ですか? テナントのクォータ制限、再構築可能な分析イベントの遅延、または状態更新の集約(coalescing)は、コードやパーティションのマイグレーションよりも迅速に入力量を削減できます。
- コンシューマーが
max.poll.interval.msを超過していませんか? 重い処理によってpollスレッドがブロックされると、リバランスが発生してラグが増幅します。単にタイムアウトを延ばすのではなく、ポーリングと処理を分離し、処理中の作業量を制限します。
30秒回答フレームワーク
「集約されたラグをパーティションごとのプロデュースレート、消費レート、ラグの傾きに分解し、それらをキーの頻度、バイト数、ハンドラー時間、リバランスログ、リーダーブローカーのメトリクスと照合します。ホットテナントは毎秒54,000レコードを生成し、1つのコンシューマーは8,000レコードを処理するため、バックログは毎秒約46,000件増加します。通常のコンシューマーを追加してもそのパーティションを加速することはできません。現時点の対策としては、ホットテナントのスロットリングまたは縮退、専用インスタンスへのパーティションの分離、そして注文ごとに直列化しつつ異なる注文を並行処理することで真の順序性境界を活用します。オフセットは連続して完了した最大のオフセットのみをコミットし、ダウンストリームへの書き込みは冪等にします。恒久対策としては、新規注文をorder_idをキーとする新しいトピックに移行し、既存注文は完了するまでレガシールートに残します。偏った負荷、コンシューマーのクラッシュ、リバランスの条件下で、パーティションレベルのラグの傾き、エンドツーエンドのp99、重複、順序違反を検証します。」
ステップバイステップの詳細回答
ステップ1: どのレイヤーがホットスポットを生み出したかを証明する
ピーク時のタイムウィンドウを使用して、以下を照合します:
- パーティションごとの毎秒プロデュースレコード数、毎秒バイト数、およびハイウォーターマークの増加量
- パーティションごとの毎秒消費レコード数、コミット済みオフセット、およびラグの傾き
- キーの頻度、バイト数、および推定処理コストの分布
- ホットパーティションのリーダーブローカーにおけるCPU、ネットワーク、ディスク待機、およびリクエストレイテンシ
- 割り当てられたコンシューマーのpoll間隔、バッチサイズ、ハンドラーレイテンシ、ガベージコレクション、エラー、およびリバランスログ
- ダウンストリームのデータベース、ストレージシステム、またはAPIにおけるパーティション相関のレイテンシやスロットリング
判断基準は数値から導き出されます。ホットパーティションには毎秒約54,000レコードが届きます。他の23パーティションで残りの66,000レコードを共有するため、残りがほぼ均等であれば1パーティションあたり平均約2,870レコードとなります。8,000レコード処理できるコンシューマーは通常のパーティションでは余力がありますが、ホットパーティションの流入には追いつきません。これによって、1つのパーティションのみが増加し他がアイドル状態になる理由が完全に説明されます。もしホットパーティションへの流入が通常通りであるにもかかわらずダウンストリームのレイテンシに連動して消費能力が低下している場合は、キー設計が根本原因であるとはまだ証明されていません。
ブローカーの配置も確認してください。リーダーが過負荷なブローカー上にあるパーティションは、プロデュースやフェッチが遅くなる可能性があります。リーダーシップの移動やレプリカの再配置によって配置のボトルネックは解消できますが、1つのtenant_idが依然として1つのパーティションにマッピングされているという事実は変わりません。
ステップ2: バックログを排出する前に増加の傾きを抑える
インシデント対応の最初の目標は以下の通りです:
hot-partition input rate ≤ hot-partition safe processing rate最も迅速な手段は受付制御(admission control)です。ホットテナントに明示的なクォータを適用する、再構築可能な分析イベントを遅延させる、最新の状態のみが重要な更新を集約する、または重要でない処理を縮退パスに回します。すべての施策において、損失、遅延、再試行のセマンティクスを明示する必要があります。レコードを通知なく破棄することはスロットリング戦略ではありません。
コンシューマー側では、制御されたメンテナンスウィンドウを設けて既存グループを停止し、排他的で重複のない明示的割り当てに置き換えることができます。十分なリソースを持つ1つのインスタンスがホットパーティションのみを担当し、残りのインスタンスが他のパーティションを担当します。2つ目の通常のコンシューマーグループを立ち上げても助けにはなりません。トピック全体を個別にフル消費してしまい、ビジネス処理の重複を引き起こすだけです。分離を行うことで、同一プロセス内でホットパーティションが通常のパーティションのリソースを奪うことを防止できますが、従来のシリアルハンドラーの処理能力が毎秒8,000レコードを超えるわけではありません。
ビジネス要件では注文内での順序のみが必要とされるため、ホットコンシューマーはorder_idによってディスパッチできます。アクティブな注文ごとに1つのシリアルキューを割り当て、異なる注文間では制限付きのワーカープールで並行処理します。プールには処理中タスク(in-flight)の上限を設ける必要があります。プールが満杯になったら、パーティションを一時停止(pause)するかワーカーへの解放量を抑え、Kafkaのラグが無制限なプロセスメモリ消費につながらないようにします。pollループは応答性を維持しなければなりません。そうしないとmax.poll.interval.msを超過してリバランスがトリガーされ、さらなる停止が発生します。
ステップ3: 連続した完了ウォーターマークをコミットする
パーティション内並行処理によって完了順序は変わりますが、コミット順序を変えてはなりません。各パーティションで次の状態を維持します:
nextCommitOffset = smallest unfinished offset
completed = offsets that finished but still have a gap before them
onComplete(offset):
add offset to completed
while completed contains nextCommitOffset:
remove nextCommitOffset from completed
increment nextCommitOffset
commit nextCommitOffsetKafkaは次に読み取るべき位置をコミットします。したがって、104、105、106がすべて完了した後に初めて、コミット位置を107に進めることができます。104が再試行中である間に105が完了した場合、コミットポイントは104にとどまります。次のコミット前にクラッシュすると完了済みのレコードが一部再実行されるため、データベースへの書き込みにはevent_idまたはその他のビジネス固有のユニークキーを使用して冪等なupsertを行う必要があります。ビジネスキーが存在しない場合は、topic-partition-offsetでソースレコードを一意に識別できます。
オフセットコミットと外部への副作用が自然にアトミックであると説明してはなりません。Kafka-to-Kafkaアプリケーションであれば、1つのKafkaトランザクション内に出力レコードと消費オフセットを含めることができます。外部データベースの場合、より一般的な規約はat-least-once消費と冪等な書き込み、または重複排除レコードとビジネス更新の両方を含むデータベーストランザクションです。
ステップ4: パーティションのスコープを真の順序性ドメインに合わせる
24パーティションの平均容量が必ずしも不足しているわけではありません。毎秒8,000レコードが測定された安全な最大値であり、計画稼働率の上限を70%とする場合、パーティションあたりの計画容量は5,600となります:
120,000 ÷ 5,600 ≈ 21.4均等に分散されていれば、24パーティションで想定ピークを適度なヘッドルーム(余裕)を持ってカバーできます。障害の原因は、トラフィックの45%をカーディナリティの低い1つのキーに集中させたことにあり、パーティションの総数にあるわけではありません。恒久的なキーは、必要とされる最小の順序性ドメインを表し、十分なカーディナリティを持ち、ピーク時にも予測どおりに分散される必要があります。ここではorder_idが自然な選択です。テナントごとの局所性に運用上の実質的な価値がある場合は、(tenant_id, order_id)の安定したエンコーディングも候補になります。
制御されたソルト(salting)の付与は、元のキー内のレコード順序が入れ替わってもよい場合、またはダウンストリームのステージで順序を復元できる場合にのみ有効です。1つの注文にランダムなソルトを付与すると、その注文が複数のパーティションから同時に届く可能性があるため、本問の不変条件に違反します。単一の注文だけで単一パーティションの容量を超える場合、キーのカーディナリティを高めても効果はありません。その注文のシリアルパスを最適化またはスロットリングするか、ビジネスプロトコルを明示的に独立したシーケンスへと再設計してください。
ステップ5: バージョン管理されたルーティングで移行する
既存のトピックにパーティションを追加し、直ちにキーを変更することには2つのリスクがあります。パーティション数が変わると、デフォルトのハッシュマッピングによって既存のキーの割り当て先が変わる可能性があります。また、切り替えの前後で同一注文のイベントが異なるパーティションに送られる可能性があり、Kafkaはパーティションをまたぐ順序性を保証しません。
より安全な設計は、order_idでパーティショニングされた新しいトピックを作成し、プロデューサーのルーティングをバージョン管理することです:
- 切り替え後に作成された注文は新しいトピックを使用する
- 既存の注文はクローズするまで古いトピックと従来のキーを使用し続ける
- すべてのプロデューサーは、ローカルクロックと切り替え時刻を比較するのではなく、同一の注文ルーティング状態を使用する
- コンシューマーは両方のルートを読み取るが、ある注文はどの瞬間においてもアクティブな1つのルートにのみ属する
- レガシー注文がすべて処理され、保持期間の要件が満たされた後に古いトピックを廃止する
注文が自然に終了しない場合は、注文ごとのマイグレーションバリアを作成します。その注文の新しいイベントを一時停止し、古いルートが記録された最終シーケンスまたはオフセットに達するまで待機し、ルーティングバージョンを変更して再開します。一時停止を伴わない代替案として、単調増加のシーケンス番号を付与してダウンストリームで両方のルートをマージする方法もありますが、バッファリング、タイムアウト、ギャップ回復の処理が導入されます。これは、ビジネス要件がその複雑さに見合う価値を持つ場合にのみ正当化されます。
ステップ6: 偏ったトラフィックと障害を用いて検証する
集約スループットだけでは十分な証拠になりません。少なくとも以下をテストしてください:
- 1つのテナントがトラフィックの45%を生成し、注文の分布が実際のピークに類似しているデータ分布
- パーティションごとの流入量、消費量、ラグの傾き、および最大ラグ
- エンドツーエンドのp50、p95、p99、および推定ドレイン(排出)時間
- 処理中のワーカー数、最も古いタスクの滞留時間、リトライ数、およびデッドレター
- 重複の影響、注文ごとの順序違反、および冪等性の競合
- 完了オフセットにギャップが存在する状態でのコンシューマーのクラッシュ
- 長時間の処理によってリバランスがトリガーされるかどうか、および復旧に要する時間
- 古いトピックと新しいトピックの境界で、注文が一方のルートにのみ現れているかどうか
合格条件には持続的な動作が含まれます。安定したピーク時にホットパーティションのラグの傾きが正(増加傾向)でなくなること、クラッシュ時に処理が再実行されてもビジネス上の影響が失われないこと、順序の狂った注文が観測されないこと、新規注文がパーティション間に分散されること、レガシー注文が予定通りに排出されることです。集約スループットが向上しても、1つの大きな注文が繰り返しホットスポットを発生させる場合、順序性ドメインまたはビジネス受付の問題が未解決のままです。
質の高い模範解答
「プロンプトから、他のコンシューマーがアイドル状態であるにもかかわらず1つのパーティションだけが遅延していることが既に分かっているため、まずコンシューマーを追加することから始めるのは避けます。最初に、パーティションごとの毎秒レコード数、毎秒バイト数、ラグの傾き、およびキーの頻度を使用してキーの偏りを証明し、同時にリーダーブローカー、リバランス、ダウンストリームのレイテンシに起因する可能性を排除します。
ホットテナントは毎秒54,000レコードを生成します。単一パーティションのコンシューマーは8,000レコードを処理するため、ラグは毎秒約46,000レコード、10分間で2,760万件増加します。残りの55%が23のパーティションにおおむね分散されている場合、各パーティションは平均で毎秒約2,870レコードとなります。これにより、ホットコンシューマーの処理能力不足と他の余力が両方説明できます。従来のコンシューマーグループでは、1つのパーティションは一度に1つのコンシューマーにしか属せないため、通常のインスタンスを増やしてもそのパーティションを加速させることはできません。
現時点の対策として、まず文書化されたテナントクォータによって流入の傾きを低減し、遅延許容または集約可能なイベントを縮退パスに移動します。制御されたメンテナンスウィンドウ内で、ホットパーティションを専用インスタンスに分離し、通常のパーティションのリソースが奪われないようにします。注文単位の順序性のみが重要であるため、order_idごとに制限付きプールへディスパッチします(注文内では直列、注文間では並行)。最も早く完了したタスクのオフセットをそのままコミットすることはしません。連続して完了した最大のオフセットを追跡し、ギャップがある箇所でコミットを止めます。クラッシュによって完了済み未コミットのレコードが再実行される可能性があるため、シンク側ではevent_idまたはビジネス固有のユニークキーを使用して冪等性を確保します。
恒久対策として、単にパーティションを追加するだけでは完全な解決になりません。計画稼働率70%の場合、8,000レコード処理可能な各パーティションは毎秒約5,600レコードに寄与するため、均等に負荷分散された24パーティションがあれば想定される120,000レコードのピークをカバーできます。問題はtenant_idによって45%が1つのパーティションに固定されていることです。そこで、order_idをキーとする新しいトピックを作成します。切り替え後の新規注文はこれを使用し、アクティブな既存注文は完了するまで古いルートに残します。共有されたルーティング状態により、1つの注文が両方のトピックにまたがることがないようにします。
本番展開の前に、同じ45%のテナント偏向を再現し、パーティションごとのスループット、ラグの傾き、エンドツーエンドのp99、重複、順序違反を検証します。オフセットの完了順序にギャップがある状態でコンシューマーをクラッシュさせて再起動時に冪等な再実行のみが行われることを確認し、リバランスを発生させて復旧時間を測定し、ルート境界にあるすべての注文が単一のトピックにのみ現れることを検証します。もし単一の注文自体が単一パーティションの容量を超える場合は、その注文のシーケンスモデルを最適化、スロットリング、または変更しない限り、コンシューマーやパーティションを増やしても解決しないというハードリミットを提示します。」
よくある間違い
- トピック全体のラグのみを見る → 平均値は単一パーティションの流入および消費の傾きを隠してしまいます → レコード数、バイト数、ラグ、キー分布をパーティションごとにグラフ化する。
- ラグが増加した際に無条件でコンシューマーを追加する → 従来のグループにおいて1つのパーティションは一度に1つのコンシューマーしか所有できません → まずパーティション数、割り当て状況、パーティションごとの容量を比較する。
- 即座にパーティションを追加する → 既存のバックログは自動的に再分散されず、デフォルトのキーマッピングが変更される可能性があります → まずバックログの増加を止め、バージョン管理されたトピックと移行境界を使用する。
- ホットキーにランダムなソルトを付与する → 1つの注文がパーティションをまたいで順序が狂って届く可能性があります → 順序の入れ替わりが許容される場合にのみソルトを使用する。本問の真の順序性スコープには
order_idを使用する。 - ワーカーが完了したオフセットを即座にコミットする → クラッシュ後に、未完了のより若いオフセットがスキップされる可能性があります → 連続して完了したウォーターマークのみをコミットする。
- オフセットコミットをexactly-once(正確に1回)処理と呼ぶ → 外部データベースの更新とKafkaオフセットは通常、単一のトランザクションにはなりません → at-least-once、冪等性、トランザクションの境界を明記する。
- 支援のために2つ目のコンシューマーグループを立ち上げる → 2つ目のグループはトピック全体の完全なコピーを個別に読み取るため、処理の重複が発生します → 排他的な明示的割り当て、または制御された処理の再設計を行う。
- 均一なトラフィックのみでテストする → 平均値が良好であっても、ホットキーの問題が解消された証明にはなりません → 現実的な偏りを再現し、平均値だけでなく最大のパーティションを検証する。
- 単一エンティティの限界を無視する → 厳密に順序付けられた1つのエンティティを、コストなしでパーティションをまたいで並列化することはできません → 順序性、スロットリング、シリアル処理の間のトレードオフを明確にする。
フォローアップの質問と回答
フォローアップ1: なぜパーティション数を24から48へ即座に増やさないのですか?
新しいパーティションは将来の並列スロットを作成しますが、古いパーティションに既に保存されているバックログを分割することはできず、1つのtenant_idが複数のパーティションにマッピングされるようにもなりません。デフォルトのキーハッシュでは、パーティション数を変更すると既存のキーが再マッピングされ、切り替えを挟んで同一の注文が異なるパーティションに配置される可能性もあります。まず順序性ドメインとマイグレーションを修正し、その上で偏りのある負荷テストを実施して総パーティション数の追加が必要かどうかを判断してください。
フォローアップ2: パーティション内並行処理を追加した後、メモリの無制限な消費をどのように防ぎますか?
パーティションごとに、処理中(in-flight)の最大レコード数と未コミットオフセットの最大ウィンドウサイズを設定します。いずれかの制限に達したら、重い処理とポーリングを分離したまま、そのパーティションを一時停止するかワーカープールへのレコード放出量を減らします。最も古い未完了オフセットの経過時間を追跡します。ある注文がブロックされ続けた場合は、後続のすべてのオフセットが無制限にメモリを消費することを許すのではなく、その注文を分離、再試行、または人手介入用へルーティングします。
フォローアップ3: ダウンストリームのデータベースが冪等なupsertをサポートしていない場合はどうしますか?
単一のデータベーストランザクション内で重複排除レコードの挿入とビジネス更新を実行します。ビジネス上のevent_idまたはtopic-partition-offsetを一意キーとして使用し、一意性制約の競合が発生した場合はそのイベントが既に適用済みであることを意味します。非トランザクショナルな外部APIの場合は、そのAPIの冪等性キー、アウトボックステーブル、またはクエリ可能な操作ステータスを使用します。冪等性の境界を構築できない場合、クラッシュ後の再実行による重複の影響がないことを保証することはできません。
フォローアップ4: 将来的にビジネス要件として厳格なテナント全体の順序が必要になった場合、何が変わりますか?
その場合、tenant_idが分割不可能な順序性ドメインとなり、注文間の並行処理は無効になります。1つのテナントが単一パーティションの容量を超える場合は、そのシリアルパスを最適化するか、テナントをスロットリングするか、どのイベントシーケンスが独立しているかを再交渉する必要があります。グローバルにシーケンス番号が付与されたシャーディングストリームとダウンストリームのマージャーを組み合わせる構成も可能ですが、順序待ち、ギャップ回復、可用性コストがコンシューマー側に転嫁されることになり、無償でスケールできるわけではありません。
フォローアップ5: 2,760万レコードのバックログを排出するのにどのくらい時間がかかりますか?
まず入力量を処理能力以下に落とす必要があります。スロットリングによってホットパーティションの入力を毎秒3,000レコードに抑え、最適化によって処理能力を12,000に引き上げた場合、純排出レートは毎秒9,000レコードになります:
27,600,000 ÷ 9,000 ≈ 3,067 seconds ≈ 51 minutesこれは一定のレートを前提とした試算です。実際の復旧計画では、リトライ、ダウンストリームのスロットリング、レコードサイズのばらつき、安全マージンを加味し、観測されたラグの傾きに基づいて継続的に試算を修正します。
フォローアップ6: 古いトピックと新しいトピックの間で、どの注文もまたがっていないことをどのように証明しますか?
注文ごとに1つのpartitioning_versionを保存します。すべてのプロデューサーはバージョン管理された同一のルーティングレコードを読み取るかキャッシュし、マイグレーションバリアが成功した後にのみバージョンが変更されます。コンシューマーは注文に対して最初に確認したトピックとバージョンを記録し、ある注文が両方のアクティブなルートに現れた場合にアラートを発報してその注文の自動進行を停止します。負荷テストやカナリアリリースの間は、単に総レコード数の一致を確認するだけでなく、プロデューサーのログ、両トピックのオフセット、およびダウンストリームの注文シーケンスを突合して検証します。