インフラ・DevOps
データ解放:Apache Kafkaのネイティブクラスターミラーリング
Data liberation: Apache Kafka's native cluster mirroring (developers.redhat.com)
要約
Apache Kafkaの新しいクラスターミラーリング機能は、外部ツールを必要とせず、ブローカーに直接組み込まれたクロス・クラスター・レプリケーションを提供します。これにより、バイト単位のデータ転送、正確なオフセットの保持、および簡単なフェイルオーバーが可能になり、ディザスターリカバリやクラスター移行シナリオが簡素化されます。
全文翻訳
データ解放:Apache Kafkaのネイティブクラスターミラーリング
クラスターミラーリング:Kafkaに組み込まれたクロス・クラスター・レプリケーション
2026年9月22日
Federico Valeri
関連トピック:オートメーションと管理, Kafka
関連製品:Red Hat OpenShift
目次:
Apache Kafkaはクラスター内でのデータ移動に優れています。リーダーはフォロワーにレプリケートし、コンシューマーは任意のレプリカからプルし、すべての仕組みが最小限の運用オーバーヘッドで実行されます。クラスター間でデータを移動することは、これほど簡単ではありませんでした。
組織が複数のKafkaクラスターを実行するのには正当な理由があります:地理的分散、コンプライアンス境界、チームの分離、バージョンの分離。しかし、データがクラスターに着地した後、別のクラスターに忠実なコピーを取得するには、常に外部ツール、慎重な調整、および運用上の予期せぬ事態に対する健全な許容度が必要でした。
KIP-1279はこの状況を変えます。クラスターミラーリングは、クロス・クラスター・レプリケーションをKafkaブローカーに直接組み込みます。外部プロセスなし、オフセット変換テーブルなし、再圧縮なし。宛先ブローカーは、フォロワーが既に利用しているのと同じフェッチプロトコルを使用して、ソースクラスターからコミットされたレコードをフェッチし、ローカルパーティションログにバイト単位で追加します。結果として、オフセット、圧縮、およびコンシューマーグループの状態を保持するミラーが得られ、フェイルオーバーがミラーの停止とクライアントのリダイレクトと同じくらい簡単になります。
この記事では、クラスターミラーリングのハイレベルなアーキテクチャ、ミラーパーティションを管理するステートマシン、およびそれらすべてをまとめる一貫性保証について説明します。次に、2つの実践的なシナリオ(ディザスターリカバリとクラスター移行)を説明し、ビデオデモで締めくくります。
新しいアプローチ:クロス・クラスター・レプリケーション
MirrorMaker 2 (MM2) は、Kafka 2.4以降、クロス・クラスター・レプリケーションの標準ツールとして機能しており、ソースクラスターから消費し宛先クラスターに生成する一連のKafka Connectワーカーとして実行されます。クラスターミラーリングは、レプリケーションをブローカーに直接組み込むことで、根本的に異なるアプローチを取ります。
ゼロ外部インフラストラクチャ:ミラーフェッチャー・スレッドは、ブローカープロセス内で実行されます。プロビジョニング、監視、またはスケーリングする必要のあるConnectワーカーはありません。単一のCLIコマンド(kafka-cluster-mirrors.sh --create)でミラーが確立され、別のコマンド(--start)でトピックのレプリケーションが開始されます。ライフサイクルの全体は、トピックやコンシューマーグループと同じAdmin APIを通じて管理されます。
バイト単位の転送:圧縮されたバッチは、解凍または再圧縮されることなく、生のバイトとしてレプリケートされます。gzip、snappy、lz4、またはzstdバッチは、元の形式で宛先に到着します。これにより、解凍/再圧縮の往復処理のCPUオーバーヘッドが排除され、プロデューサーの元の圧縮選択が維持されます。
正確なオフセットの保持:宛先ログは、トピックコンパクションによって残されたギャップを含め、ソースと同じオフセットを維持します。コンシューマーグループはオフセット変換なしでフェイルオーバーします:ソースでのコミット済みオフセットは、宛先でのコミット済みオフセットです。
ワンコマンド・フェイルオーバー:ミラーを停止する(--stop)と、パーティションは決定論的なシーケンスを通過します:フェッチャーが削除され、最後のミラーエポックが永続化され、リーダーエポックがインクリメントされ、保留中のトランザクションが中止され、新しい制御レコードがすべてのプロデューサー状態を期限切れにします。パーティションは宛先で書き込み可能になります。外部の調整やオフセットクエリは不要です。
アンクリーン・リーダー選出のサポート:ソースクラスターでアンクリーン・リーダー選出(ULE)が発生した場合、宛先はリカバリ状態に入ります。レプリケーションを再開する前に、ISRメンバーだけでなく、すべての割り当てられたレプリカが切り捨てられたオフセットに収束するのを待ちます。これにより、ソースが不完全なログを持つリーダーを選出した場合でも、クラスター間のログの一貫性が保証されます。
以下の表は、主な違いをまとめたものです。
側面 | MirrorMaker 2 | クラスターミラーリング
---|---|---
デプロイメント | 外部Connectワーカー | ブローカー組み込み
圧縮 | 解凍と再圧縮 | バイト単位パススルー
オフセット | トピック経由のロスのある変換 | クラスター間で同一
コンシューマー・フェイルオーバー | オフセット同期トピックのクエリ | 直接、変換なし
アンクリーン・リーダー選出 | 処理なし | 完全なログ収束
ソース互換性 | Kafka 2.0+ | Kafka 2.1+
監視 | Connect固有のツール | 標準ブローカーJMXメトリクス
クラスターミラーリングにより、宛先ブローカーはクロス・クラスター・レプリケーションのアクティブな参加者になります。各ブローカーは、標準フェッチプロトコルを使用してソースクラスターから直接データをフェッチし、生のレコードバッチをローカルパーティションログに追加します。ソースと宛先のパーティションは同じトピックIDを共有します。
データレプリケーションを超えて、ブローカーはメタデータ検出、構成同期、グループオフセット同期、およびACL伝播も処理します。帯域幅制御は両側で機能します。宛先ブローカーは、設定可能なレプリケーションレート制限を強制します。ソース側では、ミラーフェッチトラフィックは標準のコンシューマーリクエストとして表示されるため、既存のクライアントクォータメカニズムは変更なしで適用されます。
アーキテクチャ
各宛先ブローカー内では、3つの主要なコンポーネントが連携しています。図1は、それらが互いに、およびソースクラスターにどのように接続されているかを示しています。
MirrorMetadataManager
MirrorMetadataManager (MMM) はオーケストレーターです。すべてのブローカーで実行され、MetadataPublisherインターフェースを実装してKRaftメタデータログの変更に反応します。コントローラーがMirrorTopicStateChangeRecordを書き込むと、影響を受けるパーティションのリーダー上のMMMが、特定の操作(作成、開始、停止、一時停止、再開、リカバリ、削除)をトリガーする状態遷移を駆動します。
MMMはソースクラスターへのAdminクライアント接続も維持します。デフォルトで60秒ごとに、ソースメタデータをリフレッシュします:設定されたインクルード/エクスクルードパターンに一致する新しいトピックの検出、トピック構成の同期、コンシューマーグループオフセットの取得、およびソースクラスターIDが変更されていないことの検証。最後のチェックは、ミラーが誤って別のクラスターを指すように再構成された場合に、サイレントなデータ破損を防ぎます。
ClusterMirrorCoordinator
ClusterMirrorCoordinator (CMC) は状態の永続化を処理します。グループコーディネーターやトランザクションコーディネーターが使用するのと同じコーディネーターパターンに従い、__mirror_stateという内部コンパクション済みトピックのシャードを管理します(デフォルト:コンパクションクリーンアップポリシー、50パーティション、レプリケーションファクター3)。各ミラーパーティションの状態は、このトピック内のキー・バリューレコードとして格納され、リーダーエポックとステートエポックのフェンシングを通じて楽観的同時実行制御が行われます。
MirrorFetcherThread
MirrorFetcherThread (MFT) が重い作業を行います。KafkaのAbstractFetcherThread(クラスター内レプリケーションに使用されるのと同じ基底クラス)を拡張し、ソースからレコードをフェッチしてローカルログに追加します。各スレッドは、ミラーごとに認証資格情報を持つ専用のNetworkClientを維持し、ミラー間でSASL/SSLコンテキストを分離します。フェッチャーマネージャーは、3次元識別子(フェッチャーID、ソースブローカーエンドポイント、ミラー名)でスレッドをキーイングし、きめ細かな負荷分散とソースでのリーダー変更への迅速な応答を可能にします。
ミラーパーティションのライフサイクル
ミラーパーティションは、一連の明確に定義された状態を進行します。図2は、完全なステートマシンを示しています。
図2:ミラーパーティション・ステートマシン。
エラーが発生すると、任意の状態がFAILEDに遷移する可能性があります。
LOG_ALIGNMENT:最初のステップは、ローカルログをソースに合わせることです。最初にミラーリングする場合やソースがサポートされていないためにLME(Last Mirror Epoch)が見つからない場合、ブローカーはゼロに切り捨てて最初からミラーリングしますが、それ以外の場合はデルタのみをミラーリングします。LMEレコードは、前のセッションからの最後のミラーオフセットとエポックを格納します。ブローカーは、最後のミラーオフセットまでのオフセットと最後のミラーエポックまでのエポックのみでレコードを切り捨てるため、あらゆる不一致が解決されます。
EPOCH_FENCING:ブローカーはBumpLeaderEpochsリクエストをコントローラーに送信し、ローカルリーダーエポックを10インクリメントします。再インクリメントしきい値は3です。この保証がないと、宛先のコンシューマーはソースからのコミット済みエポックをローカルエポックを超える値で初期化する可能性があり、リーダーの拒否を引き起こします。