Quarkus中用SmallRye Kafka发送Kogito消息报错及命名规范咨询
问题解决:Quarkus + Kogito BPMN消息配置错误及命名规范
问题描述
创建了包含多带启动消息和结束消息BPMN的Quarkus应用,消息相关的application.properties配置如下:
kafka.bootstrap.servers=0.0.0.0\:9092 mp.messaging.incoming.endSales.auto.offset.reset=earliest mp.messaging.incoming.endSales.connector=smallrye-kafka mp.messaging.incoming.endSales.topic=endSales mp.messaging.incoming.endSales.value.deserializer=org.apache.kafka.common.serialization.StringDeserializer mp.messaging.outgoing.kogito-processinstances-events.connector=smallrye-kafka mp.messaging.outgoing.kogito-processinstances-events.topic=kogito-processinstances-events mp.messaging.outgoing.kogito-processinstances-events.value.serializer=org.apache.kafka.common.serialization.StringSerializer mp.messaging.outgoing.kogito-usertaskinstances-events.connector=smallrye-kafka mp.messaging.outgoing.kogito-usertaskinstances-events.topic=kogito-usertaskinstances-events mp.messaging.outgoing.kogito-usertaskinstances-events.value.serializer=org.apache.kafka.common.serialization.StringSerializer mp.messaging.outgoing.kogito-variables-events.connector=smallrye-kafka mp.messaging.outgoing.kogito-variables-events.topic=kogito-variables-events mp.messaging.outgoing.kogito-variables-events.value.serializer=org.apache.kafka.common.serialization.StringSerializer mp.messaging.outgoing.kogito_outgoing_stream.connector=smallrye-kafka mp.messaging.outgoing.kogito_outgoing_stream.topic=endSales mp.messaging.outgoing.kogito.outgoing_stream.value.serializer=org.apache.kafka.common.serialization.StringSerializer mp.messaging.outgoing.kogito_outgoing_newEnds.connector=smallrye-kafka mp.messaging.outgoing.kogito_outgoing_newEnds.topic=newEnds mp.messaging.outgoing.kogito_outgoing_newEnds.value.serializer=org.apache.kafka.common.serialization.StringSerializer mp.messaging.outgoing.kogito_outgoing_sudoku.connector=smallrye-kafka mp.messaging.outgoing.kogito_outgoing_sudoku.topic=sudoku mp.messaging.outgoing.kogito_outgoing_sudoku.value.serializer=org.apache.kafka.common.serialization.StringSerializer
构建时出现如下错误:
java.lang.IllegalArgumentException: SRMSG00071: Invalid channel configuration - the `connector` attribute must be set for channel `kogito` at io.smallrye.reactive.messaging.providers.impl.ConnectorConfig.lambda$new$0(ConnectorConfig.java:50) at java.base/java.util.Optional.orElseThrow(Optional.java:403) at io.smallrye.reactive.messaging.providers.impl.ConnectorConfig.lambda$new$1(ConnectorConfig.java:50) at java.base/java.util.Optional.orElseGet(Optional.java:364) at io.smallrye.reactive.messaging.providers.impl.ConnectorConfig.<init>(ConnectorConfig.java:49) at io.smallrye.reactive.messaging.providers.impl.ConfiguredChannelFactory.lambda$extractConfigurationFor$0(ConfiguredChannelFactory.java:85) at java.base/java.lang.Iterable.forEach(Iterable.java:75) at io.smallrye.reactive.messaging.providers.impl.ConfiguredChannelFactory.extractConfigurationFor(ConfiguredChannelFactory.java:74) at io.smallrye.reactive.messaging.providers.impl.ConfiguredChannelFactory.initialize(ConfiguredChannelFactory.java:101) at io.smallrye.reactive.messaging.providers.impl.ConfiguredChannelFactory_Subclass.initialize$$superforward1(Unknown Source) at io.smallrye.reactive.messaging.providers.impl.ConfiguredChannelFactory_Subclass$$function$$4.apply(Unknown Source) at io.quarkus.arc.impl.AroundInvokeInvocationContext.proceed(AroundInvokeInvocationContext.java:54) at io.quarkus.arc.runtime.devconsole.InvocationInterceptor.proceed(InvocationInterceptor.java:62) at io.quarkus.arc.runtime.devconsole.InvocationInterceptor.monitor(InvocationInterceptor.java:49) at io.quarkus.arc.runtime.devconsole.InvocationInterceptor_Bean.intercept(Unknown Source) at io.quarkus.arc.impl.InterceptorInvocation.invoke(InterceptorInvocation.java:41) at io.quarkus.arc.impl.AroundInvokeInvocationContext.perform(AroundInvokeInvocationContext.java:41) at io.quarkus.arc.impl.InvocationContexts.performAroundInvoke(InvocationContexts.java:32) at io.smallrye.reactive.messaging.providers.impl.ConfiguredChannelFactory_Subclass.initialize(Unknown Source) at io.smallrye.reactive.messaging.providers.impl.ConfiguredChannelFactory_ClientProxy.initialize(Unknown Source) at java.base/java.util.Iterator.forEachRemaining(Iterator.java:133) at java.base/java.util.Spliterators$IteratorSpliterator.forEachRemaining(Spliterators.java:1845) at java.base/java.util.stream.ReferencePipeline$Head.forEach(ReferencePipeline.java:762) at io.smallrye.reactive.messaging.providers.extension.MediatorManager.start(MediatorManager.java:192)
错误原因与解决方法
错误原因
错误日志明确指出通道kogito未配置connector,问题出在配置中的一行笔误:
mp.messaging.outgoing.kogito.outgoing_stream.value.serializer=org.apache.kafka.common.serialization.StringSerializer
这里将通道名kogito_outgoing_stream错误地写成了kogito.outgoing_stream,SmallRye Reactive Messaging会将点号解析为配置层级分隔符,误认为存在一个名为kogito的独立通道,但该通道没有对应的connector配置,因此抛出异常。
解决方法
将上述错误配置修改为:
mp.messaging.outgoing.kogito_outgoing_stream.value.serializer=org.apache.kafka.common.serialization.StringSerializer
即将点号替换为下划线,与前面定义的通道名kogito_outgoing_stream保持一致即可解决问题。
Kogito与Kafka通道命名规范
- Kogito自动生成的事件通道:
Kogito默认会为流程实例、用户任务实例、变量变更等事件生成通道,命名格式为kogito-{event-type}-events,例如你配置中的kogito-processinstances-events、kogito-usertaskinstances-events,这类命名是Kogito的约定规则,也可通过自定义配置修改。 - 自定义消息通道:
- 需遵循SmallRye Reactive Messaging的规则:通道名仅允许包含字母、数字、下划线、连字符,禁止使用点号(会被解析为配置层级,引发识别错误)。
- 建议使用业务相关的命名,可统一添加
kogito_前缀标识Kogito关联通道,便于维护区分,例如kogito_outgoing_newEnds、kogito_outgoing_sudoku。 - 通道对应的Kafka主题可自由指定,建议与通道名保持一致或有明确关联,降低维护成本。
内容的提问来源于stack exchange,提问作者codeforHarman
相关产品推荐
相关产品推荐

