5.1 消息模式 (Message Patterns)


文档摘要

5.1 消息模式 (Message Patterns) RabbitMQ 高级主题:5.1 消息模式 (Message Patterns) 在构建复杂的分布式系统时,仅仅依靠简单的点对点消息传递是不够的。RabbitMQ 提供了多种消息模式,允许开发者构建更灵活、可扩展和容错的应用程序。 这些模式定义了消息如何在生产者、消费者和交换机之间流动,从而解决了各种常见的消息传递问题。 发布/订阅 (Publish/Subscribe) 概念: 发布/订阅模式是一种消息传递范例,其中消息生产者(发布者)将消息发送到交换机,而多个消息消费者(订阅者)可以订阅该交换机,并接收所有发送到该交换机的消息。 这种模式实现了生产者和消费者之间的解耦,允许消费者独立地接收消息,而无需知道消息的来源。

5.1 消息模式 (Message Patterns)

RabbitMQ 高级主题:5.1 消息模式 (Message Patterns)

在构建复杂的分布式系统时,仅仅依靠简单的点对点消息传递是不够的。RabbitMQ 提供了多种消息模式,允许开发者构建更灵活、可扩展和容错的应用程序。 这些模式定义了消息如何在生产者、消费者和交换机之间流动,从而解决了各种常见的消息传递问题。

1. 发布/订阅 (Publish/Subscribe)

概念:

发布/订阅模式是一种消息传递范例,其中消息生产者(发布者)将消息发送到交换机,而多个消息消费者(订阅者)可以订阅该交换机,并接收所有发送到该交换机的消息。 这种模式实现了生产者和消费者之间的解耦,允许消费者独立地接收消息,而无需知道消息的来源。

RabbitMQ 实现:

在 RabbitMQ 中,发布/订阅模式通常使用 fanout 类型的交换机来实现。 fanout 交换机将所有接收到的消息广播到所有绑定到它的队列。

适用场景:

  • 广播通知:例如,向所有订阅者发送系统更新通知。

  • 日志聚合:将来自多个应用程序的日志消息发送到中央日志服务器。

  • 实时数据流:将实时数据(例如股票行情或传感器数据)分发给多个消费者。

代码实践 (Python):

import pika import time # 连接RabbitMQ服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明fanout类型的交换机 channel.exchange_declare(exchange='logs', exchange_type='fanout') # 发布消息 def publish_message(message): channel.basic_publish(exchange='logs', routing_key='', body=message) print(f" [x] Sent {message}") # 模拟发布消息 for i in range(5): message = f"Log message {i}" publish_message(message) time.sleep(1) connection.close()
import pika # 连接RabbitMQ服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明fanout类型的交换机 channel.exchange_declare(exchange='logs', exchange_type='fanout') # 创建一个临时队列 result = channel.queue_declare(queue='', exclusive=True) queue_name = result.method.queue # 将队列绑定到交换机 channel.queue_bind(exchange='logs', queue=queue_name) print(' [*] Waiting for logs. 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()

详解:

  • 生产者 (Publish): 生产者连接到 RabbitMQ 服务器,声明一个名为 logsfanout 类型的交换机。 然后,它将消息发布到该交换机,routing_key 为空字符串,因为 fanout 交换机忽略 routing_key

  • 消费者 (Subscribe): 消费者也连接到 RabbitMQ 服务器,声明相同的 fanout 类型的交换机。 重要的是,消费者创建一个临时队列 (exclusive=True),并在连接断开时自动删除该队列。 然后,消费者将该临时队列绑定到 logs 交换机。 最后,消费者开始消费队列中的消息,并使用 callback 函数处理接收到的消息。

Graph TD 图:

2. 路由 (Routing)

概念:

路由模式允许消息生产者将消息发送到交换机,并根据消息的路由键将消息路由到特定的队列。 消费者可以订阅具有特定路由键的队列,从而只接收他们感兴趣的消息。

RabbitMQ 实现:

