Quarkus Kafka消费者启动报错SRMSG00207:通道无下游组件
问题描述
作为Quarkus和Kafka新手,需开发消费者监听两个日志主题(记录系统内Quarkus消息流转阶段)。此前单主题消费正常,但添加两个@Incoming注解后,消费者启动时抛出SRMSG00207、SRMSG00208警告,提示stud-in和curr-in通道无下游组件,无法消费消息。
启动警告日志
2024-03-18 13:55:45,720 WARN [io.smallrye.reactive.messaging.provider] (Quarkus Main Thread) SRMSG00208: The connector 'IncomingConnector{channel:'curr-in', attribute:'mp.messaging.incoming.curr-in'}' has no downstreams 2024-03-18 13:55:45,723 WARN [io.smallrye.reactive.messaging.provider] (Quarkus Main Thread) SRMSG00207: Some components are not connected to either downstream consumers or upstream producers: - IncomingConnector{channel:'curr-in', attribute:'mp.messaging.incoming.curr-in'} has no downstream - IncomingConnector{channel:'stud-in', attribute:'mp.messaging.incoming.stud-in'} has no downstream
配置与代码片段
application.properties配置
## Kafka student log properties mp.messaging.incoming.stud-in.connector=smallrye-kafka mp.messaging.incoming.stud-in.topic=academia-student-logging-topic mp.messaging.incoming.stud-in.key.deserializer=org.apache.kafka.common.serialization.StringDeserializer mp.messaging.incoming.stud-in.value.deserializer=org.apache.kafka.common.serialization.StringDeserializer mp.messaging.incoming.stud-in.health-enabled=true mp.messaging.incoming.stud-in.health-readiness-enabled=true mp.messaging.incoming.stud-in.auto.offset.reset=earliest mp.messaging.incoming.stud-in.enable.auto.commit=false mp.messaging.incoming.stud-in.max.poll.records=1 mp.messaging.incoming.stud-in.client.id=KafkaLogConsumer mp.messaging.incoming.stud-in.retry=true mp.messaging.incoming.stud-in.retry-attempts=-1 mp.messaging.incoming.stud-in.retry-max-wait=20 ## Kafka curriculum log properties mp.messaging.incoming.curr-in.connector=smallrye-kafka mp.messaging.incoming.curr-in.topic=academia-curriculum-logging-topic mp.messaging.incoming.curr-in.key.deserializer=org.apache.kafka.common.serialization.StringDeserializer mp.messaging.incoming.curr-in.value.deserializer=org.apache.kafka.common.serialization.StringDeserializer mp.messaging.incoming.curr-in.health-enabled=true mp.messaging.incoming.curr-in.health-readiness-enabled=true mp.messaging.incoming.curr-in.auto.offset.reset=earliest mp.messaging.incoming.curr-in.enable.auto.commit=false mp.messaging.incoming.curr-in.max.poll.records=1 mp.messaging.incoming.curr-in.client.id=KafkaLogConsumer mp.messaging.incoming.curr-in.retry=true mp.messaging.incoming.curr-in.retry-attempts=-1 mp.messaging.incoming.curr-in.retry-max-wait=20
消费代码片段
@Incoming("stud-in") @Incoming("curr-in") @Retry(delay = 120, delayUnit = ChronoUnit.SECONDS, maxRetries = -1, maxDuration = 300, durationUnit = ChronoUnit.SECONDS) public void receive(ConsumerRecord<String, String> event) throws Exception { // 消费逻辑 }
pom依赖片段
<modelVersion>4.0.0</modelVersion> <groupId>sun.kafka.readers</groupId> <artifactId>integration-kafka-loglistener</artifactId> <version>1.0.0-SNAPSHOT</version> <properties> <compiler-plugin.version>3.10.1</compiler-plugin.version> <maven.compiler.release>11</maven.compiler.release> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding> <quarkus.platform.artifact-id>quarkus-bom</quarkus.platform.artifact-id> <quarkus.platform.group-id>io.quarkus.platform</quarkus.platform.group-id> <quarkus.platform.version>2.16.3.Final</quarkus.platform.version> <skipITs>true</skipITs> <surefire-plugin.version>3.0.0-M7</surefire-plugin.version> </properties> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-smallrye-reactive-messaging-kafka</artifactId> </dependency>
原因分析
SmallRye Reactive Messaging不允许在同一个方法上绑定多个@Incoming注解。框架无法将多个输入通道关联到单个消费方法,导致通道找不到下游处理逻辑,触发SRMSG00207/SRMSG00208警告。
解决方案
提供两种可行方案:
方案1:为每个通道创建独立消费方法
将两个通道的消费逻辑拆分到独立方法中,共同逻辑抽成私有方法复用:
@Incoming("stud-in") @Retry(delay = 120, delayUnit = ChronoUnit.SECONDS, maxRetries = -1, maxDuration = 300, durationUnit = ChronoUnit.SECONDS) public void receiveStudEvent(ConsumerRecord<String, String> event) throws Exception { processLogEvent(event); } @Incoming("curr-in") @Retry(delay = 120, delayUnit = ChronoUnit.SECONDS, maxRetries = -1, maxDuration = 300, durationUnit = ChronoUnit.SECONDS) public void receiveCurrEvent(ConsumerRecord<String, String> event) throws Exception { processLogEvent(event); } private void processLogEvent(ConsumerRecord<String, String> event) { // 这里编写统一的日志处理逻辑 }
方案2:使用通道合并(Merge)
通过SmallRye的合并连接器,将两个输入通道合并为一个,再用单个方法监听:
- 在application.properties中添加合并通道配置:
# 定义合并通道,将stud-in和curr-in的消息合并到logs-in通道 mp.messaging.incoming.logs-in.connector=smallrye-merge mp.messaging.incoming.logs-in.merge=stud-in,curr-in
- 修改消费代码,监听合并后的通道:
@Incoming("logs-in") @Retry(delay = 120, delayUnit = ChronoUnit.SECONDS, maxRetries = -1, maxDuration = 300, durationUnit = ChronoUnit.SECONDS) public void receiveLogEvent(ConsumerRecord<String, String> event) throws Exception { // 处理合并后的消息逻辑 }
内容的提问来源于stack exchange,提问作者Elmar Matthee
相关产品推荐
相关产品推荐

