Docker环境下Kafka与Bytewax连接被拒绝问题求助
Bytewax连接Kafka失败但Python KafkaConsumer正常的问题排查
问题现象
使用Bytewax消费Kafka流执行聚合时,持续出现连接被拒绝的错误,日志同时提示"Unknown topic or partition"。但运行在另一个容器中的Python KafkaConsumer 可以正常消费并打印数据。
相关配置与代码
docker-compose.yml
services: kafka: image: apache/kafka ports: - "9092:9092" environment: # Configure listeners for both docker and host communication KAFKA_LISTENERS: CONTROLLER://localhost:9091,HOST://0.0.0.0:9092,DOCKER://kafka:9093 KAFKA_ADVERTISED_LISTENERS: HOST://localhost:9092,DOCKER://kafka:9093 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,DOCKER:PLAINTEXT,HOST:PLAINTEXT # Settings required for KRaft mode KAFKA_NODE_ID: 1 KAFKA_PROCESS_ROLES: broker,controller KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9091 # Listener to use for broker-to-broker communication KAFKA_INTER_BROKER_LISTENER_NAME: DOCKER # Required for a single node cluster KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 healthcheck: test: ["CMD", "bash", "-c", "/opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka:9092 --list"] interval: 10s timeout: 5s retries: 5 networks: - app-network kafka-ui: image: ghcr.io/kafbat/kafka-ui:latest ports: - 8080:8080 environment: DYNAMIC_CONFIG_ENABLED: "true" KAFKA_CLUSTERS_0_NAME: local KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9093 depends_on: - kafka networks: - app-network consumer: build: context: ./kafka_consumer dockerfile: Dockerfile container_name: consumer depends_on: factory-service: condition: service_started kafka: condition: service_healthy ports: - "8099:80" networks: - app-network bytewax: build: context: ./consumer dockerfile: Dockerfile container_name: bytewax depends_on: - kafka networks: - app-network networks: app-network: driver: bridge
正常工作的consumer.py
from kafka import KafkaConsumer import json KAFKA_BROKER = "kafka:9093" KAFKA_TOPIC = ["factory_001","factory_002"] consumer = KafkaConsumer( *KAFKA_TOPIC, group_id='my-group', bootstrap_servers=KAFKA_BROKER, value_deserializer=lambda x: json.loads(x.decode("utf-8")) )
无法工作的Stream_process.py
from bytewax import operators as op from bytewax.connectors.kafka import KafkaSource from bytewax.connectors.stdio import StdOutSink from bytewax.dataflow import Dataflow KAFKA_BROKER = ["kafka:9093"] KAFKA_TOPIC = ["factory_001"] flow = Dataflow("Average Aggregation") stream = op.input("kafka-in", flow, KafkaSource(KAFKA_BROKER, KAFKA_TOPIC)) op.output("out", stream, StdOutSink())
错误日志
%3|1742203770.748|FAIL|rdkafka#producer-1| [thrd:kafka:9093/bootstrap]: kafka:9093/bootstrap: Connect to ipv4#172.19.0.2:9093 failed: Connection refused (after 1ms in state CONNECT) %3|1742203771.749|FAIL|rdkafka#producer-1| [thrd:kafka:9093/bootstrap]: kafka:9093/bootstrap: Connect to ipv4#172.19.0.2:9093 failed: Connection refused (after 0ms in state CONNECT, 1 identical error(s) suppressed) thread '<unnamed>' panicked at src/run.rs:128:17: Box<dyn Any> note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/bytewax/connectors/kafka/__init__.py", line 387, in list_parts return list(_list_parts(client, self._topics)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ %3|1742203770.748|FAIL|rdkafka#producer-1| [thrd:kafka:9093/bootstrap]: kafka:9093/bootstrap: Connect to ipv4#172.19.0.2:9093 failed: Connection refused (after 1ms in state CONNECT) %3|1742203771.749|FAIL|rdkafka#producer-1| [thrd:kafka:9093/bootstrap]: kafka:9093/bootstrap: Connect to ipv4#172.19.0.2:9093 failed: Connection refused (after 0ms in state CONNECT, 1 identical error(s) suppressed) thread '<unnamed>' panicked at src/run.rs:128:17: Box<dyn Any> note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/bytewax/connectors/kafka/__init__.py", line 387, in list_parts return list(_list_parts(client, self._topics)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ ONNECT, 1 identical error(s) suppressed) thread '<unnamed>' panicked at src/run.rs:128:17: Box<dyn Any> note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/bytewax/connectors/kafka/__init__.py", line 387, in list_parts return list(_list_parts(client, self._topics)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ Box<dyn Any> note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/bytewax/connectors/kafka/__init__.py", line 387, in list_parts return list(_list_parts(client, self._topics)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/bytewax/connectors/kafka/__init__.py", line 387, in list_parts return list(_list_parts(client, self._topics)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ return list(_list_parts(client, self._topics)) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/bytewax/connectors/kafka/__init__.py", line 172, in _list_parts raise RuntimeError(msg) RuntimeError: error listing partitions for Kafka topic `'factory_001'`: Broker: Unknown topic or partition The above exception was the direct cause of the following exception: bytewax.errors.BytewaxRuntimeError: (src/inputs.rs:252:47): error calling `FixedPartitionSource.list_parts` in step "Average Aggregation.kafka-in" The above exception was the direct cause of the following exception: bytewax.errors.BytewaxRuntimeError: (src/worker.rs:354:34): error building FixedPartitionedSource The above exception was the direct cause of the following exception: Traceback (most recent call last): File "<frozen runpy>", line 198, in _run_module_as_main File "<frozen runpy>", line 88, in _run_code File "/usr/local/lib/python3.11/site-packages/bytewax/run.py", line 355, in <module> cli_main(**kwargs) bytewax.errors.BytewaxRuntimeError: (src/worker.rs:149:10): error building production dataflow
原因分析
- 启动时机不匹配:Bytewax容器的
depends_on仅依赖Kafka容器启动,未等待Kafka完成健康检查。而正常工作的consumer容器明确指定等待Kafka处于service_healthy状态后才启动,导致Bytewax在Kafka尚未完全就绪时发起连接请求,触发连接拒绝。 - 主题元数据查询失败:Kafka未完全初始化时,Bytewax尝试查询主题分区信息,返回"Unknown topic or partition",进一步引发运行时异常。
解决方案
1. 调整Bytewax容器的启动依赖
修改docker-compose.yml中bytewax服务的depends_on配置,使其等待Kafka健康就绪:
bytewax: build: context: ./consumer dockerfile: Dockerfile container_name: bytewax depends_on: kafka: condition: service_healthy networks: - app-network
2. 增强Bytewax KafkaSource的容错配置
Bytewax底层依赖librdkafka,可添加额外配置参数提升容错性:
from bytewax import operators as op from bytewax.connectors.kafka import KafkaSource, KafkaConfig from bytewax.connectors.stdio import StdOutSink from bytewax.dataflow import Dataflow KAFKA_BROKER = ["kafka:9093"] KAFKA_TOPIC = ["factory_001"] # 添加Kafka配置参数 kafka_config = KafkaConfig( bootstrap_servers=KAFKA_BROKER, auto_offset_reset="earliest", # 自动重置偏移量到最早位置 retry_backoff_ms=1000, # 连接重试间隔 message_max_bytes=10485760 # 根据业务调整最大消息大小 ) flow = Dataflow("Average Aggregation") stream = op.input("kafka-in", flow, KafkaSource(kafka_config, KAFKA_TOPIC)) op.output("out", stream, StdOutSink())
3. 验证并创建Kafka主题
确保factory_001主题已存在:
# 进入Kafka容器 docker exec -it <kafka-container-id> bash # 列出所有主题 /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list
若主题不存在,手动创建:
/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic factory_001 --partitions 1 --replication-factor 1
内容的提问来源于stack exchange,提问作者LeXXan
相关产品推荐
相关产品推荐

