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

使用Kafka实现两Python程序双向消息传递的代码故障排查

Kafka双向通信异常问题排查

问题背景

计算机视觉项目中需要实现两端双向消息流转:图像数据从webots-controller发送到AI-model作为推理输入,AI-model输出的运动指令再回传给webots-controller执行。方案采用两个Kafka Topic分别作为两个方向的传输通道,但测试代码运行时出现消息丢失、接收阻塞的异常。

原始代码实现

AI-model端代码

# AI-model.py
from kafka import KafkaConsumer, KafkaProducer

producer = KafkaProducer(bootstrap_servers='localhost:9092')
consumer = KafkaConsumer('model-mailbox')

while(True):
    img = consumer.__next__()
    print(img.key)
    print('a-received')

    producer.send('webots-mailbox', key=b'movement', value=b'a')
    producer.flush()
    print('a-sent')

webots-controller端代码

# webots-controller.py
from kafka import KafkaConsumer, KafkaProducer

producer = KafkaProducer(bootstrap_servers='localhost:9092')
consumer = KafkaConsumer('webots-mailbox')

while True:
    producer.send('model-mailbox', key=b'image', value=b'b')
    producer.flush()
    print('b-sent')

    movement = consumer.__next__()
    print(movement.key)
    print('b-received')

异常现象

优先启动AI-model程序后,两端运行状态如下:

  • AI-model端可正常收到webots-controller发送的图像消息,完成消息发送后阻塞等待下一条输入,控制台输出:
matin@matin:~/ python AI-model.py
b'image'
a-received
a-sent
  • webots-controller端发送完第一条图像消息后,一直阻塞等待运动指令消息,无法收到AI-model返回的结果,控制台仅输出:
matin@matin:~/ python webots-controller.py
b-sent

补充测试结果

注释掉AI-model代码中a-received对应的两行打印逻辑后,消息可以正常被webots-controller接收,收发流程跑通,此时webots-controller控制台输出为:

matin@matin:~/ python webots-controller.py
b-sent
b'movement'
b-received

问题根因

异常本质是KafkaConsumer异步初始化带来的时序竞态,和打印逻辑本身无关,打印操作只是改变了代码执行的时间差,刚好触发/避开了竞态窗口:

  • kafka-python的KafkaConsumer实例创建是异步过程:执行KafkaConsumer()代码返回时,后台线程还在执行集群连接、消费组加入、分区分配、消费位点初始化的流程,并没有进入可正常消费的状态。
  • 原始代码中消费者没有显式配置参数:默认auto_offset_reset='latest'(仅消费消费者完成位点初始化之后产生的新消息),且自动生成随机消费组ID,每次启动都是全新消费组,位点初始化的时机完全不受控。
  • 时序逻辑冲突:webots-controller启动后发送完第一条消息,立刻阻塞等待返回结果;如果AI-model端因为打印操作产生延迟,发送返回消息的时间点早于webots-controller端消费者完成位点初始化的时间点,webots端的消费者会把消费起始位点设置在这条返回消息之后,永远无法收到该消息,进入永久阻塞。去掉打印后AI-model返回消息的速度更快,消息到达时webots端消费者还没完成位点初始化,拉取时就能读到这条消息,流程暂时跑通,但只要执行时序出现波动就会再次复现阻塞问题。

修复方案

从消除初始化竞态、明确消费配置两个角度修改代码:

  • 给两个消费者显式指定固定的group_id,根据业务需求配置auto_offset_reset策略(需要消费历史消息设为earliest,仅消费启动后的新消息可保留latest)。
  • 消费者创建完成后,主动调用一次短超时的poll()方法,等待消费者完成分区分配、位点初始化后再进入业务收发循环。
  • 替换直接调用__next__()的阻塞方式,用带超时的poll()方法拉取消息,避免永久阻塞。

修改后代码

AI-model.py

# AI-model.py
from kafka import KafkaConsumer, KafkaProducer

producer = KafkaProducer(bootstrap_servers='localhost:9092')
consumer = KafkaConsumer(
    'model-mailbox',
    bootstrap_servers='localhost:9092',
    group_id='ai-model-group',
    auto_offset_reset='latest',
    enable_auto_commit=True
)
# 等待消费者完成初始化、分区分配
consumer.poll(timeout_ms=2000)

while True:
    msg_batch = consumer.poll(timeout_ms=100)
    for topic_partition, messages in msg_batch.items():
        for img in messages:
            print(img.key)
            print('a-received')

            producer.send('webots-mailbox', key=b'movement', value=b'a')
            producer.flush()
            print('a-sent')

webots-controller.py

# webots-controller.py
from kafka import KafkaConsumer, KafkaProducer

producer = KafkaProducer(bootstrap_servers='localhost:9092')
consumer = KafkaConsumer(
    'webots-mailbox',
    bootstrap_servers='localhost:9092',
    group_id='webots-controller-group',
    auto_offset_reset='latest',
    enable_auto_commit=True
)
# 等待消费者完成初始化
consumer.poll(timeout_ms=2000)

while True:
    producer.send('model-mailbox', key=b'image', value=b'b')
    producer.flush()
    print('b-sent')

    msg_batch = consumer.poll(timeout_ms=1000)
    for topic_partition, messages in msg_batch.items():
        for movement in messages:
            print(movement.key)
            print('b-received')

内容的提问来源于stack exchange,提问作者Matin Zivdar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 01:39:17