6. RabbitMQ 实践应用


文档摘要

RabbitMQ 实践应用 RabbitMQ 实践应用详解 引言 在现代分布式系统中,消息队列扮演着至关重要的角色。它们允许不同的服务和应用以异步、解耦的方式进行通信,从而提高系统的可靠性、可伸缩性和灵活性。RabbitMQ 作为一款开源的消息队列中间件,凭借其强大的功能、稳定的性能和易用性,成为了众多企业和开发者的首选。 1. 任务队列 (Task Queues) 1.1 应用场景 任务队列是消息队列最经典的应用场景之一。它主要用于处理耗时、异步的任务,例如: 用户注册后的邮件发送、短信通知: 用户注册成功后,发送欢迎邮件或短信通知并不需要立即完成,可以放入任务队列异步处理,避免阻塞用户注册流程。

6. RabbitMQ 实践应用

RabbitMQ 实践应用详解

引言

在现代分布式系统中,消息队列扮演着至关重要的角色。它们允许不同的服务和应用以异步、解耦的方式进行通信,从而提高系统的可靠性、可伸缩性和灵活性。RabbitMQ 作为一款开源的消息队列中间件,凭借其强大的功能、稳定的性能和易用性,成为了众多企业和开发者的首选。

1. 任务队列 (Task Queues)

1.1 应用场景

任务队列是消息队列最经典的应用场景之一。它主要用于处理耗时、异步的任务,例如:

  • 用户注册后的邮件发送、短信通知: 用户注册成功后,发送欢迎邮件或短信通知并不需要立即完成,可以放入任务队列异步处理,避免阻塞用户注册流程。

  • 图片/视频处理: 上传图片或视频后,进行压缩、转码、添加水印等操作通常比较耗时,可以放入任务队列后台处理,提高用户上传体验。

  • 数据分析与报表生成: 复杂的统计分析、报表生成等任务,可以放入任务队列异步处理,避免影响主业务系统的性能。

  • 订单处理: 订单支付成功后,后续的库存扣减、物流信息更新、积分增加等操作,可以放入任务队列异步处理,提高订单处理效率。

1.2 代码实践 (Python + pika)

以下代码示例将演示如何使用 RabbitMQ 创建一个简单的任务队列,用于模拟耗时任务的处理。

生产者 (task_producer.py):

#!/usr/bin/env python import pika import sys import time # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明队列 (如果队列不存在则创建) channel.queue_declare(queue='task_queue', durable=True) # 声明持久化队列 message = ' '.join(sys.argv[1:]) or "Hello World!" # 设置消息持久化 properties = pika.BasicProperties(delivery_mode=2,) # make message persistent channel.basic_publish(exchange='', routing_key='task_queue', body=message, properties=properties) # 发送消息并设置持久化 print(f" [x] Sent '{message}'") connection.close()

消费者 (task_consumer.py):

#!/usr/bin/env python import pika import time # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('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(f" [x] Received {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()

