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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 13:37:03