RabbitMQ+aio_pika高并发下消息处理后无法发布问题求助
高并发下RabbitMQ消息发布阻塞问题排查与解决
环境与架构
基于RabbitMQ 3.12、Python 3.8、aio_pika 9.4.0搭建9个Docker容器的文本处理微服务架构,所有处理结果最终写入MongoDB:
- 服务A处理耗时最长1ms,生成消息体积小于1MB,需将结果发送至服务B、C后汇总到Sink节点
异常现象
- 低并发(≤100条消息)时系统运行正常
- 高并发(>1000条消息)时,服务A能正常消费并处理消息,但无法将结果发布到Exchange,消息始终处于未确认状态并被重复消费;其他处理速度更快的服务无此问题
- 少数成功流转的消息可被B、C正常处理、确认并发送至Sink;单条消息间隔10秒发送无异常
已尝试的无效方案
- 添加asyncio信号量无效果
- 接收消息时提前确认,会丢失约60%的消息
相关代码
共用AsyncPikaClient类代码
import logging from aio_pika import connect_robust, ExchangeType, Message, DeliveryMode import json from dotenv import load_dotenv import os import asyncio load_dotenv() user = os.environ["RABBITMQ_USER"] pwd = os.environ["RABBITMQ_PWD"] class AsyncPikaClient: def __init__(self, on_message_callback, i_consume_from, i_publish_to=None, max_concurrent_tasks=5): if on_message_callback : self.on_message = on_message_callback else : self.on_message = lambda x : x #identity function self.i_consume_from = i_consume_from self.i_publish_to = i_publish_to self.semaphore = asyncio.Semaphore(max_concurrent_tasks) self.publish_queue = [] self.consume_queue = None async def do_initialization(self):#, on_message_callback, i_consume_from, i_publish_to): self.connection = await connect_robust(f"amqp://{user}:{pwd}@rabbitmq:5672/") self.channel = await self.connection.channel() self.exchange = await self.channel.declare_exchange("direct", ExchangeType.DIRECT,) if self.i_consume_from: print("Initializing consume queue : ", self.i_consume_from) self.consume_queue = await self.channel.declare_queue(self.i_consume_from, durable=True) await self.consume_queue.bind(self.exchange, routing_key=self.i_consume_from) if self.i_publish_to: for name in self.i_publish_to : self.publish_queue.append(await self.channel.declare_queue(name, durable=True)) await self.publish_queue[-1].bind(self.exchange, routing_key=name) print('AIO pika connection initialized, queues are :', self.i_consume_from, self.i_publish_to) async def consume(self): await self.consume_queue.consume(self.process_incoming_message, no_ack=False) async def process_incoming_message(self, message): print(f"Got message on queue {self.i_consume_from}") body = message.body if body: try : mess_body = json.loads(body) try : data = await self.on_message(mess_body) except : print("Bug in treating the text !") print(f"Now the treatment is done in {self.i_consume_from}") await self.publish_all(data) print("Published !") await message.ack() print("Acked") except json.decoder.JSONDecodeError: print("Error : message has to be a JSON doc") await message.nack() except Exception as e: print(f"Unexpected error: {e}") await message.nack() else : print("Error : there is no body in the message") await message.nack() async def publish_all(self, message_body): if isinstance(message_body, list): publish_tasks = [] for item in message_body : publish_tasks.extend([ self.publish(item, queue) for queue in self.i_publish_to ]) else : publish_tasks = [ self.publish(message_body, queue) for queue in self.i_publish_to ] await asyncio.gather(*publish_tasks, return_exceptions=True) async def publish(self, message_body: dict, routing_key:str): """Method to publish message to RabbitMQ""" async with self.semaphore: message = Message( json.dumps(message_body).encode('utf-8'), delivery_mode=DeliveryMode.PERSISTENT, app_id=self.i_consume_from ) print(len(json.dumps(message_body).encode('utf-8'))) try : print(len(message)) except : print("") print(f"Message ready, trying to publish to queue {routing_key}") #everything prints to here await self.exchange.publish(message, routing_key=routing_key) #this does not happen print(f"Published to exchange using routing key : ", routing_key) #this never gets printed
服务A示例代码
import asyncio from src.client_pika import AsyncPikaClient async def counts_on_udpipe_annotations(received: Dict) -> None: #some treatment return result if __name__ == "__main__": pika_client = AsyncPikaClient(on_message_callback=counts_on_udpipe_annotations, i_consume_from="queueA", i_publish_to=["queueB", "queueC"]) loop = asyncio.get_event_loop() loop.run_until_complete(pika_client.do_initialization()) loop.create_task(pika_client.consume()) loop.run_forever()
问题原因分析
- 通道并发限制不足:aio_pika的Channel默认有
max_inflight_messages限制(默认值通常为10),高并发下publish请求被阻塞,导致无法完成发布,进而无法ack原消息,触发重复消费 - 信号量作用位置错误:当前信号量仅包裹publish内部逻辑,未限制消费-处理-发布的全流程并发数,导致消费速度远快于发布速度,消息堆积引发阻塞
- 异常静默掩盖问题:
publish_all使用return_exceptions=True,导致publish过程中的异常被静默,无法定位发布失败的具体原因 - 冗余队列声明占用资源:服务A声明B、C的队列并绑定Exchange,属于不必要的通道操作,额外占用RabbitMQ连接资源
解决方案
1. 提高通道并发处理上限
在创建Channel时显式设置max_inflight_messages,适配高并发场景:
async def do_initialization(self): self.connection = await connect_robust(f"amqp://{user}:{pwd}@rabbitmq:5672/") # 根据实际并发量调整,建议设置为100-200 self.channel = await self.connection.channel(max_inflight_messages=150) self.exchange = await self.channel.declare_exchange("direct", ExchangeType.DIRECT,) # 保留消费队列的声明逻辑,其余代码不变
2. 调整信号量作用范围,控制消费并发
将信号量移至process_incoming_message入口,限制同时处理的消息数量,避免消费速度远超发布能力:
async def process_incoming_message(self, message): async with self.semaphore: # 移至此处,控制并发处理的消息总数 print(f"Got message on queue {self.i_consume_from}") body = message.body if body: try : mess_body = json.loads(body) try : data = await self.on_message(mess_body) except Exception as e: print(f"Bug in treating the text: {e}") await message.nack() return print(f"Now the treatment is done in {self.i_consume_from}") await self.publish_all(data) print("Published !") await message.ack() print("Acked") except json.decoder.JSONDecodeError: print("Error : message has to be a JSON doc") await message.nack() except Exception as e: print(f"Unexpected error: {e}") await message.nack() else : print("Error : there is no body in the message") await message.nack()
同时删除publish方法中的async with self.semaphore:代码块,避免双重限制。
3. 移除冗余的队列声明逻辑
服务A无需声明B、C的队列,删除do_initialization中对应代码,由B、C各自负责声明并绑定自己的队列:
# 移除以下代码块 # if self.i_publish_to: # for name in self.i_publish_to : # self.publish_queue.append(await self.channel.declare_queue(name, durable=True)) # await self.publish_queue[-1].bind(self.exchange, routing_key=name)
4. 取消异常静默,暴露发布错误
修改publish_all的asyncio.gather调用,移除return_exceptions=True,让异常冒泡以便排查:
async def publish_all(self, message_body): if isinstance(message_body, list): publish_tasks = [] for item in message_body : publish_tasks.extend([ self.publish(item, queue) for queue in self.i_publish_to ]) else : publish_tasks = [ self.publish(message_body, queue) for queue in self.i_publish_to ] # 移除return_exceptions=True,捕获发布异常 await asyncio.gather(*publish_tasks)
5. 可选:批量确认优化性能
如果消息量极大,可使用批量ack减少网络交互,但需注意数据一致性风险:
# 在AsyncPikaClient的__init__中添加计数器和阈值 def __init__(self, on_message_callback, i_consume_from, i_publish_to=None, max_concurrent_tasks=5, batch_ack_threshold=10): # 原有代码不变 self.ack_counter = 0 self.batch_ack_threshold = batch_ack_threshold # 修改process_incoming_message中的ack逻辑 await self.publish_all(data) self.ack_counter += 1 if self.ack_counter >= self.batch_ack_threshold: await message.ack(multiple=True) self.ack_counter = 0 else: await message.ack()
内容的提问来源于stack exchange,提问作者Cyriel Mallart
相关产品推荐
相关产品推荐