1.3 代码详解

  • 生产者 (task_producer.py):

    • 使用 pika.BlockingConnection 连接到 RabbitMQ 服务器 (默认连接本地 localhost)。

    • 创建信道 (channel),用于进行 AMQP 协议的操作。

    • 使用 channel.queue_declare(queue='task_queue', durable=True) 声明一个名为 task_queue 的队列。durable=True 参数表示队列是持久化的,即使 RabbitMQ 服务器重启,队列也不会丢失。

    • 从命令行参数获取消息内容,如果没有参数则使用默认消息 "Hello World!"。

    • 创建 pika.BasicProperties(delivery_mode=2,) 对象,设置消息的 delivery_mode2,表示消息需要持久化,即使 RabbitMQ 服务器重启,消息也不会丢失。

    • 使用 channel.basic_publish 发布消息到默认交换机 (direct exchange),路由键为 task_queue,消息体为 message,并设置消息属性为 properties (包含持久化设置)。

    • 打印消息发送成功的提示信息。

    • 关闭连接。

  • 消费者 (task_consumer.py):

    • 连接 RabbitMQ 服务器并创建信道。

    • 同样使用 channel.queue_declare(queue='task_queue', durable=True) 声明队列,确保队列存在。

    • 打印等待消息的提示信息。

    • 定义回调函数 callback(ch, method, properties, body),当消费者接收到消息时,该函数会被调用:

      • 打印接收到的消息内容。

      • 使用 time.sleep(body.count(b'.')) 模拟耗时任务,根据消息内容中 '.' 的数量决定休眠时间。

      • 打印任务完成的提示信息。

      • 关键步骤: 使用 ch.basic_ack(delivery_tag = method.delivery_tag) 发送消息确认 (acknowledgement)。RabbitMQ 需要接收到消费者的确认信息后,才会认为消息已被成功处理并从队列中移除。如果消费者在处理消息过程中崩溃或未发送确认,RabbitMQ 会将消息重新放回队列,等待其他消费者处理,从而保证消息的可靠 delivery。

    • 使用 channel.basic_qos(prefetch_count=1) 设置预取计数为 1。这表示消费者一次最多接收一条消息,只有在确认处理完当前消息后,RabbitMQ 才会发送下一条消息。这可以防止消费者负载过重,提高系统的公平性和稳定性。

    • 使用 channel.basic_consume(queue='task_queue', on_message_callback=callback) 注册消费者,监听 task_queue 队列,并将接收到的消息传递给 callback 函数处理。

    • 使用 channel.start_consuming() 开始消费消息,消费者程序会一直运行,等待接收消息。

1.4 运行示例

  1. 启动消费者: 在终端中运行 python task_consumer.py

  2. 启动生产者并发送消息: 在另一个终端中运行 python task_producer.py Hello World .... (可以尝试添加不同数量的 '.' 来模拟不同耗时任务)。

您会看到消费者终端打印接收到的消息,并模拟耗时任务后打印 "Done"。多次运行生产者发送不同消息,消费者会依次处理这些任务。

1.5 Mermaid 图表

图表解释:

  • Producer (生产者): 负责将任务消息发送到 Task Queue。

  • Task Queue (任务队列): RabbitMQ 中的队列 task_queue,用于存储待处理的任务消息。

  • Consumer (消费者): 负责从 Task Queue 中接收任务消息并进行处理。

  • Ack (确认): Consumer 在完成任务处理后,向 RabbitMQ 发送确认消息,告知消息已被成功处理。

2. 发布/订阅模式 (Publish/Subscribe)

2.1 应用场景

发布/订阅模式 (Pub/Sub) 是一种消息传递模式,其中消息生产者 (Publisher) 将消息发布到交换机 (Exchange),而多个消息消费者 (Subscribers) 可以订阅感兴趣的消息,并从队列中接收消息。RabbitMQ 通过 Fanout Exchange 类型来实现 Pub/Sub 模式。常见的应用场景包括:

  • 实时消息广播: 例如,股票行情更新、新闻推送、聊天室消息广播等,需要将消息实时推送给所有订阅者。

  • 事件通知: 例如,系统状态变更、配置更新等事件发生时,需要通知多个关注该事件的服务或模块。

  • 日志收集: 将多个服务产生的日志消息发布到 Fanout Exchange,多个日志收集服务可以订阅并进行日志聚合分析。

2.2 代码实践 (Python + pika)

以下代码示例将演示如何使用 RabbitMQ 实现一个简单的发布/订阅模式,模拟消息广播。

发布者 (publisher.py):

#!/usr/bin/env python import pika import sys # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Fanout Exchange channel.exchange_declare(exchange='logs', exchange_type='fanout') message = ' '.join(sys.argv[1:]) or "info: Hello World!" channel.basic_publish(exchange='logs', routing_key='', # Fanout Exchange 忽略 routing_key body=message) print(f" [x] Sent {message}") connection.close()

订阅者 (subscriber.py):

#!/usr/bin/env python import pika # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Fanout Exchange (订阅者也需要声明,确保 Exchange 存在) channel.exchange_declare(exchange='logs', exchange_type='fanout') # 创建临时队列 (exclusive=True, auto_delete=True) result = channel.queue_declare(queue='', exclusive=True, auto_delete=True) queue_name = result.method.queue # 绑定队列到 Fanout Exchange 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] {body.decode()}") channel.basic_consume(queue=queue_name, auto_ack=True, on_message_callback=callback) channel.start_consuming()

