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

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()

问题原因分析

  1. 通道并发限制不足:aio_pika的Channel默认有max_inflight_messages限制(默认值通常为10),高并发下publish请求被阻塞,导致无法完成发布,进而无法ack原消息,触发重复消费
  2. 信号量作用位置错误:当前信号量仅包裹publish内部逻辑,未限制消费-处理-发布的全流程并发数,导致消费速度远快于发布速度,消息堆积引发阻塞
  3. 异常静默掩盖问题:publish_all使用return_exceptions=True,导致publish过程中的异常被静默,无法定位发布失败的具体原因
  4. 冗余队列声明占用资源:服务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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 03:37:07