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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 01:04:59