最新の安定バージョンについては、Spring for Apache Kafka 4.1.1 を使用してください!

非同期 @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 サスペンド関数を含む) を持つリスナー上で構成されている場合、エラーハンドラーは失敗後に呼び出されます。このエラーハンドラーとその目的の詳細については、"例外の処理" を参照してください。

DefaultErrorHandler のようなシークアフターハンドリングエラーハンドラーを使用すると、非同期リスナーは、同じパーティション上の失敗したレコードの後ろにキューイングされているレコードを再実行できます。先頭レコードが再試行サイクルに入ると、コンシューマーはオフセットまでシークバックされ、以降のすべてのポーリングで、そのパーティション上の後続レコードに対してリスナーが再フェッチおよび再呼び出しされます。リスナーの呼び出し回数は、maxAttempts だけでなく、先頭レコードの後ろにキューイングされている失敗したレコードの数に応じて増加します。

This does not affect blocking listeners, because the synchronous throw aborts the batch iteration before the later records are invoked. For async listeners (suspend / Mono / CompletableFuture), the dispatch loop has already invoked the next record before the failure callback fires.

リスナーが配信ごとに非冪等な処理(送信 HTTP 呼び出し、オフセットコミットを伴うトランザクションではない DB 書き込みなど)を実行する場合は、以下のいずれかを選択してください。

  • Use @RetryableTopic (non-blocking retry). Failing records are routed to a separate retry topic so the main partition keeps advancing and later records are not re-executed during a peer’s retry.

  • リスナーを冪等にします(Kafka のデフォルトの少なくとも 1 回配信はすでにこれを想定している)。

This choice can be made per topic rather than application-wide. Events whose correctness depends on strict per-partition ordering (for example, state transitions keyed by an entity id) can stay on a seek-based handler and rely on partition count plus key distribution for throughput. Events that do not require ordering (search index updates, notifications, downstream enrichment) can run on @RetryableTopic listeners so a failing record does not re-execute its partition peers.

非同期戻り値型では、同じパーティション上のレコードで結果が混在する場合、パーティションごとの処理順序が保持されません。失敗したレコードが再試行 / バックオフサイクル中である間に、同じパーティション上の後段オフセットの成功したレコードが先にリスナー本体を完了する可能性があります。Kafka は引き続き順序通りに配信し、コミットされたオフセットは一貫性を保ちますが、リスナー本体は並行して実行されるため、副作用がパーティションの順序から外れる可能性があります。ブロッキングリスナーは、パーティションを厳密に 1 つのレコードずつ処理するため、この問題を回避します。パーティション上の次のレコードが呼び出される前に、前段オフセットのレコードが完了 (成功または回復) します。コンシューマーがキーごとの処理順序 (エンティティ ID をキーとする冪等更新、順序付き状態遷移) に依存している場合は、どのエラーハンドラーを使用するかに関わらず、ブロッキングリスナーが適切なモデルです。