RabbitMQ Publisher阻塞问题求助:需同时实现消息发布与回复消费
解决RabbitMQ Publisher阻塞问题:循环发送消息并处理回复
我一眼就看出问题出在start_consuming()这个方法上——它是个永久阻塞的调用,一旦执行就会一直停在那里监听队列,完全不会回到你的while循环里继续发消息。除此之外,你的corr_id每次循环都用同一个,这也会导致后续的回复匹配出问题,得一起改掉。
问题拆解
start_consuming()的阻塞特性:这个方法会无限期地监听队列、处理消息,直到你主动调用stop_consuming()。所以你的代码发完第一条消息后,就卡在consume_response里了,再也走不出这个函数,自然没法继续循环发消息。- 固定的
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)
关键修改说明
current_corr_id每次循环生成:确保每个请求的ID唯一,回复能正确匹配。wait_for_response替代consume_response:用process_data_events(time_limit=1)处理事件,每次等待1秒,直到收到对应的回复(通过response_received标记控制),这样不会永久阻塞。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
相关产品推荐
相关产品推荐

