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

Docker中KafkaAdminClient创建失败及Python消息发送报错求助

问题描述

我希望在Docker容器内创建Kafka Topic,并通过容器外的Python代码向该Topic发送消息,操作步骤如下:

  1. 搭建Docker下的Kafka基础设施,docker-compose.yml内容:
version: '2'
services:
  zookeeper:
    image: "confluentinc/cp-zookeeper:5.2.1"
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000

  kafka0:
    image: "confluentinc/cp-enterprise-kafka:5.2.1"
    ports:
      - '9092:9092'
      - '29094:29094'
    depends_on:
      - zookeeper
    environment:
      KAFKA_BROKER_ID: 0
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_LISTENERS: LISTENER_BOB://kafka0:29092,LISTENER_FRED://kafka0:9092,LISTENER_ALICE://kafka0:29094
      KAFKA_ADVERTISED_LISTENERS: LISTENER_BOB://kafka0:29092,LISTENER_FRED://localhost:9092,LISTENER_ALICE://never-gonna-give-you-up:29094
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: LISTENER_BOB:PLAINTEXT,LISTENER_FRED:PLAINTEXT,LISTENER_ALICE:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: LISTENER_BOB
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false"
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 100
  1. 编写Python发送消息脚本:
import asyncio

from aiokafka import AIOKafkaProducer


async def send_to_kafka():
    producer = AIOKafkaProducer(
        bootstrap_servers='localhost:9092',
        enable_idempotence=True)
    await producer.start()
    try:
        await producer.send_and_wait("my-topic-1", b"Super message")
    finally:
        await producer.stop()

asyncio.run(send_to_kafka())

执行脚本时出现错误:

Topic my-topic-1 not found in cluster metadata

尝试在Docker容器内创建Topic,执行命令:

docker exec -ti fetcher-service-kafka0-1 bash
kafka-topics --bootstrap-server kafka0:29092 --create --if-not-exists --topic my-topic-1 --replication-factor 1 --partitions 1

出现DNS解析错误:

[2024-11-27 12:11:11,237] WARN Couldn't resolve server kafka0:29092 from bootstrap.servers as DNS resolution failed for kafka0 (org.apache.kafka.clients.ClientUtils)
fetcher-service-init-kafka-1  | Exception in thread "main" org.apache.kafka.common.KafkaException: Failed to create new KafkaAdminClient
fetcher-service-init-kafka-1  |         at org.apache.kafka.clients.admin.KafkaAdminClient.createInternal(KafkaAdminClient.java:386)
fetcher-service-init-kafka-1  |         at org.apache.kafka.clients.admin.AdminClient.create(AdminClient.java:55)
fetcher-service-init-kafka-1  |         at kafka.admin.TopicCommand$AdminClientTopicService$.createAdminClient(TopicCommand.scala:150)
fetcher-service-init-kafka-1  |         at kafka.admin.TopicCommand$AdminClientTopicService$.apply(TopicCommand.scala:154)
fetcher-service-init-kafka-1  |         at kafka.admin.TopicCommand$.main(TopicCommand.scala:55)
fetcher-service-init-kafka-1  |         at kafka.admin.TopicCommand.main(TopicCommand.scala)
fetcher-service-init-kafka-1  | Caused by: org.apache.kafka.common.config.ConfigException: No resolvable bootstrap urls given in bootstrap.servers
fetcher-service-init-kafka-1  |         at org.apache.kafka.clients.ClientUtils.parseAndValidateAddresses(ClientUtils.java:90)
fetcher-service-init-kafka-1  |         at org.apache.kafka.clients.ClientUtils.parseAndValidateAddresses(ClientUtils.java:49)
fetcher-service-init-kafka-1  |         at org.apache.kafka.clients.admin.KafkaAdminClient.createInternal(KafkaAdminClient.java:346)
解决方案

1. 修复容器内创建Topic的DNS解析问题

容器内执行kafka-topics命令时,kafka0主机名无法被解析,直接使用容器内部的监听地址localhost:29092即可解决:

kafka-topics --bootstrap-server localhost:29092 --create --if-not-exists --topic my-topic-1 --replication-factor 1 --partitions 1

如果希望保留kafka0作为主机名,可在docker-compose.yml的kafka0服务中添加extra_hosts配置,手动映射主机名到本地回环地址:

kafka0:
  # 保留原有配置不变
  extra_hosts:
    - "kafka0:127.0.0.1"

修改后重启Kafka容器,再执行创建命令即可。

2. 确保Python脚本正常发送消息

Topic创建完成后,Python脚本的配置(使用localhost:9092对应外部暴露的LISTENER_FRED)是正确的,注意以下两点:

  • 等待Kafka容器完全启动后再执行脚本,避免元数据未同步完成
  • 若仍报错,可在AIOKafkaProducer中添加metadata_max_age_ms参数,强制定期刷新元数据:
producer = AIOKafkaProducer(
    bootstrap_servers='localhost:9092',
    enable_idempotence=True,
    metadata_max_age_ms=5000  # 每5秒刷新一次元数据
)

3. 可选:开启自动创建Topic(不推荐生产环境)

若不想手动创建Topic,可将docker-compose.yml中的KAFKA_AUTO_CREATE_TOPICS_ENABLE改为"true",这样Python脚本发送消息时Kafka会自动创建my-topic-1。但生产环境建议手动管理Topic,避免意外生成不符合规范的Topic。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 23:25:56