Docker中Kafka kRaft模式SASL认证:外部Python脚本无法收发消息
Kafka kRaft集群外部Python脚本连接失败排查
问题现象
Docker中运行kRaft模式Kafka集群,容器内命令行可正常生产消费消息,但外部Python脚本连接时报错:
%6|1715163671.929|FAIL|rdkafka#producer-1| [thrd:sasl_plaintext://localhost:9094/bootstrap]: sasl_plaintext://localhost:9094/bootstrap: Disconnected while requesting ApiVersion: might be caused by incorrect security.protocol configuration (connecting to a SSL listener?) or broker version is < 0.10 (see api.version.request) (after 4ms in state APIVERSION_QUERY)
用户配置文件
docker-compose.yml
version: "3.8" services: kafka1: image: confluentinc/cp-kafka:latest hostname: kafka1 container_name: kafka1 ports: - "9092:9092" environment: KAFKA_LISTENERS: INTERNAL://kafka1:19092,EXTERNAL://0.0.0.0:9092,CONTROLLER://kafka1:29093 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka1:19092,EXTERNAL://localhost:9092 KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:SASL_PLAINTEXT,EXTERNAL:SASL_PLAINTEXT KAFKA_SASL_MECHANISM_INTER_BROKER_PROTOCOL: PLAIN KAFKA_SASL_ENABLED_MECHANISMS: PLAIN KAFKA_SASL_MECHANISM_CONTROLLER_PROTOCOL: PLAIN KAFKA_SUPER_USERS: "User:admin" KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 KAFKA_PROCESS_ROLES: 'controller,broker' KAFKA_NODE_ID: 1 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true' KAFKA_ALLOW_EVERYONE_IF_NO_ACL_FOUND: 'true' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_METADATA_LOG_SEGMENT_MS: 15000 KAFKA_METADATA_MAX_RETENTION_MS: 1200000 KAFKA_METADATA_LOG_MAX_RECORD_BYTES_BETWEEN_SNAPSHOTS: 2800 KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' KAFKA_OPTS: '-Djava.security.auth.login.config=/etc/kafka/kafka_server_jaas.conf' CLUSTER_ID: '<ClusterID>' volumes: - ./kafka_server_jaas.conf:/etc/kafka/kafka_server_jaas.conf - kafka1-data:/var/lib/kafka/data - ./clusterID:/tmp/clusterID kafka2: image: confluentinc/cp-kafka:latest hostname: kafka2 container_name: kafka2 ports: - "9093:9093" environment: KAFKA_LISTENERS: INTERNAL://kafka2:19093,EXTERNAL://0.0.0.0:9093,CONTROLLER://kafka2:29093 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka2:19093,EXTERNAL://localhost:9093 KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:SASL_PLAINTEXT,EXTERNAL:SASL_PLAINTEXT KAFKA_SASL_MECHANISM_INTER_BROKER_PROTOCOL: PLAIN KAFKA_SASL_ENABLED_MECHANISMS: PLAIN KAFKA_SASL_MECHANISM_CONTROLLER_PROTOCOL: PLAIN KAFKA_SUPER_USERS: "User:admin" KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 KAFKA_PROCESS_ROLES: 'controller,broker' KAFKA_NODE_ID: 2 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" KAFKA_ALLOW_EVERYONE_IF_NO_ACL_FOUND: 'true' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_METADATA_LOG_SEGMENT_MS: 15000 KAFKA_METADATA_MAX_RETENTION_MS: 1200000 KAFKA_METADATA_LOG_MAX_RECORD_BYTES_BETWEEN_SNAPSHOTS: 2800 KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' KAFKA_OPTS: "-Djava.security.auth.login.config=/etc/kafka/kafka_server_jaas.conf" CLUSTER_ID: '<ClusterID>' volumes: - ./kafka_server_jaas.conf:/etc/kafka/kafka_server_jaas.conf - kafka2-data:/var/lib/kafka/data - ./clusterID:/tmp/clusterID kafka3: image: confluentinc/cp-kafka:latest hostname: kafka3 container_name: kafka3 ports: - "9094:9094" environment: KAFKA_LISTENERS: INTERNAL://kafka3:19094,EXTERNAL://0.0.0.0:9094,CONTROLLER://kafka3:29093 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka3:19094,EXTERNAL://localhost:9094 KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:SASL_PLAINTEXT,EXTERNAL:SASL_PLAINTEXT KAFKA_SASL_MECHANISM_INTER_BROKER_PROTOCOL: PLAIN KAFKA_SASL_ENABLED_MECHANISMS: PLAIN KAFKA_SASL_MECHANISM_CONTROLLER_PROTOCOL: PLAIN KAFKA_SUPER_USERS: "User:admin" KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 KAFKA_PROCESS_ROLES: 'controller,broker' KAFKA_NODE_ID: 3 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true" KAFKA_ALLOW_EVERYONE_IF_NO_ACL_FOUND: 'true' KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_METADATA_LOG_SEGMENT_MS: 15000 KAFKA_METADATA_MAX_RETENTION_MS: 1200000 KAFKA_METADATA_LOG_MAX_RECORD_BYTES_BETWEEN_SNAPSHOTS: 2800 KAFKA_LOG_DIRS: '/tmp/kraft-combined-logs' KAFKA_OPTS: "-Djava.security.auth.login.config=/etc/kafka/kafka_server_jaas.conf" CLUSTER_ID: '<ClusterID>' volumes: - ./kafka_server_jaas.conf:/etc/kafka/kafka_server_jaas.conf - kafka3-data:/var/lib/kafka/data - ./clusterID:/tmp/clusterID volumes: kafka1-data: kafka2-data: kafka3-data:
kafka_server_jaas.conf
KafkaServer { org.apache.kafka.common.security.plain.PlainLoginModule required username="admin" password="secret-admin" user_usera="secret-usera"; };
修正后的Python脚本
from confluent_kafka import Producer def delivery_report(err, msg): """Delivery report callback.""" if err is not None: print(f'Message delivery failed: {err}') else: print(f'Message delivered to {msg.topic()} [{msg.partition()}]') def main(): conf = { 'bootstrap.servers': 'localhost:9092,localhost:9093,localhost:9094', 'security.protocol': 'sasl_plaintext', 'sasl.mechanism': 'PLAIN', 'sasl.username': 'usera', 'sasl.password': 'secret-usera', } producer = Producer(conf) # 生产测试消息验证连接 try: producer.produce('test_topic', key='test_key', value='test_value', callback=delivery_report) producer.flush() # 等待消息投递完成 print("消息发送成功,Kafka连接正常。") except Exception as e: print(f"错误:{e}") # 关闭生产者 producer.flush() if __name__ == '__main__': main()
错误排查与修复步骤
1. 核心错误:密码不匹配
Python脚本中sasl.password原配置为secret-users,但JAAS文件中定义的usera密码是secret-usera,认证信息不匹配导致连接失败,这是主要问题。
2. Docker配置优化
原docker-compose中KAFKA_LISTENERS存在换行和多余逗号,可能导致Kafka解析配置出错。需将配置合并为一行,并将EXTERNAL监听地址改为0.0.0.0:端口(容器内监听localhost会导致外部无法通过端口映射访问)。
3. 修复操作流程
- 修正Python脚本中的密码为
secret-usera - 更新docker-compose.yml中的环境变量格式,确保监听地址正确
- 重启Kafka集群:
docker-compose down -v docker-compose up -d - 运行修正后的Python脚本测试连接
内容的提问来源于stack exchange,提问作者dvasilevich
相关产品推荐
相关产品推荐

