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

RabbitMQ Publisher阻塞问题求助:需同时实现消息发布与回复消费

解决RabbitMQ Publisher阻塞问题:循环发送消息并处理回复

我一眼就看出问题出在start_consuming()这个方法上——它是个永久阻塞的调用,一旦执行就会一直停在那里监听队列,完全不会回到你的while循环里继续发消息。除此之外,你的corr_id每次循环都用同一个,这也会导致后续的回复匹配出问题,得一起改掉。

问题拆解

  1. start_consuming()的阻塞特性:这个方法会无限期地监听队列、处理消息,直到你主动调用stop_consuming()。所以你的代码发完第一条消息后,就卡在consume_response里了,再也走不出这个函数,自然没法继续循环发消息。
  2. 固定的corr_id:每次发送请求都用同一个correlation ID,当你后续发送新消息时,Subscriber的回复会和之前的ID匹配,导致on_response里的判断失效,没法正确处理新回复。

解决方案

我们需要做两个关键修改:

  • 用非阻塞的消息处理方式代替start_consuming(),处理完回复就退出监听,回到循环继续发消息。
  • 每次循环生成新的corr_id,确保每个请求和回复一一对应。

修改后的Publisher代码

#!/usr/bin/env python
import pika
import sys
import uuid
import time

connection = pika.BlockingConnection(
    pika.ConnectionParameters(host='localhost'))
channel = connection.channel()

channel.queue_declare(queue='task_queue', durable=True)
channel.queue_declare(queue='reply_queue', durable=True)

message = "Hello World!"
response_received = False
current_corr_id = ""

def on_response(ch, method, properties, body):
    global response_received
    print("----- On_response -----")
    print("Received CORR_ID : ", properties.correlation_id)
    if current_corr_id == properties.correlation_id:
        resp = body.decode('utf-8')
        print("RESPONSE : ", resp)
        response_received = True
    ch.basic_ack(delivery_tag=method.delivery_tag)

def wait_for_response(channel):
    global response_received
    response_received = False
    # 注册回调,不启动永久监听
    consumer_tag = channel.basic_consume(queue='reply_queue', on_message_callback=on_response)
    # 循环处理事件,直到收到对应的回复
    while not response_received:
        # 处理待处理的事件,每次等待1秒,避免永久阻塞
        connection.process_data_events(time_limit=1)
    # 处理完回复后取消消费者,避免重复注册
    channel.basic_cancel(consumer_tag=consumer_tag)

while True:
    current_corr_id = str(uuid.uuid4())
    channel.basic_publish(
        exchange='',
        routing_key='task_queue',
        body=message,
        properties=pika.BasicProperties(
            reply_to='reply_queue',
            correlation_id=current_corr_id
        ))
    print(" [x] Sent %r (CORR_ID: %s)" % (message, current_corr_id))
    # 等待回复,处理完就回到循环
    wait_for_response(channel)

关键修改说明

  1. current_corr_id每次循环生成:确保每个请求的ID唯一,回复能正确匹配。
  2. wait_for_response替代consume_response:用process_data_events(time_limit=1)处理事件,每次等待1秒,直到收到对应的回复(通过response_received标记控制),这样不会永久阻塞。
  3. channel.basic_cancel:每次处理完回复后取消消费者,避免重复注册导致回调被多次触发。

运行效果

修改后,你的流程就能正常循环了:

Publisher → 发送消息"Hello World!"(带新的CORR_ID)→ Subscriber
Subscriber → 处理消息,回复带相同CORR_ID的消息到reply_queue → Publisher
Publisher收到回复并打印,然后回到循环,继续发送下一条消息

额外小提示

如果担心process_data_events的超时问题,也可以用channel.basic_get()手动获取单条消息,逻辑会更直观:

def wait_for_response(channel):
    global response_received
    response_received = False
    while not response_received:
        method_frame, header_frame, body = channel.basic_get(queue='reply_queue')
        if method_frame:
            on_response(channel, method_frame, header_frame, body)
        else:
            # 没收到消息时短暂休眠,避免CPU空转
            time.sleep(0.1)

内容的提问来源于stack exchange,提问作者Shital Jadhav

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 09:22:11