Quarkus中Smallrye Kafka按通道配置Broker后无法消费事件
问题排查与修复方案
你的配置存在几个关键问题,导致无法消费Kafka事件,逐一说明:
1. 通用配置前缀错误
你当前的通用配置使用mp.messaging.xxx前缀,这会导致这些配置无法正确关联到SmallRye Kafka连接器。正确的通用配置前缀应该是mp.messaging.connector.smallrye-kafka.xxx——因为你已经通过mp.messaging.connector=smallrye-kafka指定了连接器,所有通道共享的配置需要附加连接器名称作为前缀,否则通道不会继承反序列化器、auto.offset.reset等关键配置,导致消费逻辑无法正常工作。
2. JAAS配置字符串格式错误
你的SASL配置字符串中,username和password之间缺少空格:
private static final String SASL_CONFIGURATION_STRING = "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"%s\"password=\"%s\";";
这会导致JAAS解析失败、认证无效。正确格式需在两者间添加空格,否则客户端无法通过Broker身份验证,即使日志显示连接成功,也没有权限消费消息。
3. 无关配置冗余
mp.messaging.connections.max.idle.ms=-1是Kafka生产者的配置,消费者不需要该参数,可移除避免不必要的配置干扰。
修改后的完整配置示例
application.properties
# Kafka common kafka.health-enabled=false kafka.session.timeout.ms=45000 # Common configuration used for all channels (修正前缀) mp.messaging.connector=smallrye-kafka mp.messaging.connector.smallrye-kafka.value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer mp.messaging.connector.smallrye-kafka.value-deserialization-failure-handler=kafka-value-failure-handler mp.messaging.connector.smallrye-kafka.fail-on-deserialization-failure=false mp.messaging.connector.smallrye-kafka.auto.offset.reset=earliest mp.messaging.connector.smallrye-kafka.use.latest.version=true mp.messaging.connector.smallrye-kafka.auto.register.schemas=false mp.messaging.connector.smallrye-kafka.failure-strategy=ignore # Topic 1 mp.messaging.incoming.my-channel-1.topic=my-topic-1 mp.messaging.incoming.my-channel-1.group.id=my-consumer-group mp.messaging.incoming.my-channel-1.kafka-configuration=my-configuration-1 # Topic 2 mp.messaging.incoming.my-channel-2.topic=my-topic-2 mp.messaging.incoming.my-channel-2.group.id=my-consumer-group mp.messaging.incoming.my-channel-2.kafka-configuration=my-configuration-2
KafkaAuthConfiguration.java (修正JAAS字符串)
import io.smallrye.common.annotation.Identifier; import jakarta.enterprise.inject.Produces; import jakarta.inject.Singleton; import lombok.extern.slf4j.Slf4j; import java.util.HashMap; import java.util.Map; @Singleton @Slf4j public class KafkaAuthConfiguration { // 修正username与password间的空格 private static final String SASL_CONFIGURATION_STRING = "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"%s\" password=\"%s\";"; @Produces @Identifier("my-configuration-1") @Singleton public Map<String, Object> config1() { final Map<String, Object> config = createConfig("clusterUrl_1", "clusterApiKey_1", "clusterApiSecret_1"); log.info("Created kafka config 1"); return config; } @Produces @Identifier("my-configuration-2") @Singleton public Map<String, Object> config2() { final Map<String, Object> config = createConfig("clusterUrl_2", "clusterApiKey_2", "clusterApiSecret_2"); log.info("Created kafka config 2"); return config; } private Map<String, Object> createConfig(final String clusterUrl, final String apiKey, final String apiSecret) { final HashMap<String, Object> config = new HashMap<>(4); config.put("bootstrap.servers", clusterUrl); config.put("security.protocol", "SASL_SSL"); config.put("sasl.mechanism", "PLAIN"); config.put("ssl.endpoint.identification.algorithm", "https"); config.put("sasl.jaas.config", String.format(SASL_CONFIGURATION_STRING, apiKey, apiSecret)); return config; } }
额外验证点
- 确认Kafka Topic存在且有未消费的消息
- 检查Broker的ACL配置,确保你的消费者组有对应Topic的
READ权限 - 开启SmallRye Kafka调试日志(添加
quarkus.log.category."io.smallrye.kafka".level=DEBUG),查看消费过程中的详细错误信息
内容的提问来源于stack exchange,提问作者AnarchoEnte
相关产品推荐
相关产品推荐

