RabbitMQ 实践应用 RabbitMQ 实践应用详解 引言 在现代分布式系统中,消息队列扮演着至关重要的角色。它们允许不同的服务和应用以异步、解耦的方式进行通信,从而提高系统的可靠性、可伸缩性和灵活性。RabbitMQ 作为一款开源的消息队列中间件,凭借其强大的功能、稳定的性能和易用性,成为了众多企业和开发者的首选。 1. 任务队列 (Task Queues) 1.1 应用场景 任务队列是消息队列最经典的应用场景之一。它主要用于处理耗时、异步的任务,例如: 用户注册后的邮件发送、短信通知: 用户注册成功后,发送欢迎邮件或短信通知并不需要立即完成,可以放入任务队列异步处理,避免阻塞用户注册流程。
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_mode 为 2,表示消息需要持久化,即使 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 运行示例
启动消费者: 在终端中运行 python task_consumer.py。
启动生产者并发送消息: 在另一个终端中运行 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 运行示例
启动多个订阅者: 在多个终端中分别运行 python subscriber.py,启动多个订阅者实例。
启动发布者并发送消息: 在另一个终端中运行 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: python routing_consumer.py info warning (订阅 info 和 warning 级别的日志)
终端 2: python routing_consumer.py error (订阅 error 级别的日志)
终端 3: python routing_consumer.py warning error (订阅 warning 和 error 级别的日志)
启动生产者并发送不同级别的日志消息:
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 和消息内容。
注册消费者并开始消费消息。