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

Kafka Broker无法接收Python Producer消息的排查求助

Kafka消费者无法接收消息的排查方案

问题背景

使用Docker Compose部署Kafka环境,生产者能发送消息但消费者收不到,相关配置和代码如下:

Docker Compose配置(compose_kafka.yml)

version: '3'

services:
  zookeeper:
    image: wurstmeister/zookeeper
    container_name: zookeeper
    ports:
      - "2181:2181"
    environment:
      ZOO_MY_ID: 1
  kafka:
    image: wurstmeister/kafka
    container_name: kafka
    ports:
      - "9092:9092"
    environment:
      KAFKA_ADVERTISED_HOST_NAME: 192.168.1.10
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
  kafka_manager:
    image: kafkamanager/kafka-manager
    container_name: kafka-manager
    restart: always
    ports:
      - "9000:9000"
    environment:
      ZK_HOSTS: "192.168.1.10:2181"
      APPLICATION_SECRET: "random-secret"

生产者相关代码

Generate.py

from faker import Faker

fake = Faker()

class Registered_user:

    def get_registered_user():
        return {
            "name": fake.name(),
            "address": fake.address(),
            "created_at": fake.year()
        }

Producer_registered_user.py

import time
import json
from kafka import KafkaProducer
from fake_data import Generate

  
def json_serializer(data):
    return json.dumps(data).encode("utf-8")


producer = KafkaProducer(bootstrap_servers='192.168.1.10:9092',
                         value_serializer=json_serializer)


if __name__ == '__main__':
    while 1 == 1:
        user = Generate.Registered_user.get_registered_user()
        producer.send('registered_user', user)
        print(user)
        time.sleep(4)

消费者代码(Consumer_registered_user.py)

import json
from kafka import KafkaConsumer

if __name__ == '__main__':
    consumer = KafkaConsumer(
        bootstrap_servers='192.168.1.10:9092',
        auto_offset_reset="from-beginning",
        group_id="consumer-group-a"
    )

    for message in consumer:
        print("User = {}".format(json.loads(message.value)))

排查步骤

1. 消费者未订阅目标Topic

当前消费者代码未指定要订阅的registered_user Topic,KafkaConsumer初始化后需显式订阅才能接收消息。修改消费者代码:

import json
from kafka import KafkaConsumer

if __name__ == '__main__':
    consumer = KafkaConsumer(
        'registered_user',  # 指定订阅的Topic
        bootstrap_servers='192.168.1.10:9092',
        auto_offset_reset="from-beginning",
        group_id="consumer-group-a"
    )

    for message in consumer:
        print("User = {}".format(json.loads(message.value)))

也可在初始化后调用consumer.subscribe(['registered_user'])完成订阅。

2. 生产者消息未成功提交

send方法为异步操作,可能消息未完成发送就进入下一次循环。可添加回调并强制刷新确认发送状态:

import time
import json
from kafka import KafkaProducer
from fake_data import Generate

def json_serializer(data):
    return json.dumps(data).encode("utf-8")

def on_send_success(record_metadata):
    print(f"消息已发送到Topic: {record_metadata.topic}, Partition: {record_metadata.partition}, Offset: {record_metadata.offset}")

def on_send_error(excp):
    print(f"消息发送失败: {excp}")

producer = KafkaProducer(bootstrap_servers='192.168.1.10:9092',
                         value_serializer=json_serializer)

if __name__ == '__main__':
    while 1 == 1:
        user = Generate.Registered_user.get_registered_user()
        producer.send('registered_user', user).add_callback(on_send_success).add_errback(on_send_error)
        producer.flush()  # 确保消息提交到Kafka
        print(user)
        time.sleep(4)

3. Kafka Topic配置异常

检查Topic的副本同步状态,进入Kafka容器执行命令:

docker exec -it kafka /opt/kafka/bin/kafka-topics.sh --describe --topic registered_user --bootstrap-server localhost:9092

若ReplicationFactor为1但副本不在ISR(同步副本)列表中,消息无法持久化。可在Docker Compose的Kafka环境变量中添加:

environment:
  KAFKA_ADVERTISED_HOST_NAME: 192.168.1.10
  KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
  KAFKA_DEFAULT_REPLICATION_FACTOR: 1
  KAFKA_MIN_IN_SYNC_REPLICAS: 1

4. 消费者反序列化配置缺失

生产者用JSON序列化消息,但消费者未指定反序列化规则,手动解码可能出现格式问题。可显式配置反序列化器:

import json
from kafka import KafkaConsumer

if __name__ == '__main__':
    consumer = KafkaConsumer(
        'registered_user',
        bootstrap_servers='192.168.1.10:9092',
        auto_offset_reset="from-beginning",
        group_id="consumer-group-a",
        value_deserializer=lambda m: json.loads(m.decode('utf-8'))
    )

    for message in consumer:
        print("User = {}".format(message.value))

5. 网络连通性验证

确认宿主机与Kafka容器的网络连通:

telnet 192.168.1.10 9092

若不通,检查Docker端口映射是否正常,或容器内9092端口是否监听:

docker exec -it kafka netstat -tulpn | grep 9092

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 22:07:33