2.12 消息路由流程 (Message Routing)


文档摘要

2.12 消息路由流程 (Message Routing) 2.12 消息路由流程 (Message Routing) 详解 在 RabbitMQ 消息队列系统中,消息路由 (Message Routing) 是核心概念之一,它决定了生产者 (Producer) 发送的消息如何被正确地投递到一个或多个队列 (Queue) 中,最终被消费者 (Consumer) 接收和处理。 消息路由的效率和灵活性直接关系到整个消息系统的性能和可靠性。 2.12.1 消息路由的核心组件 要理解消息路由,首先需要了解 RabbitMQ 中参与路由的关键组件: 生产者 (Producer): 消息的发送者,负责创建和发布消息。生产者在发送消息时,会将消息发送到指定的 交换机 (Exchange)。

2.12 消息路由流程 (Message Routing)

2.12 消息路由流程 (Message Routing) 详解

在 RabbitMQ 消息队列系统中,消息路由 (Message Routing) 是核心概念之一,它决定了生产者 (Producer) 发送的消息如何被正确地投递到一个或多个队列 (Queue) 中,最终被消费者 (Consumer) 接收和处理。 消息路由的效率和灵活性直接关系到整个消息系统的性能和可靠性。

2.12.1 消息路由的核心组件

要理解消息路由,首先需要了解 RabbitMQ 中参与路由的关键组件:

  • 生产者 (Producer): 消息的发送者,负责创建和发布消息。生产者在发送消息时,会将消息发送到指定的 交换机 (Exchange)

  • 交换机 (Exchange): 消息的接收和路由中心。生产者发布的消息首先到达交换机,交换机根据预设的路由规则将消息路由到一个或多个队列。交换机本身不存储消息,它的主要职责是接收消息并根据路由规则分发消息。

  • 队列 (Queue): 消息的存储容器。队列接收来自交换机路由的消息,并等待消费者来消费。消息在队列中按照先进先出 (FIFO) 的原则进行存储和消费。

  • 绑定 (Binding): 交换机和队列之间的关联关系。绑定定义了交换机如何将消息路由到队列的规则。一个交换机可以与多个队列绑定,一个队列也可以与多个交换机绑定 (尽管不常见)。绑定由 路由键 (Routing Key) 和可选的 参数 (Arguments) 组成,路由键是路由规则的核心。

  • 路由键 (Routing Key): 生产者在发送消息时附加的一个属性,用于指示消息的路由目标。交换机根据路由键和绑定规则来决定消息应该被路由到哪些队列。

  • 消费者 (Consumer): 消息的接收者,从队列中获取消息并进行处理。消费者订阅队列,并接收队列中的消息。

可以用下图来表示这些组件之间的关系:

2.12.2 消息路由流程详解

消息路由流程的核心在于交换机如何根据消息的属性 (主要是路由键和交换机类型) 以及预定义的绑定规则,将消息投递到正确的队列。 详细的路由流程如下:

  1. 生产者发布消息: 生产者构建消息,并指定要发送到的 交换机名称 (Exchange Name)路由键 (Routing Key) (可选,取决于交换机类型)。然后将消息发送给 RabbitMQ 服务器。

  2. 交换机接收消息: RabbitMQ 服务器接收到生产者发送的消息,并将其投递到指定的交换机。

  3. 交换机根据类型和绑定规则路由消息: 交换机根据自身的 类型 (Type) 和已配置的 绑定 (Bindings) 规则,以及消息的 路由键 (Routing Key) (如果适用),来判断消息应该被路由到哪些队列。 RabbitMQ 提供了四种主要的交换机类型,每种类型有不同的路由策略:

    • Direct Exchange (直连交换机): 根据消息的路由键 完全匹配 绑定键的队列。

    • Fanout Exchange (扇形交换机): 将消息 广播 到所有绑定到该交换机的队列,忽略路由键。

    • Topic Exchange (主题交换机): 根据消息的路由键和绑定键的 模式匹配 规则,将消息路由到匹配的队列。支持通配符匹配。

    • Headers Exchange (headers交换机): 根据消息的 headers 属性进行路由,而非路由键。可以进行更复杂的属性匹配。

  4. 消息投递到队列: 交换机根据路由规则,将消息投递到一个或多个匹配的队列中。

  5. 消费者消费消息: 订阅了相关队列的消费者,从队列中接收并处理消息。

