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

Azure环境下Docker部署带客户端认证的Kafka Broker问题排查

Kafka容器持续重启问题修复方案

问题背景

在Azure Docker环境部署Kafka Broker,无认证版本运行正常,但添加SASL客户端认证配置后,Broker陷入持续重启循环。


现有配置

docker-compose.yaml

version: '3'

services:
  zookeeper:
    image: wurstmeister/zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
  kafka:
    image: wurstmeister/kafka
    container_name: kafka
    ports:
      - "9092:9092"
    environment:
      KAFKA_ADVERTISED_HOST_NAME: [ContainerIp]
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_OPTS: "-Djava.security.auth.login.config=/etc/kafka/kafka_server_jaas.conf"
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'PLAINTEXT:PLAINTEXT,SASL_PLAINTEXT:SASL_PLAINTEXT'
      KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://[ContainerIp]:9096,SASL_PLAINTEXT://[ContainerIp]:9092'
      KAFKA_LISTENERS: 'PLAINTEXT://[ContainerIp]:9096,SASL_PLAINTEXT://[ContainerIp]:9092'
      KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
      KAFKA_SASL_ENABLED_MECHANISMS: 'SASL_PLAINTEXT'
      KAFKA_SECURITY_PROTOCOL: 'SASL_PLAINTEXT'
    volumes:
      - kafkaconfig:/etc/kafka/kafka_server_jaas.conf

volumes:
  kafkaconfig:
    driver: azure_file
    driver_opts:
      share_name: [azureshare]
      storage_account_name: [shareaccount]
      storage_account_key: [accountkey]

kafka_server_jaas.conf(Azure文件共享存储)

KafkaServer {
};

KafkaClient {
    org.apache.kafka.common.security.plain.PlainLoginModule required
    username="Admin" \
    password="Password";
};

本地consumer.py

import json 
from kafka import KafkaConsumer

if __name__ == '__main__':
    # Kafka Consumer 
    consumer = KafkaConsumer(
        'pumpdatalive',
        bootstrap_servers='[ContainerIP]:9092',
        auto_offset_reset='earliest',
        sasl_mechanism='PLAIN',
        sasl_plain_username='Admin',
        sasl_plain_password='Password',
        security_protocol = "PLAINTEXT"
    )
    for message in consumer:
        print(json.loads(message.value))

本地producer.py

import time 
import json 
import random 
from datetime import datetime
from data_generator import generate_message
from kafka import KafkaProducer

# Messages will be serialized as JSON 
def serializer(message):
    return json.dumps(message).encode('utf-8')

# Kafka Producer
producer = KafkaProducer(
    bootstrap_servers=['[ContainerIP]:9092'],
    value_serializer=serializer,
    sasl_mechanism='PLAIN',
    sasl_plain_username='Admin',
    sasl_plain_password='Password',
    security_protocol = "PLAINTEXT"
)

if __name__ == '__main__':
    # Infinite loop - runs until you kill the program
    while True:
        # Generate a message
        dummy_message = generate_message()
        
        # Send it to our 'messages' topic
        print(f'Producing message @ {datetime.now()} | Message = {str(dummy_message)}')
        producer.send('pumpdatalive', dummy_message)
        
        # Sleep for a random number of seconds
        time_to_sleep = random.randint(1, 11)
        time.sleep(time_to_sleep)

问题根源与修复步骤

1. JAAS配置文件缺失核心认证规则

当前KafkaServer段为空,Kafka Broker无法加载认证模块,导致启动失败。需补充SASL PLAIN认证的用户配置:

修改后的kafka_server_jaas.conf:

KafkaServer {
    org.apache.kafka.common.security.plain.PlainLoginModule required
    username="Admin"
    password="Password"
    user_Admin="Password";
};

KafkaClient {
    org.apache.kafka.common.security.plain.PlainLoginModule required
    username="Admin"
    password="Password";
};
  • username="Admin"和password="Password":Broker间通信使用的凭据
  • user_Admin="Password":允许连接的客户端用户(Admin)及其密码,需与客户端配置一致

2. Kafka环境变量配置修正

  • KAFKA_SASL_ENABLED_MECHANISMS应设置为PLAIN(认证机制),而非SASL_PLAINTEXT(安全协议)
  • 移除全局KAFKA_SECURITY_PROTOCOL,监听器协议已通过LISTENER_SECURITY_PROTOCOL_MAP指定
  • 监听器绑定0.0.0.0,避免依赖固定容器IP,适配Azure动态网络环境

修改后的kafka服务环境变量:

environment:
  KAFKA_ADVERTISED_HOST_NAME: [ContainerIp]
  KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
  KAFKA_OPTS: "-Djava.security.auth.login.config=/etc/kafka/kafka_server_jaas.conf"
  KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'PLAINTEXT:PLAINTEXT,SASL_PLAINTEXT:SASL_PLAINTEXT'
  KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://[ContainerIp]:9096,SASL_PLAINTEXT://[ContainerIp]:9092'
  KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9096,SASL_PLAINTEXT://0.0.0.0:9092'
  KAFKA_INTER_BROKER_LISTENER_NAME: 'PLAINTEXT'
  KAFKA_SASL_ENABLED_MECHANISMS: 'PLAIN'

3. 客户端代码安全协议修正

客户端连接的是SASL认证监听器,需将security_protocol改为SASL_PLAINTEXT:

consumer.py修改部分

consumer = KafkaConsumer(
    'pumpdatalive',
    bootstrap_servers='[ContainerIP]:9092',
    auto_offset_reset='earliest',
    sasl_mechanism='PLAIN',
    sasl_plain_username='Admin',
    sasl_plain_password='Password',
    security_protocol = "SASL_PLAINTEXT"
)

producer.py修改部分

producer = KafkaProducer(
    bootstrap_servers=['[ContainerIP]:9092'],
    value_serializer=serializer,
    sasl_mechanism='PLAIN',
    sasl_plain_username='Admin',
    sasl_plain_password='Password',
    security_protocol = "SASL_PLAINTEXT"
)

4. 验证挂载文件有效性

进入容器确认JAAS配置文件已正确加载:

docker exec -it kafka cat /etc/kafka/kafka_server_jaas.conf

验证流程

  1. 重启Kafka服务:
    docker-compose down && docker-compose up -d
    
  2. 查看容器日志确认启动成功:
    docker logs kafka
    
  3. 运行本地producer和consumer,验证消息收发正常。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 16:34:50