RabbitMQ延迟插件的安装和使用
背景介绍
在消息队列应用中,常需要实现"延时任务"场景(如订单30分钟未支付自动取消)。RabbitMQ本身可通过死信队列(DLX) 结合消息TTL实现延迟,但配置复杂且不够灵活。官方延迟插件 rabbitmq_delayed_message_exchange 提供了一种更优雅的方案,它通过一种特殊的交换机类型(x-delayed-message)直接支持延迟消息路由,无需依赖死信机制,简化了开发流程。
完整操作步骤与代码
以下基于原文(RabbitMQ 3.8版本 + PHP php-amqplib扩展)记录完整过程。
环境准备
默认根目录已安装 composer 并下载安装 rabbitmq(php-amqplib/php-amqplib)扩展
RabbitMQ 版本:3.8
1. 插件下载与安装
获取插件:
RabbitMQ 插件地址:https://www.rabbitmq.com/community-plugins
找到延迟插件
rabbitmq_delayed_message_exchange,去 GitHub 下载原文对应版本 3.8 下载地址:
https://github.com/rabbitmq/rabbitmq-rtopic-exchange/releases/download/v3.8.0/rabbitmq_rtopic_exchange-3.8.0.ez
上传并启用插件:
# 切换到插件目录(根据实际安装路径调整) cd /usr/lib/rabbitmq/lib/rabbitmq_server-3.8.19/plugins # 启动插件 rabbitmq-plugins enable rabbitmq_delayed_message_exchange # 查看插件是否启动成功(对应的插件有 E*,说明启动成功) rabbitmq-plugins list # 重启 RabbitMQ 服务 systemctl restart rabbitmq-server
2. 生产者代码(delay_pub.php)
完整代码如下:
<?php require_once __DIR__ . '/vendor/autoload.php'; use PhpAmqpLib\Connection\AMQPStreamConnection; use PhpAmqpLib\Message\AMQPMessage; use PhpAmqpLib\Wire\AMQPTable; $dbName = 'sanqing'; $dbPwd = '111111'; $tableName = 'order'; $connection = new AMQPStreamConnection('localhost', 5672, $dbName, $dbPwd, $tableName); $channel = $connection->channel(); $exc_name = 'delay_exc_pay'; $routing_key = 'delay_route_pay'; $queue_name = 'delay_queue_pay'; $ttl = 20000; // 延迟时间,单位毫秒(20秒) $channel->exchange_declare($exc_name, 'x-delayed-message', false, true, false); $args = new AMQPTable(['x-delayed-type' => 'direct']); $channel->queue_declare($queue_name, false, true, false, false, false, $args); $channel->queue_bind($queue_name, $exc_name, $routing_key); $data = 'this is ' . $routing_key . ' delay message'; $arr = [ 'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT, 'application_headers' => new AMQPTable(['x-delay' => $ttl]) ]; $msg = new AMQPMessage($data, $arr); $channel->basic_publish($msg, $exc_name, $routing_key); $channel->close(); $connection->close();
3. 消费者代码(delay_work.php)
完整代码如下:
<?php require_once __DIR__ . '/vendor/autoload.php'; use PhpAmqpLib\Connection\AMQPStreamConnection; $dbName = 'sanqing'; $dbPwd = '111111'; $tableName = 'order'; $connection = new AMQPStreamConnection('localhost', 5672, $dbName, $dbPwd, $tableName); $channel = $connection->channel(); $exc_name = 'delay_exc_pay'; $routing_key = 'delay_route_pay'; $queue_name = 'delay_queue_pay'; $channel->exchange_declare($exc_name, 'x-delayed-message', false, true, false); $channel->queue_bind($queue_name, $exc_name, $routing_key); $callback = function($msg) { echo 'received ' . $msg->body . "\n"; $msg->ack(); }; $channel->basic_qos(null, 1, null); $channel->basic_consume($queue_name, '', false, false, false, false, $callback); while ($channel->is_open()) { $channel->wait(); } $channel->close(); $connection->close();
4. 运行验证
按顺序执行以下命令:
# 终端1:启动消费者(保持运行) php delay_work.php # 终端2:运行生产者(发送延迟消息) php delay_pub.php
因为消息 TTL 过期时间设置的是 20 秒,所以消费者在生产者启动 20 秒后才能收到数据。
5. RabbitMQ 管理界面配置(重要)
在 RabbitMQ 管理界面新建交换器时,需要正确配置参数,否则启动生产者 delay_pub.php 可能会报错:x-delayed-type must be an existing exchange type。
配置步骤:
登录管理界面(默认 http://localhost:15672)
进入 Exchanges 选项卡,点击 "Add a new exchange"
填写以下信息:
| 配置项 | 值 |
|---|---|
| Virtual host | order(或你的实际 vhost) |
| Name | delay_exc_pay |
| Type | x-delayed-message |
| Durability | Durable |
| Auto delete | No |
| Internal | No |
| Arguments | x-delayed-type = direct(String 类型) |
6. 验证插件生效
在 RabbitMQ 管理界面的 Exchanges 列表中,可以看到:
类型为
x-delayed-message的交换器带有 "DM" 标记(表示 Delayed Message 插件已生效)
示例:
delay_exc_pay | x-delayed-message | D DM
关键注意事项总结
| 注意事项 | 说明 |
|---|---|
| 版本匹配 | 插件版本需与 RabbitMQ 服务版本严格对应 |
| 交换机参数 | 必须设置 x-delayed-type 参数指定内部路由类型(如 direct) |
| 延迟单位 | x-delay 参数单位为毫秒,如 20000 = 20秒 |
| 管理界面配置 | 若通过代码声明失败,可先在管理界面手动创建交换器 |
| 消费者持久运行 | 消费者需保持运行状态,才能接收延迟后的消息 |
| 消息持久化 | 示例中设置了 DELIVERY_MODE_PERSISTENT,保证消息可靠性 |
请先 登录后发表评论 ~