非同期 @KafkaListener 戻り型
バージョン 3.2 以降では、@KafkaListener (および @KafkaHandler) メソッドに非同期戻り値の型を指定して、応答を非同期に送信できるようになりました。戻り値の型には、CompletableFuture<?>、Mono<?>、Kotlin suspend 関数が含まれます。
@KafkaListener(id = "myListener", topics = "myTopic")
public CompletableFuture<String> listen(String data) {
...
CompletableFuture<String> future = new CompletableFuture<>();
future.complete("done");
return future;
}@KafkaListener(id = "myListener", topics = "myTopic")
public Mono<Void> listen(String data) {
...
return Mono.empty();
}AckMode は自動的に MANUAL を設定し、非同期戻り値の型が検出されるとアウトオブオーダーコミットを有効にします。代わりに、非同期操作が完了すると、非同期完了が確認応答されます。非同期結果がエラーで完了した場合、メッセージが回復されるかどうかはコンテナーのエラーハンドラーによって異なります。非同期結果オブジェクトの作成を妨げる何らかの例外がリスナーメソッド内で発生した場合は、その例外をキャッチし、メッセージを確認応答または回復させる適切な戻りオブジェクトを返さなければなりません。 |
KafkaListenerErrorHandler が非同期戻り値の型 (Kotlin サスペンド関数を含む) を持つリスナー上で構成されている場合、エラーハンドラーは失敗後に呼び出されます。このエラーハンドラーとその目的の詳細については、"例外の処理" を参照してください。
リスナーが配信ごとに非冪等な処理(送信 HTTP 呼び出し、オフセットコミットを伴うトランザクションではない DB 書き込みなど)を実行する場合は、以下のいずれかを選択してください。
この選択は、アプリケーション全体ではなく、トピックごとに行うことができます。パーティションごとの厳密な順序付けが正当性を保証するイベント(エンティティ ID をキーとする状態遷移など)は、シークベースのハンドラー上で実行され、スループットはパーティション数とキーの分散に依存します。順序付けを必要としないイベント(検索インデックスの更新、通知、ダウンストリームのエンリッチメントなど)は、 非同期戻り値型では、同じパーティション上のレコードで結果が混在する場合、パーティションごとの処理順序が保持されません。失敗したレコードが再試行 / バックオフサイクル中である間に、同じパーティション上の後段オフセットの成功したレコードが先にリスナー本体を完了する可能性があります。Kafka は引き続き順序通りに配信し、コミットされたオフセットは一貫性を保ちますが、リスナー本体は並行して実行されるため、副作用がパーティションの順序から外れる可能性があります。ブロッキングリスナーは、パーティションを厳密に 1 つのレコードずつ処理するため、この問題を回避します。パーティション上の次のレコードが呼び出される前に、前段オフセットのレコードが完了 (成功または回復) します。コンシューマーがキーごとの処理順序 (エンティティ ID をキーとする冪等更新、順序付き状態遷移) に依存している場合は、どのエラーハンドラーを使用するかに関わらず、ブロッキングリスナーが適切なモデルです。 |