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

Docker Compose环境下confluent-kafka连接失败排查求助

Docker Compose环境下Kafka消费者连接被拒绝问题排查

问题详情

在Docker Compose环境中搭建Kafka消费者时出现连接拒绝错误,错误日志如下:

notification-service-notification-service-1 | %3|1693292196.192|FAIL|rdkafka#consumer-1| [thrd:sasl_plaintext://kafka:9092/bootstrap]: sasl_plaintext://kafka:9092/bootstrap: Connect to ipv4#172.28.0.4:9092 failed: Connection refused (after 8ms in state CONNECT)

环境配置

  • 使用Docker Compose管理Kafka、ZooKeeper服务
  • 通过confluent-kafka库实现消费者
  • 环境变量配置Kafka连接信息:
KAFKA_SERVERS=kafka:9092
KAFKA_SECURITY_PROTOCOL=SASL_PLAINTEXT
KAFKA_SASL_MECHANISMS=PLAIN

消费者代码

class KafkaConsumer:
    def __init__(self, topic):
        config = get_config()

        self.consumer_config = {
            "bootstrap.servers": config.kafka_servers,
            "group.id": config.kafka_consumer_group_id,
            "security.protocol": config.kafka_security_protocol,
            "sasl.mechanism": config.kafka_sasl_mechanisms,
            "sasl.username": config.kafka_username,
            "sasl.password": config.kafka_password,
            "auto.offset.reset": "earliest",
        }

        self.consumer = Consumer(self.consumer_config)
        self.topic = topic

    def deserializer(self, serialized):
        try:
            json_data = json.loads(serialized)
            notification = Notification(**json_data)
        except Exception as e:
            logger.error(e)
            return

        return notification

    async def start(self):
        self.consumer.subscribe([self.topic])
        try:
            while True:
                msg = self.consumer.poll(1.0)
                if msg is None:
                    continue
                if msg.error():
                    if msg.error().code() == KafkaError._PARTITION_EOF:
                        continue
                    else:
                        logger.error(msg.error())
                else:
                    notification = self.deserializer(msg.value())
                    if notification:
                        try:
                            trigger_push_notification(notification)
                        except Exception as e:
                            logger.error(e)
        finally:
            self.consumer.close()

排查步骤

1. 确认Kafka容器的运行状态与端口监听

  • 执行docker-compose ps查看kafka容器状态,确保容器处于Up状态,没有频繁重启
  • 进入kafka容器内部:docker-compose exec kafka bash,然后执行nc -zv localhost 9092,确认broker在容器内正常监听9092端口
  • 重点检查Kafka的KAFKA_ADVERTISED_LISTENERS配置,这是Docker环境下Kafka连接的核心配置,必须设置为容器可被访问的地址(比如SASL_PLAINTEXT://kafka:9092),否则消费者会拿到错误的broker地址导致连接失败

2. 验证SASL认证配置一致性

  • 确认Kafka broker端已开启SASL认证:检查kafka容器的配置(通常是server.properties或环境变量),确保SASL_ENABLED_MECHANISMS=PLAIN、SASL_MECHANISM_INTER_BROKER_PROTOCOL=PLAIN
  • 检查broker的JAAS配置文件,确认其中的用户名、密码与消费者环境变量中的KAFKA_USERNAME、KAFKA_PASSWORD完全一致
  • 如果broker未开启SASL认证,消费者端的security.protocol不能设为SASL_PLAINTEXT,否则会导致连接失败

3. 检查Docker网络连通性

  • 确认消费者服务与Kafka服务在同一个Docker网络中:查看docker-compose.yml,两者是否配置了相同的networks(默认情况下同一compose文件的服务会在同一网络)
  • 进入消费者容器:docker-compose exec notification-service bash,执行ping kafka确认主机名能解析到日志中的IP(172.28.0.4),再执行nc -zv kafka 9092测试端口是否能连通

4. 验证代码与环境变量加载

  • 在消费者初始化代码中添加日志,打印self.consumer_config的完整内容,确认bootstrap.servers、security.protocol等参数与环境变量配置一致
  • 检查get_config()函数是否正确读取环境变量,比如是否存在大小写错误、变量名映射错误(比如把KAFKA_SERVERS映射成了其他字段)

内容的提问来源于stack exchange,提问作者cole

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 05:14:50