在 RabbitMQ 中,路由模式通常使用 direct 类型的交换机来实现。 direct 交换机将消息路由到路由键完全匹配的队列。

适用场景:

  • 优先级消息处理:根据消息的优先级(例如,errorwarninginfo)将消息路由到不同的队列。

  • 特定用户通知:根据用户的 ID 将通知消息路由到特定的用户队列。

  • 事件驱动架构:根据事件类型将事件消息路由到相应的处理程序。

代码实践 (Python):

import pika import sys # 连接RabbitMQ服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明direct类型的交换机 channel.exchange_declare(exchange='direct_logs', exchange_type='direct') # 获取路由键 severity = sys.argv[1] if len(sys.argv) > 1 else 'info' message = ' '.join(sys.argv[2:]) or 'Hello World!' # 发布消息 channel.basic_publish(exchange='direct_logs', routing_key=severity, body=message) print(f" [x] Sent {severity}:{message}") connection.close()
import pika import sys # 连接RabbitMQ服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明direct类型的交换机 channel.exchange_declare(exchange='direct_logs', exchange_type='direct') # 获取路由键 severity = sys.argv[1] if len(sys.argv) > 1 else 'info' # 创建队列 result = channel.queue_declare(queue='', exclusive=True) queue_name = result.method.queue # 将队列绑定到交换机,并指定路由键 channel.queue_bind(exchange='direct_logs', queue=queue_name, routing_key=severity) print(f' [*] Waiting for logs with severity: {severity}. 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()

详解:

  • 生产者 (Publish): 生产者连接到 RabbitMQ 服务器,声明一个名为 direct_logsdirect 类型的交换机。 它从命令行参数获取路由键(severity)和消息内容。 然后,它将消息发布到交换机,并指定路由键。

  • 消费者 (Subscribe): 消费者也连接到 RabbitMQ 服务器,声明相同的 direct 类型的交换机。 它从命令行参数获取路由键(severity),并创建一个临时队列。 然后,消费者将该临时队列绑定到 direct_logs 交换机,并指定路由键。 这意味着该队列只会接收路由键与指定 severity 匹配的消息。 最后,消费者开始消费队列中的消息,并使用 callback 函数处理接收到的消息。

Graph TD 图:

3. 主题 (Topics)

概念:

主题模式是路由模式的更高级版本,它允许使用通配符进行更灵活的路由。 消息生产者将消息发送到交换机,并根据消息的主题将消息路由到匹配的队列。 消费者可以订阅具有特定主题模式的队列,从而接收与该模式匹配的消息。

RabbitMQ 实现:

在 RabbitMQ 中,主题模式使用 topic 类型的交换机来实现。 topic 交换机使用以下通配符:

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

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

适用场景:

  • 复杂的事件路由:根据事件的多个属性(例如,区域、类型、严重性)将事件消息路由到相应的处理程序。

  • 灵活的日志过滤:根据日志消息的多个属性(例如,模块、级别、主机)将日志消息路由到不同的队列。

  • 服务集成:将来自不同服务的消息路由到相应的消费者,而无需修改生产者代码。

代码实践 (Python):

import pika import sys # 连接RabbitMQ服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明topic类型的交换机 channel.exchange_declare(exchange='topic_logs', exchange_type='topic') # 获取路由键 routing_key = sys.argv[1] if len(sys.argv) > 1 else 'anonymous.info' message = ' '.join(sys.argv[2:]) or 'Hello World!' # 发布消息 channel.basic_publish(exchange='topic_logs', routing_key=routing_key, body=message) print(f" [x] Sent {routing_key}:{message}") connection.close()
import pika import sys # 连接RabbitMQ服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明topic类型的交换机 channel.exchange_declare(exchange='topic_logs', exchange_type='topic') # 获取绑定键 binding_key = sys.argv[1] if len(sys.argv) > 1 else '#' # 创建队列 result = channel.queue_declare(queue='', exclusive=True) queue_name = result.method.queue # 将队列绑定到交换机,并指定绑定键 channel.queue_bind(exchange='topic_logs', queue=queue_name, routing_key=binding_key) print(f' [*] Waiting for logs with binding key: {binding_key}. 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()

