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

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

原因分析

  1. 启动时机不匹配:Bytewax容器的depends_on仅依赖Kafka容器启动,未等待Kafka完成健康检查。而正常工作的consumer容器明确指定等待Kafka处于service_healthy状态后才启动,导致Bytewax在Kafka尚未完全就绪时发起连接请求,触发连接拒绝。
  2. 主题元数据查询失败: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 21:28:13