2.3 代码详解

  • 发布者 (publisher.py):

    • 连接 RabbitMQ 服务器并创建信道。

    • 使用 channel.exchange_declare(exchange='logs', exchange_type='fanout') 声明一个名为 logs 的 Fanout Exchange。exchange_type='fanout' 指定交换机类型为 Fanout。

    • 从命令行参数获取消息内容,如果没有参数则使用默认消息 "info: Hello World!"。

    • 使用 channel.basic_publish 发布消息到 logs Exchange,routing_key 设置为空字符串,因为 Fanout Exchange 会忽略 routing_key,将消息广播到所有绑定到该 Exchange 的队列。

    • 打印消息发送成功的提示信息。

    • 关闭连接。

  • 订阅者 (subscriber.py):

    • 连接 RabbitMQ 服务器并创建信道。

    • 同样使用 channel.exchange_declare(exchange='logs', exchange_type='fanout') 声明 logs Exchange。

    • 使用 channel.queue_declare(queue='', exclusive=True, auto_delete=True) 创建一个临时队列。

      • queue='' 表示由 RabbitMQ 自动生成一个随机队列名称。

      • exclusive=True 表示该队列是排他队列,只允许当前连接访问,当连接关闭时队列自动删除。

      • auto_delete=True 表示当最后一个消费者取消订阅后,队列自动删除。

    • 使用 channel.queue_bind(exchange='logs', queue=queue_name) 将临时队列绑定到 logs Exchange。这意味着发布到 logs Exchange 的消息都会被路由到该队列。

    • 打印等待日志消息的提示信息。

    • 定义回调函数 callback(ch, method, properties, body),当订阅者接收到消息时,该函数会被调用:

      • 打印接收到的消息内容。
    • 使用 channel.basic_consume(queue=queue_name, auto_ack=True, on_message_callback=callback) 注册消费者,监听临时队列 queue_name,并将接收到的消息传递给 callback 函数处理。auto_ack=True 表示自动消息确认,消费者接收到消息后会自动发送确认,适用于对消息可靠性要求不高的场景。

    • 使用 channel.start_consuming() 开始消费消息。

2.4 运行示例

  1. 启动多个订阅者: 在多个终端中分别运行 python subscriber.py,启动多个订阅者实例。

  2. 启动发布者并发送消息: 在另一个终端中运行 python publisher.py "This is a log message!"

您会看到所有订阅者终端都打印接收到的消息 "This is a log message!"。多次运行发布者发送不同消息,所有订阅者都会收到相同的消息,实现了消息广播。

2.5 Mermaid 图表

图表解释:

  • Publisher (发布者): 负责将消息发布到 Fanout Exchange。

  • Fanout Exchange (Fanout 交换机): 接收发布者发送的消息,并将消息广播到所有绑定到该 Exchange 的队列。

  • Queue 1, Queue 2 (队列 1, 队列 2): 多个临时队列,分别绑定到 Fanout Exchange,用于接收广播的消息。

  • Subscriber 1, Subscriber 2 (订阅者 1, 订阅者 2): 多个消费者,分别从各自的队列中接收消息。

3. 路由模式 (Routing)

3.1 应用场景

路由模式允许消息生产者将消息发送到 Direct Exchange,并根据消息的 Routing Key 将消息路由到特定的队列。消费者可以根据 Routing Key 订阅感兴趣的消息类型。常见的应用场景包括:

  • 不同类型的日志消息处理: 例如,将不同级别的日志消息 (info, warning, error) 通过不同的 Routing Key 发送到不同的队列,分别由不同的日志处理服务进行处理。

  • 订单系统中的订单类型路由: 例如,根据订单类型 (普通订单、VIP 订单、秒杀订单) 将订单消息路由到不同的队列,由不同的订单处理服务进行处理。

  • 微服务架构中的服务间通信: 不同的微服务可以通过 Direct Exchange 和 Routing Key 进行精确的消息传递。

