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

Docker容器中运行Confluent Kafka时无法创建Topic求助

问题:Docker中运行Confluent Kafka创建Topic失败,出现rdkafka日志Purging 1 unserved events from background queue

环境配置

Docker Compose配置(Kafka部分)

kafka:
  image: confluentinc/cp-kafka:latest
  networks:
    - kafka_network
  depends_on:
    - zookeeper
  ports:
    - 29092:29092
    - 29093:29093
  environment:
    KAFKA_BROKER_ID: 1
    KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
    KAFKA_LISTENERS: EXTERNAL_SAME_HOST://:29092,EXTERNAL_DIFFERENT_HOST://:29093,INTERNAL://:9092
    KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka:9092,EXTERNAL_SAME_HOST://localhost:29092,EXTERNAL_DIFFERENT_HOST://172.23.0.3:29093
    KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL_SAME_HOST:PLAINTEXT,EXTERNAL_DIFFERENT_HOST:PLAINTEXT
    KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
    KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1

Python创建Topic代码

def create_kafka_topic(bootstrap_servers, topic_name, num_partitions=1, replication_factor=1):
    # Configure the admin client with the bootstrap servers
    admin_client = AdminClient({'bootstrap.servers': bootstrap_servers})

    # Create a NewTopic object with topic configuration
    new_topic = NewTopic(topic_name, num_partitions=num_partitions, replication_factor=replication_factor)

    # Create the topic
    admin_client.create_topics([new_topic])

if __name__ == "__main__":
    # Configure your Kafka broker address here
    bootstrap_servers = "localhost:29092"

    # Topic details
    topic_name = config.TOPIC
    num_partitions = config.NUM_PARTITIONS
    replication_factor = 1
    
    # Create the Kafka topic for transactions
    create_kafka_topic(bootstrap_servers, topic_name_transactions, num_partitions, replication_factor)

错误日志

%6|1691170181.360|BGQUEUE|rdkafka#producer-1| [thrd:background]: Purging 1 unserved events from background queue

排查与解决方案

1. 修复Python代码变量名错误

代码中定义了topic_name = config.TOPIC,但调用创建函数时使用了未定义的topic_name_transactions,会直接触发NameError导致程序提前终止,rdkafka的日志只是程序中断后的资源清理输出。

修正后主函数部分:

if __name__ == "__main__":
    bootstrap_servers = "localhost:29092"

    topic_name = config.TOPIC
    num_partitions = config.NUM_PARTITIONS
    replication_factor = 1
    
    # 使用正确的变量名topic_name
    create_kafka_topic(bootstrap_servers, topic_name, num_partitions, replication_factor)

2. 等待AdminClient操作完成并处理结果

admin_client.create_topics()返回Future对象,必须等待操作完成并检查结果,否则程序可能在Kafka响应前就退出,导致Topic创建操作未执行成功。

修改后的创建函数:

def create_kafka_topic(bootstrap_servers, topic_name, num_partitions=1, replication_factor=1):
    admin_client = AdminClient({'bootstrap.servers': bootstrap_servers})
    new_topic = NewTopic(topic_name, num_partitions=num_partitions, replication_factor=replication_factor)
    
    # 发起请求并等待结果
    futures = admin_client.create_topics([new_topic])
    for topic, future in futures.items():
        try:
            future.result()
            print(f"Topic {topic} created successfully")
        except Exception as e:
            print(f"Failed to create topic {topic}: {e}")

3. 验证Kafka服务可用性

  • 检查容器状态:执行docker-compose ps,确认ZooKeeper和Kafka容器均处于Up状态
  • 进入Kafka容器,用官方工具测试Topic创建:
docker-compose exec kafka kafka-topics --create --topic test-topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1

如果该命令能成功创建Topic,说明Kafka服务本身无问题,问题集中在Python代码层面。

4. 检查本地与Kafka的网络连通性

确认本地能访问29092端口:

nc -zv localhost 29092

若连接失败,检查Docker端口映射是否生效,或本地防火墙是否拦截了该端口。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 03:47:09