资讯动态

AMQP异步实现aio_pika 2.0

发布时间:2026/9/10 22:47:16 来源:尧图企业网站定制
AMQPAdvanced Message Queuing Protocol高级消息队列协议是一个为面向消息的中间件设计的应用层标准协议旨在实现跨平台、跨语言的异步消息通信。核心定位与背景AMQP 最初由摩根大通于 2003 年提出后于 2011 年成为 OASIS 标准2014 年获得 ISO/IEC 国际认证。该协议的核心目标是让不同语言、不同平台开发的客户端只要遵循同一套协议就能互相收发消息实现真正的互操作性。四个核心概念的角色1. Producer生产者角色定义‌消息的发送方通常是应用程序的一部分。核心职责‌创建消息并将其发送到 RabbitMQ Broker服务器。工作逻辑‌生产者不直接将消息发送给队列而是发送给 ‌Exchange交换机‌。在发送消息时生产者可以指定一个 ‌Routing Key路由键‌这是一个字符串标签用于帮助交换机决定如何将消息路由到特定的队列。生产者不知道消息最终会被哪个消费者接收实现了发送端与接收端的解耦。2. Consumer消费者角色定义‌消息的接收方通常是另一个应用程序或服务。核心职责‌从队列中获取消息并进行业务处理。工作逻辑‌消费者连接到 RabbitMQ 服务器并订阅特定的 ‌Queue队列‌。消费者只关心从队列中取出消息它不知道消息是由哪个生产者发送的也不需要知道交换机的路由规则。一旦消息被消费者成功确认Ack该消息通常会从队列中删除除非配置了持久化或其他特殊策略。3. Exchange交换机角色定义‌消息的路由中心相当于邮局的分拣中心。核心职责‌接收生产者发送的消息并根据指定的规则将消息路由到一个或多个队列中。关键特性‌不存储消息‌交换机本身不保存消息如果消息无法路由到任何队列它可能会被丢弃或返回给生产者取决于配置。路由类型‌交换机有多种类型决定了不同的路由行为Direct‌精确匹配 Routing Key。Fanout‌广播模式将消息路由到所有绑定的队列。Topic‌主题匹配支持通配符* 和 #进行模糊匹配。Headers‌根据消息头属性进行匹配。4. Queue队列角色定义‌消息的存储容器相当于邮局的信箱。核心职责‌缓存消息直到消费者将其取走。关键特性‌先进先出FIFO‌默认情况下消息按照进入队列的顺序被消费。持久化‌队列和消息可以配置为持久化即使 RabbitMQ 重启数据也不会丢失。多消费者支持‌多个消费者可以监听同一个队列RabbitMQ 会以轮询等方式将消息分发给不同的消费者实现负载均衡。import asyncio import aio_pika import logging import json logging.basicConfig( levellogging.INFO, format%(asctime)s - %(levelname)s - %(message)s ) logger logging.getLogger(async_amqp_sdk) class RabbitMQConnector: def __init__( self, rabbitmq_ip, rabbitmq_port, usernameguest, passwordguest, heartbeat60, ): self.rabbitmq_ip rabbitmq_ip self.rabbitmq_port rabbitmq_port self.username username self.password password self.heartbeat heartbeat self.connection None self.channel None async def on_connection_blocked(self): logger.warning(Connection to RabbitMQ was blocked) async def on_connection_unblocked(self): logger.info(Connection to RabbitMQ was unblocked) async def connect(self): while not self.channel: try: # 创建Connection self.connection await aio_pika.connect_robust( famqp://{self.username}:{self.password}{self.rabbitmq_ip}:{self.rabbitmq_port}/, heartbeatself.heartbeat, blocked_connection_callbackself.on_connection_blocked, unblocked_connection_callbackself.on_connection_unblocked, ) # 创建 Channel self.channel await self.connection.channel() logger.info(Connected to RabbitMQ) except Exception as ex: logger.error(fFailed to connect to RabbitMQ: {ex}) await asyncio.sleep(5) async def close(self): if self.connection and not self.connection.is_closed: await self.connection.close() logger.info(Connection closed) class RabbitMQProducer: def __init__( self, RabbitMQConnector, exchange_nameeason_exchange, exchange_typetopic, ): self.RabbitMQConnector RabbitMQConnector self.exchange_name exchange_name self.exchange_type exchange_type # 声明交换机 async def initialize(self): 初始化交换机 if not self.RabbitMQConnector.channel: await self.RabbitMQConnector.connect() self.exchange await self.RabbitMQConnector.channel.declare_exchange( self.exchange_name, typeself.exchange_type, durableTrue ) return self async def publish(self, message, routing_keys: list): # 确保连接和交换机已初始化 if not self.exchange: await self.initialize() # 发布消息 for routing_key in routing_keys: try: await self.exchange.publish( aio_pika.Message( bodyjson.dumps(message).encode(utf-8), content_typeapplication/json, delivery_modeaio_pika.DeliveryMode.PERSISTENT, # 持久化消息 ), routing_keyrouting_key, ) logger.info( fSent message to {self.exchange_name} routing_key: {routing_key} ) except Exception as ex: logger.error(fFailed to send message: {ex}) class RabbitMQConsumer: def __init__( self, RabbitMQConnector, exchange_name, queue_name, routing_key, ttl21600000, ): self.RabbitMQConnector RabbitMQConnector self.exchange_name exchange_name self.queue_name queue_name self.routing_key routing_key self.consumer_task None self.ttl ttl async def initialize(self): 初始化队列和绑定 if not self.RabbitMQConnector.channel: await self.RabbitMQConnector.connect() # 声明交换机 self.exchange await self.RabbitMQConnector.channel.declare_exchange( self.exchange_name, typeaio_pika.ExchangeType.TOPIC, durableTrue ) # 声明队列 args {x-message-ttl: self.ttl} if self.ttl else {} self.queue await self.RabbitMQConnector.channel.declare_queue( self.queue_name, durableTrue, argumentsargs, ) # 绑定队列到交换机 await self.queue.bind(self.exchange, self.routing_key) return self async def callback(self, on_message_callback, ackFalse): 内部方法用于实际消费消息。 try: # 开始消费队列中的消息 await self.queue.consume( on_message_callback, no_ackack, ) # no_ackFalse 表示需要手动确认 logger.info(fStart to consume from queue: {self.queue_name}) except Exception as ex: logger.error(ex) async def consume(self, on_message_callback): 启动消费任务。 try: # 确保连接和队列已初始化 if not self.queue: await self.initialize() if self.consumer_task is None or self.consumer_task.done(): self.consumer_task asyncio.create_task( self.callback(on_message_callback, False) ) logger.info(Consuming task started) await self.consumer_task # 保持程序运行以接收消息 await asyncio.Future() # 创建一个永远不会完成的 Future except Exception as ex: logger.error(ex) # 消息处理机制 async def on_message(message: aio_pika.IncomingMessage): deal with message async with message.process(): # 确保消息在处理完成后自动处理确认或拒绝 try: logger.info( fReceived message : {message.body.decode()} ) # 模拟消息处理机制 # await message.ack() #无需手动确认 except Exception as ex: # 如果需要可以手动拒绝消息并重新入队 await message.reject(requeueTrue) logger.error(fError processing message: {ex}) async def main(): try: rabbitmq_ip 10.146.212.85 rabbitmq_port 30025 exchange_name eason_exchange queue_name eason_queue routing_key eason_routing message {send message amqp:: this is a test message aio_pika} # 创建连接 rabbitMQConnector RabbitMQConnector( rabbitmq_iprabbitmq_ip, rabbitmq_portrabbitmq_port, ) await rabbitMQConnector.connect() # 测试发布消息 publisher RabbitMQProducer(rabbitMQConnector, exchange_nameexchange_name) await publisher.initialize() # 发布测试消息 await publisher.publish(message, routing_key) # 消费测试消息 consumer RabbitMQConsumer( rabbitMQConnector, exchange_nameexchange_name, queue_namequeue_name, routing_keyrouting_key, ) await consumer.initialize() # 启动消费者在后台运行 consumer_task asyncio.create_task(consumer.consume(on_message)) # 等待一段时间让消费者运行 await asyncio.sleep(10) # 取消消费者任务 consumer_task.cancel() try: await consumer_task except asyncio.CancelledError: logger.error(Consumer task cancelled) # 关闭连接 await rabbitMQConnector.close() except Exception as ex: logger.error(fmain function error: {ex}) if __name__ __main__: asyncio.run(main())

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价