プロンプトとコンテキスト
このデータエンジニアリングおよびストリーム処理の質問は、データエンジニア、ストリーミングプラットフォームエンジニア、バックエンドのデータインフラストラクチャ担当者を対象としています。コンシューマーがクラッシュしたりリバランスが発生したりする可能性があり、プロデューサーは応答の喪失後にリトライすることがあります。結果はまず Kafka に戻されます。回答では、exactly-once をシステム全体の魔法のような保証として扱うのではなく、Kafka のアトミックな処理を任意の外部副作用と明確に区別する必要があります。
面接官が評価するポイント
- exactly-once を出力の可視性と入力オフセットのアトミックな進行に分解して説明できるか。
- トランザクショナルプロデューサー、
transactional.id、read_committed、および手動オフセットコミットを正しく使用できるか。 - アボート、再起動、フェンシング、リバランス時のリカバリを説明できるか。
- データベース、検索エンジン、または HTTP サービスには、独自のトランザクション、冪等性キー、または照合(reconciliation)が必要であることを認識しているか。
回答前の明確化のための質問
まず、出力が Kafka 内にとどまるか、Kafka Streams を使用しているか、デプロイ中にコンシューマーグループがローリングアップデートされるか、外部システムが Kafka と同じトランザクションでコミットする必要があるかを確認します。次に、「入力ごとに1つの可視出力」と「1つの外部副作用」を区別します。後者には出力先との連携が必要です。レイテンシ、トランザクションのバッチサイズ、リトライウィンドウ、許容されるラグについて確認します。
30秒の回答フレームワーク
入力レコード、変換後の出力、およびコンシューマーオフセットを1つの Kafka トランザクションに含めます。プロデューサーは transactional.id を有効にし、コンシューマーは自動コミットを無効にして read_committed を使用します。コミットが成功すると3つすべてが同時に可視化され、アボートされた場合は出力が隠され、オフセットはトランザクション前の状態のままになります。再起動後は、安定した一意のトランザクション ID によって古いインスタンスがフェンシングされます。これは Kafka の読み取りと書き込みのみをカバーします。外部データベースには、結果とオフセットを1つのストレージトランザクションに含めるか、冪等性、アウトボックス、および照合が必要です。
ステップバイステップの解決策
1. 保証の境界を提示する
Kafka の設計では、topic 間の exactly-once を、1つのトランザクション内で出力レコードとコンシューマーの位置をアトミックに更新することと定義しています。これは、関数が1回だけ実行されることや、任意の HTTP リクエストが1回だけ届くことを意味するものではありません。構成を説明する前に、この境界を明確にします。
2. トランザクショナルプロデューサーで出力とオフセットをアトミックに書き込む
自動コミットを無効にし、バッチを処理して出力を送信し、トランザクションの一部としてそのバッチのオフセットを送信します。コアとなる疑似コードは次のとおりです。
producer.initTransactions();
while (running) {
ConsumerRecords<String, Order> records = consumer.poll(timeout);
producer.beginTransaction();
try {
for (ConsumerRecord<String, Order> r : records) {
producer.send(new ProducerRecord<>("orders-enriched", r.key(), enrich(r.value())));
}
producer.sendOffsetsToTransaction(offsets(records), groupMetadata);
producer.commitTransaction();
} catch (AbortableException e) {
producer.abortTransaction();
consumer.seekToCommitted();
}
}出力とオフセットを一緒にコミットすることで、「オフセット前の出力」障害による目に見える重複を防ぎ、「出力前のオフセット」障害によるデータ損失を防ぎます。デプロイされているバージョンに従ってクライアントの例外クラスを処理します。すべての例外をむやみにリトライすることは安全ではありません。
3. コンシューマーがコミット済みトランザクションのみを参照するようにする
isolation.level=read_committed を設定し、enable.auto.commit=false を維持します。read_uncommitted はアボートされたトランザクションのレコードを公開するため、コンシューマーがロールバックされるべき出力を伝播してしまう可能性があります。read_committed はトランザクションマーカーを使用して、可視性をコミット結果と一致させます。
4. 再起動、フェンシング、リバランスを処理する
アクティブなコンシューマーインスタンスごとに、クラスタ全体で一意かつ安定した transactional.id を割り当てます。新しいインスタンスが同じ ID で登録されると、Kafka は古いインスタンスの処理中トランザクションをアボートしてフェンシングし、両方がコミットするのを防ぎます。アボート後、アプリケーションはコンシューマーの位置を再作成するか明示的に巻き戻してバッチを再処理する必要があります。アボートされたトランザクション内で進んだローカルカーソルから処理を継続してはなりません。パーティションの割り当てにより、グループメンバーが一度に1つのパーティションを排他的に所有することが保証されます。
5. 外部システムが連携しなければならない理由を説明する
結果が PostgreSQL に送信される場合、Kafka のトランザクションにデータベースのコミットは自動的に含まれません。データベースが先にクラッシュすると Kafka のオフセットが未コミットのままになるため、リトライには一意キーまたは条件付きバージョン更新が必要です。オフセットを先にコミットすると、書き込みが失われる可能性があります。より堅牢な設計では、結果とオフセットを1つのデータベーストランザクションに含めるか、再生可能なアウトボックス/コネクタレコードを書き込み、送信先で重複排除と照合を行わせます。送信先の連携がない場合は、エンドツーエンドの exactly-once ではなく、検出可能な重複を伴う at-least-once 実行を保証します。
6. 設定のスナップショットではなく障害注入で検証する
sendOffsetsToTransaction の前後でクラッシュさせ、コミット応答を喪失させ、リバランスを発生させ、古いインスタンスをフェンシングし、トランザクションを期限切れにし、シンクを使用不可にします。read_committed で出力をコンシュームし、入力 ID ごとの可視出力数、最終オフセット、再起動後のラグ、およびアボートされたトランザクションのレコードを確認します。外部シンクについては、一意キーの競合、リプレイ、および照合を個別にテストします。アボート率、コンシューマーラグ、処理レイテンシ、フェンシング回数を監視します。
質の高い回答例
まず、主張の適用範囲を限定します。入力と出力の両方が Kafka である場合、Kafka Streams または同等のトランザクショナルな consume-transform-produce ループを使用します。コンシューマーは自動コミットを無効にし、プロデューサーは安定した一意の transactional.id を持ち、各バッチの出力とオフセットは sendOffsetsToTransaction でコミットされます。ダウンストリームのコンシューマーは read_committed を使用するため、アボートされた出力は不可視になります。再起動時には、同じトランザクション ID によって古いインスタンスがフェンシングされ、アボートが発生した場合は最後にコミットされたオフセットまで巻き戻す必要があります。送信先が PostgreSQL の場合、Kafka の保証をエンドツーエンドの主張に拡大することはせず、結果とオフセットを1つのデータベーストランザクションに含めるか、アウトボックス、冪等性キー、照合を使用します。その後、クラッシュ、応答喪失、リバランス、フェンシングを注入し、入力 ID ごとの可視出力、オフセット、送信先の状態を検証します。
よくある間違い
- 冪等プロデューサーを exactly-once と呼ぶ → 主にプロデューサーのリトライ時の重複ログエントリを防ぐだけであり、出力とオフセットをアトミックにするわけではない → トランザクションとオフセットコミットを追加する。
read_committedのみを設定する → 可視性は変わるが、トランザクションをコミットしたりクラッシュから回復したりするわけではない → 自動コミットをオフにして完全なトランザクションライフサイクルを実装する。- アクティブなインスタンス間で1つの
transactional.idを共有する → インスタンス同士がフェンシングし合い、スループットが不安定になる → アクティブなインスタンスごとに安定した一意の ID を割り当てる。 - アボート後にローカルの位置から継続する → その位置は未コミットのバッチ内にある可能性がある → コミットされたオフセットを再ロードするか、明示的にシークする。
- Kafka によってデータベースの副作用が1回だけ発生すると主張する → 2つのシステム間には自動的なアトミックコミットが存在しない → 送信先トランザクション、冪等書き込み、アウトボックス、または照合を使用する。
フォローアップの質問と回答
exactly-once はビジネス機能が1回だけ実行されることを意味しますか?
いいえ。関数は実行された後、コミット失敗時に再度実行される可能性があります。保証されるのは、Kafka の可視出力とオフセットコミットが一致することです。可能な限り関数に外部副作用を持たせないようにするか、それらの副作用を再現可能かつ照合可能にしておきます。
read_committed が依然としてレイテンシを追加する可能性があるのはなぜですか?
コンシューマーはアボートされたレコードをスキップし、トランザクションが完了するのを待つ必要があるためです。未完了のトランザクションや長いトランザクションは、可視化レイテンシとラグを増加させます。バッチサイズとタイムアウトを制限し、トランザクション期間を監視し、低レイテンシが求められるパスでは at-least-once 処理とのトレードオフを検討します。
トランザクションの途中でリバランスが発生した場合はどうなりますか?
新しい所有者は、最後にコミットされたオフセットから再開します。現在のトランザクションはアボートされるべきであり、古いインスタンスは解放されるかフェンシングされ、新しいインスタンスが未コミットのバッチを再処理します。現在のメンバーが引き続きパーティションを所有していない限り、オフセットをコミットしてはなりません。
Kafka の出力を PostgreSQL に送信するにはどうすればよいですか?
ビジネス結果とコンシューマーの位置を1つの PostgreSQL トランザクション内で書き込むか、一意のイベント ID を持つアウトボックス行を書き込み、信頼性の高いリレーで配信します。トランザクションを共有できない場合は、at-least-once 配信、一意性制約、条件付きバージョン更新、照合を採用し、これを exactly-once としてアピールしないようにします。
設計が正しく機能していることを証明するメトリクスは何ですか?
入力イベント ID ごとに read_committed の一意の可視出力数をカウントし、コミット済みオフセット、アボート数およびフェンシング数、コンシューマーラグ、再起動後のリプレイ量、シンクでの重複キー競合と相関させます。障害注入の前後でサンプルを保持します。設定のスナップショットだけでは不変条件を証明できません。