3.2 代码实践 (Python + pika)

以下代码示例将演示如何使用 RabbitMQ 实现一个简单的路由模式,模拟不同类型的日志消息路由。

日志生产者 (routing_producer.py):

#!/usr/bin/env python import pika import sys # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Direct Exchange 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, # 使用 severity 作为 Routing Key body=message) print(f" [x] Sent {severity}:{message}") connection.close()

日志消费者 (routing_consumer.py):

#!/usr/bin/env python import pika import sys # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Direct Exchange (消费者也需要声明) channel.exchange_declare(exchange='direct_logs', exchange_type='direct') # 创建临时队列 result = channel.queue_declare(queue='', exclusive=True, auto_delete=True) queue_name = result.method.queue severities = sys.argv[1:] # 从命令行参数获取需要订阅的 severity 级别 if not severities: sys.stderr.write(f"Usage: {sys.argv[0]} [info] [warning] [error]\n") sys.exit(1) for severity in severities: # 绑定队列到 Direct Exchange,并指定 Routing Key channel.queue_bind(exchange='direct_logs', queue=queue_name, routing_key=severity) print(' [*] Waiting for logs. To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] {method.routing_key}:{body.decode()}") channel.basic_consume(queue=queue_name, auto_ack=True, on_message_callback=callback) channel.start_consuming()

3.3 代码详解

  • 日志生产者 (routing_producer.py):

    • 连接 RabbitMQ 服务器并创建信道。

    • 使用 channel.exchange_declare(exchange='direct_logs', exchange_type='direct') 声明一个名为 direct_logs 的 Direct Exchange。exchange_type='direct' 指定交换机类型为 Direct。

    • 从命令行参数获取日志级别 severity (例如:info, warning, error) 和消息内容。

    • 使用 channel.basic_publish 发布消息到 direct_logs Exchange,并将 severity 作为 routing_key

    • 打印消息发送成功的提示信息,包含 severity 和消息内容。

    • 关闭连接。

  • 日志消费者 (routing_consumer.py):

    • 连接 RabbitMQ 服务器并创建信道。

    • 同样声明 direct_logs Exchange。

    • 创建临时队列。

    • 从命令行参数获取需要订阅的 severities 列表 (例如:info, warning, error)。

    • 遍历 severities 列表,使用 channel.queue_bind(exchange='direct_logs', queue=queue_name, routing_key=severity) 将临时队列绑定到 direct_logs Exchange,并为每个 severity 指定不同的 routing_key。这意味着队列只会接收 Routing Key 与绑定时指定的 Routing Key 完全匹配的消息。

    • 打印等待日志消息的提示信息。

    • 定义回调函数 callback(ch, method, properties, body),打印接收到的消息,包含 Routing Key 和消息内容。

    • 注册消费者并开始消费消息。

3.4 运行示例

  1. 启动消费者并订阅不同级别的日志:

    • 终端 1: python routing_consumer.py info warning (订阅 info 和 warning 级别的日志)

    • 终端 2: python routing_consumer.py error (订阅 error 级别的日志)

    • 终端 3: python routing_consumer.py warning error (订阅 warning 和 error 级别的日志)

  2. 启动生产者并发送不同级别的日志消息:

    • python routing_producer.py info "This is an info log."

    • python routing_producer.py warning "This is a warning log."

    • python routing_producer.py error "This is an error log."

您会看到不同消费者终端只接收到与其订阅级别匹配的日志消息。例如,订阅 "info warning" 的消费者会收到 info 和 warning 级别的日志,但不会收到 error 级别的日志。

3.5 Mermaid 图表

图表解释:

  • Producer (生产者): 负责将消息和 Routing Key 发送到 Direct Exchange。

  • Direct Exchange (Direct 交换机): 根据消息的 Routing Key,将消息路由到 Routing Key 完全匹配的队列。

  • Queue 1, Queue 2 (队列 1, 队列 2): 不同的队列,分别绑定到 Direct Exchange,并指定不同的 Routing Key。

  • Consumer 1, Consumer 2 (消费者 1, 消费者 2): 不同的消费者,分别从各自的队列中接收消息,只接收 Routing Key 与其订阅匹配的消息。

