news 2026/9/10 22:47:10

AMQP异步实现aio_pika 2.0

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
AMQP异步实现aio_pika 2.0

AMQP(Advanced 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( level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s" ) logger = logging.getLogger("async_amqp_sdk") class RabbitMQConnector: def __init__( self, rabbitmq_ip, rabbitmq_port, username="guest", password="guest", heartbeat=60, ): 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( f"amqp://{self.username}:{self.password}@{self.rabbitmq_ip}:{self.rabbitmq_port}/", heartbeat=self.heartbeat, blocked_connection_callback=self.on_connection_blocked, unblocked_connection_callback=self.on_connection_unblocked, ) # 创建 Channel self.channel = await self.connection.channel() logger.info("Connected to RabbitMQ") except Exception as ex: logger.error(f"Failed 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_name="eason_exchange", exchange_type="topic", ): 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, type=self.exchange_type, durable=True ) 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( body=json.dumps(message).encode("utf-8"), content_type="application/json", delivery_mode=aio_pika.DeliveryMode.PERSISTENT, # 持久化消息 ), routing_key=routing_key, ) logger.info( f"Sent message to {self.exchange_name} routing_key: {routing_key}" ) except Exception as ex: logger.error(f"Failed to send message: {ex}") class RabbitMQConsumer: def __init__( self, RabbitMQConnector, exchange_name, queue_name, routing_key, ttl=21600000, ): 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, type=aio_pika.ExchangeType.TOPIC, durable=True ) # 声明队列 args = {"x-message-ttl": self.ttl} if self.ttl else {} self.queue = await self.RabbitMQConnector.channel.declare_queue( self.queue_name, durable=True, arguments=args, ) # 绑定队列到交换机 await self.queue.bind(self.exchange, self.routing_key) return self async def callback(self, on_message_callback, ack=False): """ 内部方法,用于实际消费消息。 """ try: # 开始消费队列中的消息 await self.queue.consume( on_message_callback, no_ack=ack, ) # no_ack=False 表示需要手动确认 logger.info(f"Start 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( f"Received message : {message.body.decode()}" ) # 模拟消息处理机制 # await message.ack() #无需手动确认 except Exception as ex: # 如果需要,可以手动拒绝消息并重新入队 await message.reject(requeue=True) logger.error(f"Error 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_ip=rabbitmq_ip, rabbitmq_port=rabbitmq_port, ) await rabbitMQConnector.connect() # 测试发布消息 publisher = RabbitMQProducer(rabbitMQConnector, exchange_name=exchange_name) await publisher.initialize() # 发布测试消息 await publisher.publish(message, routing_key) # 消费测试消息 consumer = RabbitMQConsumer( rabbitMQConnector, exchange_name=exchange_name, queue_name=queue_name, routing_key=routing_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(f"main function error: {ex}") if __name__ == "__main__": asyncio.run(main())
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 22:45:02

从“事后追溯”到“实时可溯”:无感定位技术如何重塑战场数据链 技术白皮书

1 概述1.1 技术背景智能化联合作战的制胜核心,已从火力对抗、兵力对抗全面转向数据链速度、认知链路精度、战场数据可信度的体系对抗。战场数据链作为指挥、感知、机动、打击、评估的核心神经网络,决定了OODA作战循环的闭环效率与战局博弈主动权。传统战…

作者头像 李华
网站建设 2026/9/10 22:43:00

深度相机技术:机器人软件开发的感知核心艺术与应用实战

在当今机器人技术快速发展的浪潮中,感知能力成为其实现智能化的基石。深度相机作为一项突破性创新,凭借其精准捕捉三维场景的能力,正重塑机器人软件开发的格局。它不仅为导航避障提供关键数据,还为物体识别、交互行为等场景带来革命性影响。本文将深入剖析深度相机的技术本…

作者头像 李华
网站建设 2026/9/10 22:42:24

企业数据中台建设:核心痛点与qData商业版解决方案

1. 企业数据中台的市场现状与核心痛点 当前企业数字化转型已进入深水区,数据资产的价值挖掘成为核心竞争力。根据行业调研数据显示,超过78%的中大型企业在2023年已将数据中台建设列入战略优先级,但实际落地效果参差不齐。传统自建数据平台往往…

作者头像 李华
网站建设 2026/9/10 22:40:43

自由学习记录:构建个性化知识管理系统的实践指南

1. 项目概述:自由学习记录的本质与实践价值 "自由学习记录(144)"这个看似简单的标题背后,隐藏着一个持续学习者的完整知识管理体系。作为坚持到第144篇的学习记录,它代表着一种持续性的个人成长方法论。这类…

作者头像 李华