Docker中KafkaAdminClient创建失败及Python消息发送报错求助
问题描述
我希望在Docker容器内创建Kafka Topic,并通过容器外的Python代码向该Topic发送消息,操作步骤如下:
- 搭建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
- 编写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
相关产品推荐
相关产品推荐

