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