You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Reactive Kafka消费者订阅失效问题排查求助

Reactive Kafka消费者无法正常运行的问题排查

我已实现可正常工作的Reactive Kafka生产者,但编写的Reactive Kafka消费者无法正常运行,同时使用@KafkaListener注解的监听器可以和生产者正常配合。

正常工作的生产者代码

KafkaPayload payload = new KafkaPayload(1, "value 1");
String topic = "a.kafka.topic";
Map<String, Object> props = kafkaProperties.buildProducerProperties();
ReactiveKafkaProducerTemplate<String, KafkaPayload> reactiveKafkaProducerTemplate = new ReactiveKafkaProducerTemplate<>(SenderOptions.create(props));
Mono<SenderResult<Void>> senderMono = reactiveKafkaProducerTemplate.send(topic, payload);
senderMono.log().subscribe();

存在问题的消费者代码

String topic = "a.kafka.topic";

ReceiverOptions<String, KafkaPayload> kafkaReceiverOptions = ReceiverOptions.create(kafkaProperties.buildConsumerProperties());
kafkaReceiverOptions.consumerProperty(ConsumerConfig.GROUP_ID_CONFIG, "kafka.id.1");
kafkaReceiverOptions.subscription(Collections.singletonList(topic));

ReactiveKafkaConsumerTemplate<String, KafkaPayload> reactiveKafkaConsumerTemplate = new ReactiveKafkaConsumerTemplate<>(kafkaReceiverOptions);

Flux<KafkaPayload> kafkaFlux = reactiveKafkaConsumerTemplate
                .receiveAutoAck()
                // .delayElements(Duration.ofSeconds(2L)) // BACKPRESSURE
                .doOnNext(
                    consumerRecord -> System.out.println(
                            "received key=" + consumerRecord.key() + ", " +
                            "value="+ consumerRecord.value() + ", " +
                            "from topic=" + consumerRecord.topic() + " ," +
                            "offset=" + consumerRecord.offset()))
                .map(ConsumerRecord::value)
                .doOnNext(kafkaPaload -> System.out
                                .println("successfully consumed KafkaPlayload="
                                                + kafkaPaload.getValue()))
                .doOnError(throwable -> System.out.println(
                                "something bad happened while consuming" +
                                                throwable.getMessage()))
                .log();
                                       
// kafkaFlux.subscribe();
StepVerifier.create(kafkaFlux)
                .consumeNextWith(payload -> {
                        System.out.println(
                                        "-------ONNEXT AFTER SUBSCRIPTION ----------");
                                        System.out.println("id: " + payload.getId() + " value: " + payload.getValue())
                })
                .verifyComplete();

正常工作的@KafkaListener代码

@KafkaListener(id = "kafka.id.1", topics = "a.kafka.topic")
public void listner(KafkaPayload payload) {
         System.out.println("id: " + payload.getId() + " value: " + payload.getValue());
}

问题原因及修正方案

1. ReceiverOptions配置未生效

ReceiverOptions是不可变对象,调用consumerProperty()和subscription()方法时,会返回新的实例而非修改原对象。你的代码中没有将新实例赋值给原变量,导致后续创建的ReactiveKafkaConsumerTemplate使用的是未配置groupId和订阅主题的初始ReceiverOptions。

修正写法:

ReceiverOptions<String, KafkaPayload> kafkaReceiverOptions = ReceiverOptions.create(kafkaProperties.buildConsumerProperties())
        .consumerProperty(ConsumerConfig.GROUP_ID_CONFIG, "kafka.id.1")
        .subscription(Collections.singletonList(topic));

2. StepVerifier使用错误

Reactive Kafka消费者是无限流(持续监听主题,不会主动结束),但你使用了verifyComplete()方法,该方法会等待流正常终止,这会导致StepVerifier超时失败,表现为消费者似乎没有工作。

修正写法:
使用thenCancel()主动终止流,避免超时:

StepVerifier.create(kafkaFlux)
        .consumeNextWith(payload -> {
            System.out.println("-------ONNEXT AFTER SUBSCRIPTION ----------");
            System.out.println("id: " + payload.getId() + " value: " + payload.getValue());
        })
        .thenCancel()
        .verify();

3. 潜在的时序问题

生产者发送消息是异步操作,若消费者还未完成订阅流程,生产者就已经发送了消息,会导致消费者无法接收到这条消息。测试时可以给消费者预留启动时间,或者确保消费者先于生产者启动。


内容的提问来源于stack exchange,提问作者Michela Dennis

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 01:15:40