使用Python Pika连接双RabbitMQ服务器时连接异常的解决方案咨询
解决Pika多连接/HTTP场景下RabbitMQ连接心跳超时问题
你遇到的问题本质是同步阻塞的连接模型导致其他RabbitMQ连接的心跳无法被及时处理,RabbitMQ服务器因长时间未收到心跳包,主动关闭了连接。以下是几个可行的解决方案,附具体实现方式:
1. 使用Pika异步连接(推荐)
Pika的SelectConnection(基于select多路复用)或AsyncioConnection(基于asyncio)专为多连接/多IO场景设计,不会阻塞在单个连接的消费操作上,而是通过事件循环同时监听所有连接的IO事件(包括心跳包收发)。
多RabbitMQ连接示例(SelectConnection)
核心思路是为每个RabbitMQ服务器创建独立的连接和通道,将所有连接的IO事件注册到同一事件循环,确保心跳和消息处理都能被及时调度:
import pika from pika.connection import SelectConnection from pika.channel import Channel # 服务器A配置 CONFIG_A = { 'host': 'server_a_ip', 'credentials': pika.PlainCredentials('user_a', 'pass_a') } # 服务器B配置 CONFIG_B = { 'host': 'server_b_ip', 'credentials': pika.PlainCredentials('user_b', 'pass_b') } channel_b = None def on_message_a(channel: Channel, method, properties, body): # 处理消息逻辑 processed_data = body.decode('utf-8').upper() # 发送到服务器B的交换器 if channel_b and channel_b.is_open: channel_b.basic_publish( exchange='exchange_b', routing_key='routing_key_b', body=processed_data.encode('utf-8') ) channel.basic_ack(delivery_tag=method.delivery_tag) def setup_channel_a(connection: SelectConnection): channel = connection.channel() channel.exchange_declare(exchange='exchange_a1', exchange_type='direct') queue = channel.queue_declare(queue='queue_a1', exclusive=True) channel.queue_bind(exchange='exchange_a1', queue=queue.method.queue, routing_key='key_a1') channel.basic_consume(queue=queue.method.queue, on_message_callback=on_message_a) def setup_channel_b(connection: SelectConnection): global channel_b channel_b = connection.channel() channel_b.exchange_declare(exchange='exchange_b', exchange_type='direct') def on_connection_open(connection: SelectConnection, setup_channel_func): connection.channel(on_open_callback=lambda ch: setup_channel_func(connection)) def main(): # 创建服务器A的连接 conn_a = SelectConnection(pika.ConnectionParameters(**CONFIG_A), on_open_callback=lambda conn: on_connection_open(conn, setup_channel_a)) # 创建服务器B的连接 conn_b = SelectConnection(pika.ConnectionParameters(**CONFIG_B), on_open_callback=lambda conn: on_connection_open(conn, setup_channel_b)) # 合并两个连接的事件循环 try: while True: # 处理所有连接的IO事件,超时时间设为心跳间隔的一半(5秒) conn_a.process_data_events(time_limit=5) conn_b.process_data_events(time_limit=5) except KeyboardInterrupt: conn_a.close() conn_b.close() if __name__ == '__main__': main()
2. 多线程分离连接(兼容同步模式)
如果更习惯使用同步的BlockingConnection,可以给每个RabbitMQ连接单独分配线程,线程间用线程安全队列传递消息:
import pika import threading from queue import Queue # 线程安全队列,传递待发送到B的消息 message_queue = Queue(maxsize=100) def consume_from_a(): conn_a = pika.BlockingConnection(pika.ConnectionParameters( host='server_a_ip', credentials=pika.PlainCredentials('user_a', 'pass_a') )) channel_a = conn_a.channel() channel_a.exchange_declare(exchange='exchange_a1', exchange_type='direct') queue = channel_a.queue_declare(queue='queue_a1', exclusive=True) channel_a.queue_bind(exchange='exchange_a1', queue=queue.method.queue, routing_key='key_a1') def callback(ch, method, properties, body): processed_data = body.decode('utf-8').upper() message_queue.put(processed_data) ch.basic_ack(delivery_tag=method.delivery_tag) channel_a.basic_consume(queue=queue.method.queue, on_message_callback=callback) try: channel_a.start_consuming() except KeyboardInterrupt: conn_a.close() def publish_to_b(): while True: try: conn_b = pika.BlockingConnection(pika.ConnectionParameters( host='server_b_ip', credentials=pika.PlainCredentials('user_b', 'pass_b') )) channel_b = conn_b.channel() channel_b.exchange_declare(exchange='exchange_b', exchange_type='direct') while True: # 队列取消息时设置超时,确保连接能处理心跳 message = message_queue.get(timeout=5) channel_b.basic_publish( exchange='exchange_b', routing_key='routing_key_b', body=message.encode('utf-8') ) message_queue.task_done() except (pika.exceptions.ConnectionClosed, pika.exceptions.ChannelClosed): # 连接断开自动重连 continue except KeyboardInterrupt: conn_b.close() break if __name__ == '__main__': t_consume = threading.Thread(target=consume_from_a, daemon=True) t_consume.start() t_publish = threading.Thread(target=publish_to_b, daemon=True) t_publish.start() t_consume.join() t_publish.join()
这里的关键是:发送线程的message_queue.get(timeout=5)不会永久阻塞,每隔5秒唤醒一次,此时BlockingConnection会自动处理心跳包,即使队列无消息,连接也能保持活跃。
3. 同步模式手动处理心跳(仅兼容旧代码,不推荐)
若必须在单线程中使用同步连接,可在消费A的消息时设置超时,定期处理B连接的心跳事件:
import pika CONFIG_A = { 'host': 'server_a_ip', 'credentials': pika.PlainCredentials('user_a', 'pass_a') } CONFIG_B = { 'host': 'server_b_ip', 'credentials': pika.PlainCredentials('user_b', 'pass_b') } def main(): conn_a = pika.BlockingConnection(pika.ConnectionParameters(**CONFIG_A)) channel_a = conn_a.channel() channel_a.exchange_declare(exchange='exchange_a1', exchange_type='direct') queue = channel_a.queue_declare(queue='queue_a1', exclusive=True) channel_a.queue_bind(exchange='exchange_a1', queue=queue.method.queue, routing_key='key_a1') conn_b = pika.BlockingConnection(pika.ConnectionParameters(**CONFIG_B)) channel_b = conn_b.channel() channel_b.exchange_declare(exchange='exchange_b', exchange_type='direct') while True: try: # 消费A的消息,设置超时时间5秒 method_frame, properties, body = channel_a.basic_get(queue=queue.method.queue, auto_ack=False) if method_frame: processed_data = body.decode('utf-8').upper() channel_b.basic_publish( exchange='exchange_b', routing_key='routing_key_b', body=processed_data.encode('utf-8') ) channel_a.basic_ack(delivery_tag=method_frame.delivery_tag) else: # 无消息时处理B连接的心跳事件 conn_b.process_data_events(time_limit=1) except (pika.exceptions.ConnectionClosed, pika.exceptions.ChannelClosed): # 重连B conn_b = pika.BlockingConnection(pika.ConnectionParameters(**CONFIG_B)) channel_b = conn_b.channel() channel_b.exchange_declare(exchange='exchange_b', exchange_type='direct') except KeyboardInterrupt: conn_a.close() conn_b.close() break if __name__ == '__main__': main()
内容的提问来源于stack exchange,提问作者Dan H
相关产品推荐
相关产品推荐