如果消息无法被路由到任何队列 (例如,没有匹配的绑定规则),则消息会被 丢弃 (默认行为) 或根据交换机的配置进行处理 (例如,发送到死信队列)。

2.12.3 交换机类型详解及代码实践

下面我们详细介绍四种交换机类型,并提供代码示例进行演示。 代码示例将使用 Python 的 pika 客户端库。

2.12.3.1 Direct Exchange (直连交换机)

路由规则: Direct Exchange 的路由规则非常简单直接。 它会将消息的路由键与绑定键进行 完全匹配。 如果消息的路由键与某个绑定的绑定键完全一致,则消息会被路由到该绑定对应的队列。

适用场景: 点对点通信,任务分发 (根据任务类型路由到不同的队列)。

Mermaid 图示:

Python 代码实践:

生产者 (direct_producer.py):

import pika # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Direct Exchange exchange_name = 'direct_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='direct') # 定义路由键 routing_keys = ['info', 'warning', 'error'] # 发送不同路由键的消息 for routing_key in routing_keys: message = f"This is a {routing_key} message." channel.basic_publish(exchange=exchange_name, routing_key=routing_key, body=message.encode()) print(f" [x] Sent '{routing_key}':'{message}'") connection.close()

消费者 (direct_consumer_info.py - 接收 'info' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'direct_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='direct') # 声明队列 (临时队列,exclusive=True, auto_delete=True) queue_name = 'info_queue' channel.queue_declare(queue=queue_name) # 绑定队列到 Direct Exchange,指定路由键 'info' channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='info') print(' [*] Waiting for messages with routing key "info". To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{method.routing_key}':'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

消费者 (direct_consumer_warning.py - 接收 'warning' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'direct_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='direct') queue_name = 'warning_queue' channel.queue_declare(queue=queue_name) # 绑定队列到 Direct Exchange,指定路由键 'warning' channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='warning') print(' [*] Waiting for messages with routing key "warning". To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{method.routing_key}':'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

消费者 (direct_consumer_error.py - 接收 'error' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'direct_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='direct') queue_name = 'error_queue' channel.queue_declare(queue=queue_name) # 绑定队列到 Direct Exchange,指定路由键 'error' channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='error') print(' [*] Waiting for messages with routing key "error". To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{method.routing_key}':'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

运行步骤:

  1. 分别运行 direct_consumer_info.py, direct_consumer_warning.py, direct_consumer_error.py 三个消费者脚本。

  2. 运行 direct_producer.py 生产者脚本。

预期结果:

  • direct_consumer_info.py 将接收并打印路由键为 "info" 的消息。

  • direct_consumer_warning.py 将接收并打印路由键为 "warning" 的消息。

  • direct_consumer_error.py 将接收并打印路由键为 "error" 的消息。

  • 每个消费者只接收到与其绑定路由键相匹配的消息,体现了 Direct Exchange 的精确匹配路由特性。

2.12.3.2 Fanout Exchange (扇形交换机)

路由规则: Fanout Exchange 会将接收到的消息 广播所有 绑定到该交换机的队列,忽略消息的路由键。 无论消息的路由键是什么,所有绑定到 Fanout Exchange 的队列都会收到一份消息的副本。

适用场景: 广播消息,例如:发布/订阅模式,群发通知、日志广播等。

Mermaid 图示:

Python 代码实践:

生产者 (fanout_producer.py):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Fanout Exchange exchange_name = 'fanout_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='fanout') message = "Hello everyone! This is a broadcast message." # 发送消息到 Fanout Exchange,路由键会被忽略 channel.basic_publish(exchange=exchange_name, routing_key='', body=message.encode()) print(f" [x] Sent '{message}'") connection.close()

消费者 (fanout_consumer_1.py):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'fanout_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='fanout') # 声明队列 (临时队列) queue_name = 'fanout_queue_1' channel.queue_declare(queue=queue_name) # 绑定队列到 Fanout Exchange (路由键为空字符串,Fanout Exchange 忽略路由键) channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='') print(' [*] Waiting for messages for fanout queue 1. To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

消费者 (fanout_consumer_2.py):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'fanout_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='fanout') queue_name = 'fanout_queue_2' channel.queue_declare(queue=queue_name) # 绑定队列到 Fanout Exchange channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='') print(' [*] Waiting for messages for fanout queue 2. To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

运行步骤:

  1. 分别运行 fanout_consumer_1.py, fanout_consumer_2.py 两个消费者脚本。

  2. 运行 fanout_producer.py 生产者脚本。

预期结果:

  • fanout_consumer_1.pyfanout_consumer_2.py 都会接收到相同的广播消息。

  • 无论生产者发送消息时指定的路由键是什么 (本例中为空字符串),所有绑定到 Fanout Exchange 的队列都会收到消息,体现了 Fanout Exchange 的广播特性。

2.12.3.3 Topic Exchange (主题交换机)

路由规则: Topic Exchange 使用 模式匹配 的方式进行路由。 消息的路由键和队列绑定的绑定键都支持使用通配符:

  • * (星号): 匹配 一个 单词。

  • # (井号): 匹配 零个或多个 单词。

  • 单词之间用 . (点号) 分隔。

适用场景: 根据主题或模式进行消息路由,实现更灵活的消息订阅和过滤。 例如:订阅特定类型的日志 (info., error.#),订阅特定区域的新闻 (news.us., news.europe.#)。

Mermaid 图示:

Python 代码实践:

生产者 (topic_producer.py):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Topic Exchange exchange_name = 'topic_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='topic') # 定义路由键和消息 routing_keys_messages = { 'quick.orange.rabbit': 'Quick orange rabbit.', 'lazy.orange.elephant': 'Lazy orange elephant.', 'quick.brown.fox': 'Quick brown fox.', 'lazy.pink.rabbit': 'Lazy pink rabbit.', 'lazy.brown.fox': 'Lazy brown fox.', 'anonymous.orange.rabbit': 'Anonymous orange rabbit.', 'lazy.yellow.rabbit': 'Lazy yellow rabbit.', } # 发送不同路由键的消息 for routing_key, message in routing_keys_messages.items(): channel.basic_publish(exchange=exchange_name, routing_key=routing_key, body=message.encode()) print(f" [x] Sent '{routing_key}':'{message}'") connection.close()

消费者 (topic_consumer_1.py - 订阅 'order.created.*' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'topic_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='topic') queue_name = 'topic_queue_1' channel.queue_declare(queue=queue_name) # 绑定队列到 Topic Exchange,使用模式 'order.created.*' channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='order.created.*') print(' [*] Waiting for messages with routing key "order.created.*". To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{method.routing_key}':'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

消费者 (topic_consumer_2.py - 订阅 'order.#' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'topic_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='topic') queue_name = 'topic_queue_2' channel.queue_declare(queue=queue_name) # 绑定队列到 Topic Exchange,使用模式 'order.#' channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='order.#') print(' [*] Waiting for messages with routing key "order.#". To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{method.routing_key}':'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

消费者 (topic_consumer_3.py - 订阅 '.payment.' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'topic_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='topic') queue_name = 'topic_queue_3' channel.queue_declare(queue=queue_name) # 绑定队列到 Topic Exchange,使用模式 '*.payment.*' channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='*.payment.*') print(' [*] Waiting for messages with routing key "*.payment.*". To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received '{method.routing_key}':'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

运行步骤:

  1. 分别运行 topic_consumer_1.py, topic_consumer_2.py, topic_consumer_3.py 三个消费者脚本。

  2. 运行 topic_producer.py 生产者脚本。

预期结果:

  • topic_consumer_1.py 将接收到路由键匹配 order.created.* 模式的消息。

  • topic_consumer_2.py 将接收到路由键匹配 order.# 模式的消息 (包括 order.created.* 的消息,因为 # 匹配零个或多个单词)。

  • topic_consumer_3.py 将接收到路由键匹配 *.payment.* 模式的消息。

  • 不同消费者根据绑定的模式接收到不同的消息,体现了 Topic Exchange 的灵活路由特性。

2.12.3.4 Headers Exchange (headers交换机)

路由规则: Headers Exchange 不使用路由键进行路由,而是使用消息的 headers 属性 进行匹配。 绑定时,可以指定一组 headers 键值对作为匹配规则。

  • 可以指定 x-match 属性来控制匹配模式:

    • x-match='any' (默认): 消息 headers 中 任意一个 header 键值对与绑定规则匹配即可。

    • x-match='all': 消息 headers 中 所有 header 键值对都必须与绑定规则匹配。

适用场景: 基于消息属性进行路由,实现更复杂的路由逻辑,例如:根据消息类型、优先级、来源等属性进行路由。

Mermaid 图示:

Python 代码实践:

生产者 (headers_producer.py):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Headers Exchange exchange_name = 'headers_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='headers') # 定义消息 headers message_headers_list = [ {'format': 'pdf', 'type': 'report'}, {'format': 'excel', 'type': 'report'}, {'format': 'pdf', 'category': 'document'}, {'type': 'log', 'level': 'error'} ] # 发送不同 headers 的消息 for headers in message_headers_list: message = f"Message with headers: {headers}" properties = pika.BasicProperties(headers=headers) channel.basic_publish(exchange=exchange_name, routing_key='', body=message.encode(), properties=properties) print(f" [x] Sent '{headers}':'{message}'") connection.close()

消费者 (headers_consumer_pdf_reports.py - 订阅 format='pdf' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'headers_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='headers') queue_name = 'headers_queue_pdf_reports' channel.queue_declare(queue=queue_name) # 绑定队列到 Headers Exchange,匹配 header 'format=pdf' (默认 x-match='any') channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='', arguments={'format': 'pdf'}) print(' [*] Waiting for messages with header "format=pdf". To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received headers:{properties.headers} body:'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

消费者 (headers_consumer_all_reports.py - 订阅 type='report' 且 format='pdf' 消息):

import pika connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() exchange_name = 'headers_exchange_example' channel.exchange_declare(exchange=exchange_name, exchange_type='headers') queue_name = 'headers_queue_all_reports' channel.queue_declare(queue=queue_name) # 绑定队列到 Headers Exchange,匹配 header 'type=report' 和 'format=pdf' (x-match='all') channel.queue_bind(exchange=exchange_name, queue=queue_name, routing_key='', arguments={'type': 'report', 'format': 'pdf', 'x-match': 'all'}) print(' [*] Waiting for messages with headers "type=report" and "format=pdf" (x-match="all"). To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] Received headers:{properties.headers} body:'{body.decode()}'") channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) channel.start_consuming()

运行步骤:

  1. 分别运行 headers_consumer_pdf_reports.py, headers_consumer_all_reports.py 两个消费者脚本。

  2. 运行 headers_producer.py 生产者脚本。

预期结果:

  • headers_consumer_pdf_reports.py 将接收到 headers 中包含 format='pdf' 的消息 (无论 x-matchany 还是 all,因为默认是 any)。

  • headers_consumer_all_reports.py 将只接收到 headers 中 同时 包含 type='report'format='pdf' 的消息 (因为指定了 x-match='all')。

  • 体现了 Headers Exchange 基于消息 headers 属性进行路由的特性,以及 x-match 属性对匹配模式的影响。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U