启用客户端心跳后,Pika/RabbitMQ线程消费者长任务处理故障
问题:RabbitMQ RPC长任务导致连接断开、响应无法返回
问题背景
- 运行可配置线程消费者(20-30线程),已守护进程化,支持CLI启停,基于RabbitMQ实现RPC模式
- 仅消费者端启用心跳(600秒),用于避免服务器OS内核断开空闲连接
- 短任务场景下运行正常,但遇到长任务(甚至需数天执行)时,生产者无法收到响应:Broker会关闭连接并将消息重新入队,其他消费者处理时重复该问题,最终所有消费者连接关闭,生产者持续等待
- 尝试过线程实现和asyncio异步调用两种方案,均出现相同故障
消费者代码实现
def mq_connection_factory(queue_name): '''global method to create and return mq connection and channel''' mq_host, mq_port, mq_usr, mq_passd = conf.get('some-config-file') credentials = pika.PlainCredentials(mq_usr, mq_passd) parameters = pika.ConnectionParameters(host=mq_host,port=mq_port, credentials=credentials, heartbeat = 600) connection = pika.BlockingConnection(parameters) channel = connection.channel() channel.queue_declare(queue='queue_name,durable=True) return connection, channel class ThreadedConsumer(threading.Thread): def __init__(self,queue_name): threading.Thread.__init__(self) self.threads = conf.get('daemon','max_parallel_thread') self.queue_name = queue_name self.channel = None self.connection = None def app_process(self,json_conf,choice): '''The main serever function based on incoming message''' #long running process #calls different modules to process the message and create multiline string output #to be sent back to producer return "with great power comes great responsibilities" def app_process_response(self, ch, delivery_tag, props, json_body, choice, connection): json_conf = json_body thread_id = threading.get_ident() log.info(f'Thread id: {thread_id} Delivery tag: {delivery_tag}') try: response = self.app_process(json_conf,choice) except Exception as e: response = f'app process failed with the error : {e}' log.info("App process response captured.") cb = functools.partial(self.ack_app_response, ch, delivery_tag, props, response) connection.add_callback_threadsafe(cb) def ack_app_response(self, ch, delivery_tag, props, response): log.info(" [x] Done") if ch.is_open: log.info('Channel is open. Going ahead with writing acknowledgement back to channel.') ch.basic_publish(exchange='', routing_key=props.reply_to, properties=pika.BasicProperties(correlation_id = \ props.correlation_id), body=str(response)) ch.basic_ack(delivery_tag=delivery_tag) else: log.error('''The Channel between Client APP and server APP is already closed. Acknowledgement can't be sent back to Client APP.''') raise Exception('''The Channel between Client APP and server APP is already closed. Acknowledgement can't be sent back to Client APP.''') log.info(f'Response was sent to client APP.') def on_request(self,ch, method, props, body, args): json_body = json.loads(body) log.info(" [x] Received %r" % json_body) choice = "make someone proud" (connection, threads) = args delivery_tag = method.delivery_tag t = threading.Thread(target = self.app_process_response, args= (ch, delivery_tag, props, json_body, choice, connection)) t.start() threads.append(t) def run(self): try: log.info('Connecting to Rabbit MQ Host.') self.connection, self.channel = mq_connection_factory(self.queue_name) log.info('Connection to rabbit MQ Host was successful.') self.channel.basic_qos(prefetch_count=self.threads*10) # i have tried 1 as well threads = [] on_message_callback = functools.partial(self.on_request, args = (self.connection, threads)) self.channel.basic_consume(self.queue_name, on_message_callback=on_message_callback) log.info('starting thread to listen to client request...') self.channel.start_consuming() except Exception as e: raise Exception('Failed to establish connetion to mq. Please check with your rabbit MQ broker.') from e #doubtful about this as my daemon is pid file based and keeps running until i run `stop daemon`` for thread in threads: thread.join() def initialize_thread(queue_name): '''Intialise threaded consumer for app to receive request and send them back to server core app''' threads = conf.get('daemon','max_parallel_thread') for i in range(threads): log.info(f'launching thread : {i}') td = ThreadedConsumer(queue_name=queue_name) td.start()
问题核心分析
- 心跳阻塞:使用
BlockingConnection时,start_consuming运行在单线程循环中,长任务会阻塞该线程,导致无法发送心跳包,Broker超过心跳超时后判定连接失效并关闭。 - 线程模型冗余:每个
ThreadedConsumer线程创建独立MQ连接,又在on_request中再开子线程处理任务,既浪费Broker连接资源,又无法保证连接心跳的持续维护。 - 连接失效后无处理:连接断开后,后续尝试发送响应/确认消息时,因通道已关闭直接失败,没有重连或降级处理逻辑。
解决方案
1. 替换为异步连接维护心跳
使用pika.SelectConnection替代BlockingConnection,它会在后台线程处理I/O和心跳,不会被长任务阻塞:
def mq_connection_factory(queue_name): mq_host, mq_port, mq_usr, mq_passd = conf.get('some-config-file') credentials = pika.PlainCredentials(mq_usr, mq_passd) parameters = pika.ConnectionParameters( host=mq_host, port=mq_port, credentials=credentials, heartbeat=600, blocked_connection_timeout=300 ) # 异步连接,后台维护I/O和心跳 connection = pika.SelectConnection(parameters) channel = connection.channel() channel.queue_declare(queue=queue_name, durable=True) return connection, channel
- 配合
blocked_connection_timeout设置,避免连接被Broker阻塞后无响应。 - 在
run方法中启动异步I/O循环:def run(self): try: self.connection, self.channel = mq_connection_factory(self.queue_name) self.channel.basic_qos(prefetch_count=self.threads) self.channel.basic_consume(self.queue_name, on_message_callback=self.on_request) # 启动异步I/O循环 self.connection.ioloop.start() except Exception as e: log.error(f"连接MQ失败: {e}") if self.connection and self.connection.is_open: self.connection.close()
2. 分离任务执行与MQ通信线程
使用线程池处理长任务,MQ回调线程仅负责接收消息和提交任务,避免阻塞心跳:
from concurrent.futures import ThreadPoolExecutor class ThreadedConsumer(threading.Thread): def __init__(self, queue_name): super().__init__() self.queue_name = queue_name self.connection = None self.channel = None # 初始化线程池,复用线程资源 self.executor = ThreadPoolExecutor(max_workers=conf.get('daemon','max_parallel_thread')) def on_task_complete(self, future, ch, delivery_tag, props): try: response = future.result() except Exception as e: response = f'app process failed with the error : {e}' if ch.is_open and self.connection.is_open: # 线程安全地发送响应 cb = functools.partial(self._send_response, ch, delivery_tag, props, response) self.connection.add_callback_threadsafe(cb) else: log.error("连接/通道已关闭,无法发送响应") def _send_response(self, ch, delivery_tag, props, response): ch.basic_publish( exchange='', routing_key=props.reply_to, properties=pika.BasicProperties(correlation_id=props.correlation_id), body=str(response) ) ch.basic_ack(delivery_tag=delivery_tag) def on_request(self, ch, method, props, body): json_body = json.loads(body) log.info(f"Received {json_body}") choice = "make someone proud" # 提交长任务到线程池,完成后回调发送响应 future = self.executor.submit(self.app_process, json_body, choice) future.add_done_callback( lambda f: self.on_task_complete(f, ch, method.delivery_tag, props) )
3. 优化连接管理
- 减少连接数量:当前
initialize_thread创建20个独立MQ连接,改为使用连接池或单连接多通道,降低Broker资源消耗。 - 添加重连逻辑:注册连接关闭回调,自动尝试重连,保证消费者可用性:
def on_connection_closed(connection, reason): log.warning(f"连接关闭: {reason},5秒后尝试重连") connection.ioloop.call_later(5, lambda: mq_connection_factory('queue_name')) # 创建连接时注册回调 connection.add_on_close_callback(on_connection_closed)
4. 调整长任务策略
- 拆分长任务:将数天级任务拆分为多个短任务,通过MQ分步执行,每步完成后向生产者发送中间状态,避免单个任务占用连接过长。
- 延长心跳超时:如果业务允许,将
heartbeat设置为更大值(如3600秒),同时调整Broker的heartbeat_timeout配置(默认是心跳间隔的2倍)。
内容的提问来源于stack exchange,提问作者Biswadeep27
相关产品推荐
相关产品推荐

