Quarkus双@Incoming注解的Kafka事件处理器触发SRMSG00212空指针异常
问题分析与解决方案
环境信息
- Quarkus版本:2.16.4.Final
- 响应式扩展:quarkus-smallrye-reactive-messaging-kafka
问题重现
在Kafka消息处理方法上添加多个@Incoming注解时,启动触发SRMSG00212错误,抛出NullPointerException。示例代码:
@ActivateRequestContext @Acknowledgment(Acknowledgment.Strategy.POST_PROCESSING) @Incoming("WEIGHTS_REQUEST") @Incoming("CREATION_WEIGHTS_REQUEST") public Uni<Void> process2(ConsumerRecord<String, WeightsEvent> record) { WeightsEvent event = record.value(); LOGGER.infof("Get Weight for Weighting Procedure %s", event.weight_proc); String key = record.key(); String topic = record.topic(); int partition = record.partition(); return Uni.createFrom().voidItem(); }
启动错误日志:
08:20:35 ERROR traceId=, parentId=, spanId=, sampled= [io.sm.re.me.provider] (Quarkus Main Thread) SRMSG00212: Unable to initialize mediator: process2: java.lang.NullPointerException: Bean is null at java.base/java.util.Objects.requireNonNull(Objects.java:233) at io.quarkus.arc.impl.BeanManagerImpl.getReference(BeanManagerImpl.java:57) at io.smallrye.reactive.messaging.providers.extension.MediatorManager.createMediator(MediatorManager.java:176) at io.smallrye.reactive.messaging.providers.extension.MediatorManager_ClientProxy.createMediator(Unknown Source) at io.smallrye.reactive.messaging.providers.wiring.Wiring$SubscriberMediatorComponent.materialize(Wiring.java:626) at io.smallrye.reactive.messaging.providers.wiring.Graph.lambda$materialize$10(Graph.java:100) at java.base/java.util.ArrayList.forEach(ArrayList.java:1511) at io.smallrye.reactive.messaging.providers.wiring.Graph.materialize(Graph.java:99) at io.smallrye.reactive.messaging.providers.extension.MediatorManager.start(MediatorManager.java:219) at io.smallrye.reactive.messaging.providers.extension.MediatorManager_ClientProxy.start(Unknown Source) at io.quarkus.smallrye.reactivemessaging.runtime.SmallRyeReactiveMessagingLifecycle.onApplicationStart(SmallRyeReactiveMessagingLifecycle.java:52) at io.quarkus.smallrye.reactivemessaging.runtime.SmallRyeReactiveMessagingLifecycle_Observer_onApplicationStart_7f54e4b27c1b49e5e062caa58f1e82797fa01393.notify(Unknown Source) at io.quarkus.arc.impl.EventImpl$Notifier.notifyObservers(EventImpl.java:328) at io.quarkus.arc.impl.EventImpl$Notifier.notify(EventImpl.java:310) at io.quarkus.arc.impl.EventImpl.fire(EventImpl.java:78) at io.quarkus.arc.runtime.ArcRecorder.fireLifecycleEvent(ArcRecorder.java:131) at io.quarkus.arc.runtime.ArcRecorder.handleLifecycleEvents(ArcRecorder.java:100) at io.quarkus.deployment.steps.LifecycleEventsBuildStep$startupEvent1144526294.deploy_0(Unknown Source) at io.quarkus.deployment.steps.LifecycleEventsBuildStep$startupEvent1144526294.deploy(Unknown Source) at io.quarkus.runner.ApplicationImpl.doStart(Unknown Source) at io.quarkus.runtime.Application.start(Application.java:101) at io.quarkus.runtime.ApplicationLifecycleManager.run(ApplicationLifecycleManager.java:108) at io.quarkus.runtime.Quarkus.run(Quarkus.java:71) at io.quarkus.runtime.Quarkus.run(Quarkus.java:44) at io.quarkus.runtime.Quarkus.run(Quarkus.java:124) at io.quarkus.runner.GeneratedMain.main(Unknown Source) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77) at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.base/java.lang.reflect.Method.invoke(Method.java:568) at io.quarkus.runner.bootstrap.StartupActionImpl$1.run(StartupActionImpl.java:104) at java.base/java.lang.Thread.run(Thread.java:833)
问题原因
Quarkus 2.16.x版本搭配的SmallRye Reactive Messaging,不支持在单个方法上直接用多个@Incoming注解绑定不同通道。这种用法会导致mediator初始化时无法正确获取Bean引用,触发空指针异常。
解决方案
方案1:拆分方法复用逻辑
为每个@Incoming通道单独创建处理方法,核心业务逻辑抽离到公共方法复用:
@ActivateRequestContext @Acknowledgment(Acknowledgment.Strategy.POST_PROCESSING) @Incoming("WEIGHTS_REQUEST") public Uni<Void> processWeightsRequest(ConsumerRecord<String, WeightsEvent> record) { return handleRecord(record); } @ActivateRequestContext @Acknowledgment(Acknowledgment.Strategy.POST_PROCESSING) @Incoming("CREATION_WEIGHTS_REQUEST") public Uni<Void> processCreationWeightsRequest(ConsumerRecord<String, WeightsEvent> record) { return handleRecord(record); } private Uni<Void> handleRecord(ConsumerRecord<String, WeightsEvent> record) { WeightsEvent event = record.value(); LOGGER.infof("Get Weight for Weighting Procedure %s", event.weight_proc); String key = record.key(); String topic = record.topic(); int partition = record.partition(); return Uni.createFrom().voidItem(); }
方案2:升级Quarkus版本
升级到Quarkus 3.x及以上版本,SmallRye Reactive Messaging在新版本中已支持单个方法绑定多个@Incoming注解的用法,可直接保留原代码结构。
验证说明
- 单个
@Incoming注解时功能正常,说明通道配置本身无问题; - 拆分方法后,服务启动正常,两个通道的消息均可被正确处理;
- 升级到Quarkus 3.x后,原多
@Incoming注解的代码可正常运行。
内容的提问来源于stack exchange,提问作者mizmauz
相关产品推荐
相关产品推荐

