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
验证流程
- 重启Kafka服务:
docker-compose down && docker-compose up -d - 查看容器日志确认启动成功:
docker logs kafka - 运行本地producer和consumer,验证消息收发正常。
内容的提问来源于stack exchange,提问作者GameFire
相关产品推荐
相关产品推荐