详解:

  • 生产者 (Publish): 生产者连接到 RabbitMQ 服务器,声明一个名为 topic_logstopic 类型的交换机。 它从命令行参数获取路由键(routing_key)和消息内容。 然后,它将消息发布到交换机,并指定路由键。

  • 消费者 (Subscribe): 消费者也连接到 RabbitMQ 服务器,声明相同的 topic_logs 类型的交换机。 它从命令行参数获取绑定键(binding_key),并创建一个临时队列。 然后,消费者将该临时队列绑定到 topic_logs 交换机,并指定绑定键。 这意味着该队列只会接收路由键与指定 binding_key 模式匹配的消息。 最后,消费者开始消费队列中的消息,并使用 callback 函数处理接收到的消息。

示例:

  • kern.*: 接收所有 kern 开头的消息。

  • *.critical: 接收所有以 critical 结尾的消息。

  • kern.#: 接收所有以 kern. 开头的消息,包括 kern.info.debug 等。

  • #: 接收所有消息。

Graph TD 图:

4. 请求/回复 (Request/Reply)

概念:

请求/回复模式允许客户端向服务器发送请求,并期望收到回复。 这是一种同步的消息传递模式,客户端在发送请求后会阻塞,直到收到回复。

RabbitMQ 实现:

在 RabbitMQ 中,请求/回复模式通常使用以下机制实现:

  1. reply_to 属性: 客户端在发送请求时,将一个唯一的队列名称作为 reply_to 属性添加到消息中。

  2. correlation_id 属性: 客户端在发送请求时,生成一个唯一的 correlation_id,并将其添加到消息中。

  3. 服务器回复: 服务器在收到请求后,处理请求并将回复发送到 reply_to 属性指定的队列。 服务器还将原始消息的 correlation_id 复制到回复消息中。

  4. 客户端匹配: 客户端接收回复消息,并使用 correlation_id 将回复与原始请求进行匹配。

适用场景:

  • 远程过程调用 (RPC): 客户端调用服务器上的方法,并等待结果。

  • 同步数据查询:客户端向服务器发送查询请求,并等待返回数据。

  • 身份验证和授权:客户端向服务器发送身份验证请求,并等待验证结果。

代码实践 (Python):

import pika import uuid class FibonacciRpcClient(object): def __init__(self): self.connection = pika.BlockingConnection( pika.ConnectionParameters(host='localhost')) self.channel = self.connection.channel() result = self.channel.queue_declare(queue='', exclusive=True) self.callback_queue = result.method.queue self.channel.basic_consume( queue=self.callback_queue, on_message_callback=self.on_response, auto_ack=True) self.response = None self.corr_id = None def on_response(self, ch, method, props, body): if self.corr_id == props.correlation_id: self.response = body.decode() def call(self, n): self.response = None self.corr_id = str(uuid.uuid4()) self.channel.basic_publish( exchange='', routing_key='rpc_queue', properties=pika.BasicProperties( reply_to=self.callback_queue, correlation_id=self.corr_id, ), body=str(n)) self.connection.process_data_events(time_limit=None) while self.response is None: self.connection.process_data_events() return int(self.response) fibonacci_rpc = FibonacciRpcClient() print(" [x] Requesting fib(30)") response = fibonacci_rpc.call(30) print(" [.] Got %r" % response)
import pika import time connection = pika.BlockingConnection( pika.ConnectionParameters(host='localhost')) channel = connection.channel() channel.queue_declare(queue='rpc_queue') def fib(n): if n == 0: return 0 elif n == 1: return 1 else: return fib(n - 1) + fib(n - 2) def on_request(ch, method, props, body): n = int(body) print(" [.] fib(%s)" % n) response = fib(n) ch.basic_publish(exchange='', routing_key=props.reply_to, properties=pika.BasicProperties(correlation_id = \ props.correlation_id), body=str(response)) ch.basic_ack(delivery_tag=method.delivery_tag) channel.basic_qos(prefetch_count=1) channel.basic_consume(queue='rpc_queue', on_message_callback=on_request) print(" [x] Awaiting RPC requests") channel.start_consuming()