4. 主题模式 (Topics)

4.1 应用场景

主题模式是路由模式的扩展,它使用 Topic Exchange,允许消费者使用 通配符 (wildcards) 订阅符合一定模式的消息。Topic Exchange 使用 Routing Key 的模式匹配来进行消息路由,比 Direct Exchange 更加灵活。常见的应用场景包括:

  • 更细粒度的日志消息路由: 例如,可以根据日志级别和模块进行更精细的路由,例如 kern.critical, auth.info, app.debug 等。

  • 复杂的事件路由: 例如,根据事件类型、事件来源等多个维度进行事件路由。

  • 多条件的消息订阅: 允许消费者订阅符合多种条件的消息。

4.2 代码实践 (Python + pika)

以下代码示例将演示如何使用 RabbitMQ 实现一个简单的主题模式,模拟使用通配符订阅日志消息。

主题日志生产者 (topic_producer.py):

#!/usr/bin/env python import pika import sys # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Topic Exchange 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, # 使用 routing_key body=message) print(f" [x] Sent {routing_key}:{message}") connection.close()

主题日志消费者 (topic_consumer.py):

#!/usr/bin/env python import pika import sys # 连接 RabbitMQ 服务器 connection = pika.BlockingConnection(pika.ConnectionParameters('localhost')) channel = connection.channel() # 声明 Topic Exchange (消费者也需要声明) channel.exchange_declare(exchange='topic_logs', exchange_type='topic') # 创建临时队列 result = channel.queue_declare(queue='', exclusive=True, auto_delete=True) queue_name = result.method.queue binding_keys = sys.argv[1:] # 从命令行参数获取 binding keys (订阅模式) if not binding_keys: sys.stderr.write(f"Usage: {sys.argv[0]} <binding_key>...\n") sys.exit(1) for binding_key in binding_keys: # 绑定队列到 Topic Exchange,并指定 binding key (订阅模式) channel.queue_bind(exchange='topic_logs', queue=queue_name, routing_key=binding_key) print(' [*] Waiting for logs. To exit press CTRL+C') def callback(ch, method, properties, body): print(f" [x] {method.routing_key}:{body.decode()}") channel.basic_consume(queue=queue_name, auto_ack=True, on_message_callback=callback) channel.start_consuming()

4.3 代码详解

  • 主题日志生产者 (topic_producer.py):

    • 连接 RabbitMQ 服务器并创建信道。

    • 使用 channel.exchange_declare(exchange='topic_logs', exchange_type='topic') 声明一个名为 topic_logs 的 Topic Exchange。exchange_type='topic' 指定交换机类型为 Topic。

    • 从命令行参数获取 routing_key (例如:kern.critical, auth.info, app.debug) 和消息内容。

    • 使用 channel.basic_publish 发布消息到 topic_logs Exchange,并将 routing_key 作为路由键。

    • 打印消息发送成功的提示信息,包含 routing_key 和消息内容。

    • 关闭连接。

  • 主题日志消费者 (topic_consumer.py):

    • 连接 RabbitMQ 服务器并创建信道。

    • 同样声明 topic_logs Exchange。

    • 创建临时队列。

    • 从命令行参数获取 binding_keys 列表 (例如:kern.*, *.critical, auth.#),这些 binding keys 定义了消费者订阅的模式。

    • 遍历 binding_keys 列表,使用 channel.queue_bind(exchange='topic_logs', queue=queue_name, routing_key=binding_key) 将临时队列绑定到 topic_logs Exchange,并为每个 binding key 指定不同的订阅模式。

      • 通配符:

        • * (星号): 匹配一个单词。例如 kern.* 可以匹配 kern.info, kern.warning, kern.critical,但不匹配 kern.critical.level

        • # (井号): 匹配零个或多个单词。例如 auth.# 可以匹配 auth.info, auth.warning, auth.info.user.login

    • 打印等待日志消息的提示信息。

    • 定义回调函数 callback(ch, method, properties, body),打印接收到的消息,包含 Routing Key 和消息内容。

    • 注册消费者并开始消费消息。


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