升级Micronaut至4.7.6后,Kafka生产者在K8s环境发送消息报错
问题描述
将Micronaut版本从4.2.1升级至4.7.6后,Kafka生产者在Kubernetes环境发送消息时抛出异常,但本地环境运行正常。
错误栈信息
2025-06-28T06:39:45.567247784Z java.lang.IllegalStateException: Cannot convert publisher result: io.micronaut.core.async.publisher.Publishers$JustThrowPublisher@4b137f48 to 'reactor.core.publisher.Mono' 2025-06-28T06:39:45.567252042Z at io.micronaut.aop.internal.intercepted.PublisherInterceptedMethod.lambda$convertPublisherResult$2(PublisherInterceptedMethod.java:118) 2025-06-28T06:39:45.567255027Z at java.base/java.util.Optional.orElseThrow(Optional.java:403) 2025-06-28T06:39:45.567257833Z at io.micronaut.aop.internal.intercepted.PublisherInterceptedMethod.convertPublisherResult(PublisherInterceptedMethod.java:118) 2025-06-28T06:39:45.567260528Z at io.micronaut.aop.internal.intercepted.ReactorInterceptedMethod.convertPublisherResult(ReactorInterceptedMethod.java:47) 2025-06-28T06:39:45.567262962Z at io.micronaut.aop.internal.intercepted.PublisherInterceptedMethod.handleException(PublisherInterceptedMethod.java:100) 2025-06-28T06:39:45.567265287Z at io.micronaut.configuration.kafka.intercept.KafkaClientIntroductionAdvice.intercept(KafkaClientIntroductionAdvice.java:155) 2025-06-28T06:39:45.567268864Z at io.micronaut.aop.chain.MethodInterceptorChain.proceed(MethodInterceptorChain.java:143) 2025-06-28T06:39:45.567271429Z at com..dew.outbound.offer.eligibility.ldso.kafka.DewOfferEligibilityEventProducer$Intercepted.send(Unknown Source) 2025-06-28T06:39:45.567273793Z at com..dew.outbound.offer.eligibility.ldso.kafka.DewOfferEligibilityEventProducer$Intercepted.send(Unknown Source) 2025-06-28T06:39:45.567276117Z at com..dew.core.kafka.listener.DewAbstractKafkaListener.lambda$produceMessage$9(DewAbstractKafkaListener.java:196) 2025-06-28T06:39:45.567278592Z at reactor.core.publisher.FluxConcatMapNoPrefetch$FluxConcatMapNoPrefetchSubscriber.onNext(FluxConcatMapNoPrefetch.java:183) 2025-06-28T06:39:45.567280986Z at reactor.core.publisher.FluxIterable$IterableSubscription.slowPath(FluxIterable.java:335) 2025-06-28T06:39:45.567283511Z at reactor.core.publisher.FluxIterable$IterableSubscription.request(FluxIterable.java:294) 2025-06-28T06:39:45.567286256Z at reactor.core.publisher.FluxConcatMapNoPrefetch$FluxConcatMapNoPrefetchSubscriber.request(FluxConcatMapNoPrefetch.java:337) 2025-06-28T06:39:45.567289162Z at reactor.core.publisher.MonoIgnoreElements$IgnoreElementsSubscriber.onSubscribe(MonoIgnoreElements.java:72) 2025-06-28T06:39:45.567292057Z at reactor.core.publisher.FluxConcatMapNoPrefetch$FluxConcatMapNoPrefetchSubscriber.onSubscribe(FluxConcatMapNoPrefetch.java:164) 2025-06-28T06:39:45.567306745Z at reactor.core.publisher.FluxIterable.subscribe(FluxIterable.java:201) 2025-06-28T06:39:45.567309240Z at reactor.core.publisher.FluxIterable.subscribe(FluxIterable.java:83) 2025-06-28T06:39:45.567311534Z at reactor.core.publisher.Mono.subscribe(Mono.java:4512)
Kafka生产者接口代码
@KafkaClient(id = "test") public interface EventProducer extends DewAbstractKafkaProducer<MyDTO> { @Topic("${test}") @Override Mono<Object> send(@KafkaKey String key, MyDTO value, Headers headers); @Override default String getProducerName() { return KafkaProducerConstants.LDSO; } }
解决方案
这个问题源于Micronaut 4.3+版本对Kafka客户端Reactor类型处理逻辑的变更:当发送消息触发异常时,框架生成的JustThrowPublisher无法被正确转换为Mono,而Kubernetes环境因网络或集群配置更容易触发发送异常,所以升级后问题仅在该环境显现。
修复步骤:
- 调整返回类型:将
send方法的返回类型从Mono<Object>改为Mono<Void>,Kafka生产者发送成功通常无需返回具体对象,空类型更符合语义,也能避免类型转换冲突。 - 对齐依赖版本:确保项目中
micronaut-kafka与micronaut-reactor的版本和Micronaut 4.7.6完全匹配,避免版本不一致引发的兼容性问题。 - 添加异常处理:在调用
send方法时显式处理异常,比如用onErrorResume将异常转换为合法的Mono实例,绕过框架层面的类型转换失败场景。
修改后的示例代码:
@Topic("${test}") @Override Mono<Void> send(@KafkaKey String key, MyDTO value, Headers headers);
内容的提问来源于stack exchange,提问作者Anshu Pandey
相关产品推荐
相关产品推荐

