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
相关产品推荐
相关产品推荐

