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
相关产品推荐
相关产品推荐