详解:

  • 客户端 (Request): 客户端创建一个临时队列用于接收回复,并生成一个唯一的 correlation_id。 它将请求消息发送到 rpc_queue 队列,并将 reply_to 设置为临时队列的名称,并将 correlation_id 设置为生成的 ID。 客户端然后阻塞,直到收到具有匹配 correlation_id 的回复消息。

  • 服务器 (Reply): 服务器监听 rpc_queue 队列上的请求。 当收到请求时,它执行计算并将回复消息发送到 reply_to 属性指定的队列。 服务器还将原始消息的 correlation_id 复制到回复消息中。

Graph TD 图:

5. 工作队列 (Work Queues)

概念:

工作队列 (又称任务队列) 的目的是为了避免立即执行资源密集型任务,而不得不等待它完成。 我们可以把任务封装成消息并发送到队列中。 在后台运行的工作进程会从队列中取出任务并最终执行这些任务。

RabbitMQ 实现:

在 RabbitMQ 中,工作队列通常使用默认的交换机和队列来实现。 生产者将任务消息发送到队列,而多个工作进程(消费者)可以从队列中获取任务并执行。

适用场景:

  • 后台任务处理:例如,处理图像、视频或生成报告。

  • 批量处理:将大量数据分割成小任务,并分发给多个工作进程并行处理。

  • 延迟任务:将任务延迟到稍后的时间执行。

代码实践 (Python):

import pika import time import sys connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost')) channel = connection.channel() channel.queue_declare(queue='task_queue', durable=True) message = ' '.join(sys.argv[1:]) or "Hello World!" channel.basic_publish(exchange='', routing_key='task_queue', body=message, properties=pika.BasicProperties( delivery_mode=2, # make message persistent )) print(" [x] Sent %r" % message) connection.close()
import pika import time connection = pika.BlockingConnection(pika.ConnectionParameters(host='localhost')) channel = connection.channel() channel.queue_declare(queue='task_queue', durable=True) print(' [*] Waiting for messages. To exit press CTRL+C') def callback(ch, method, properties, body): print(" [x] Received %r" % body.decode()) time.sleep(body.count(b'.')) print(" [x] Done") ch.basic_ack(delivery_tag=method.delivery_tag) channel.basic_qos(prefetch_count=1) channel.basic_consume(queue='task_queue', on_message_callback=callback) channel.start_consuming()

详解:

  • 生产者 (Publish): 生产者连接到 RabbitMQ 服务器,声明一个名为 task_queue 的队列,并设置 durable=True 以确保队列在服务器重启后仍然存在。 它将任务消息发送到该队列,并设置 delivery_mode=2 以确保消息在服务器重启后仍然存在。

  • 消费者 (Subscribe): 消费者也连接到 RabbitMQ 服务器,声明相同的 task_queue 队列,并设置 durable=True。 它设置 prefetch_count=1 以确保每个消费者一次只接收一个任务,避免某些消费者过载。 然后,消费者开始消费队列中的消息,并使用 callback 函数处理接收到的消息。 在 callback 函数中,消费者模拟任务处理,并使用 ch.basic_ack(delivery_tag=method.delivery_tag) 确认消息已成功处理。

Graph TD 图:

这些只是 RabbitMQ 中一些常用的消息模式。 通过组合和定制这些模式,您可以构建满足各种需求的复杂消息传递解决方案。 在选择消息模式时,请考虑您的应用程序的具体需求,例如可靠性、可扩展性、性能和复杂性。 记住,在实际应用中,这些模式经常被组合使用,以满足更复杂的需求。 深入理解这些模式是构建健壮、可扩展的分布式系统的关键。


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