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

升级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环境因网络或集群配置更容易触发发送异常,所以升级后问题仅在该环境显现。

修复步骤:

  1. 调整返回类型:将send方法的返回类型从Mono<Object>改为Mono<Void>,Kafka生产者发送成功通常无需返回具体对象,空类型更符合语义,也能避免类型转换冲突。
  2. 对齐依赖版本:确保项目中micronaut-kafka与micronaut-reactor的版本和Micronaut 4.7.6完全匹配,避免版本不一致引发的兼容性问题。
  3. 添加异常处理:在调用send方法时显式处理异常,比如用onErrorResume将异常转换为合法的Mono实例,绕过框架层面的类型转换失败场景。

修改后的示例代码:

@Topic("${test}")
@Override
Mono<Void> send(@KafkaKey String key, MyDTO value, Headers headers);

内容的提问来源于stack exchange,提问作者Anshu Pandey

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 18:43:16