fiberphp / queue-rabbitmq
🐰 FiberPHP RabbitMQ 队列驱动 —— AMQP 协议实现,支持连接池、延迟消息、消息确认。
Requires
- php: >=8.3
- bunny/bunny: ^0.5
- fiberphp/framework: dev-master
- fiberphp/queue: dev-master
- workerman/workerman: ^5.0
Requires (Dev)
- phpunit/phpunit: ^11.0
This package is auto-updated.
Last update: 2026-08-23 15:33:03 UTC
README
FiberPHP 框架的 RabbitMQ (AMQP 0-9-1) 队列驱动子包。自实现 AMQP 协议编解码(基于 Workerman AsyncTcpConnection,不依赖阻塞型 ext-amqp),提供连接池 + 通道池管理、协程挂起/恢复、延迟消息(DLX + TTL / delayed-message-exchange 插件)与死信队列,通过 config/queue-rabbitmq.php 声明式管理多连接。
特性
- 自实现 AMQP 协议:
RabbitmqProtocol实现 Workerman 协议接口(input/decode/encode),内部委托bunny/bunny的 ProtocolReader/ProtocolWriter 完成 Frame 与字节互转,非阻塞、全协程化 - 连接池 + 通道池:
ConnectionsManagement统一管理多连接的连接池与通道池,支持 pool / pool-less 双模式;channel 池耗尽时影子模式换连接重试 - 协程挂起:
Client::await()注册等待者并挂起当前协程,帧到达后wakeup()恢复,使异步 AMQP 调用像同步代码 - 消费隔离:
default(发布)与consumer(消费)使用独立连接,destroy 互不影响,无借还竞态 - 延迟消息:
x-delayed-message交换机(需 broker 插件)或 DLX + TTL 两种路径 - 死信队列:每个队列自动声明 DLX +
.failed死信队列,nack(requeue=false) 转死信 - 拓扑缓存:
RabbitmqDriver缓存已声明的 exchange/queue/binding,避免重复 declare - 心跳保活:跟随 broker 协商的心跳间隔定时发送 Heartbeat 帧,发送失败标记连接死亡并唤醒所有等待协程
- 事务 / 发布确认:Channel 支持
select()/commit()/rollback()事务模式与confirm()发布确认模式
环境要求
- PHP >= 8.3
ext-pcntl、ext-posixbunny/bunny^0.5workerman/workerman^5.0fiberphp/frameworkdev-masterfiberphp/queuedev-master- RabbitMQ Server(延迟消息需安装
rabbitmq_delayed_message_exchange插件)
安装
composer require fiberphp/queue-rabbitmq
安装后 PackageInstaller::discover 自动把 config/queue-rabbitmq.php 拷贝到应用 config/queue-rabbitmq.php(幂等不覆盖),并通过 PackageManifest 注册 RabbitmqQueueProvider。
配置
config/queue-rabbitmq.php:
return [
// AMQP 连接名(对应 connections 中的键)
'connection' => 'default',
// 每次推送消息数上限
'prefetch_count' => 1,
// 是否启用延迟消息(需 broker 安装 rabbitmq_delayed_message_exchange 插件)
'enable_delayed' => true,
// 延迟交换机类型(x-delayed-type 参数)
'delayed_plugin_args' => ['x-delayed-type' => 'direct'],
// fail_fast:boot 时初始化连接池并 ping
'fail_fast' => true,
// 日志 LoggerInterface | LoggerInterface::class
'logger' => null,
// AMQP 连接配置
'connections' => [
'default' => [ /* 发布连接 + 连接池 + 通道池 */ ],
'consumer' => [ /* 消费专用连接(pool-less 长连接) */ ],
],
];
单连接字段(connections.*.config)
| 字段 | 说明 |
|---|---|
host / port / vhost | broker 地址、端口、虚拟主机 |
username / password | 认证账号密码 |
mechanism | SASL 机制(PLAIN / AMQPLAIN) |
timeout | 连接超时(秒) |
restart_interval | 重启间隔 |
debug | 是否开启协议二进制 dump |
channels_pool | 通道池(max_connections 限制单连接最大通道数) |
client_properties | 客户端标识(name/version) |
heartbeat_callback | 心跳回调 callable |
连接池(connections_pool)
| 字段 | 说明 |
|---|---|
enable | true=连接池;false=pool-less 专用长连接 |
min_connections / max_connections | 最小/最大连接数 |
idle_timeout / wait_timeout | 空闲超时 / 借用等待超时 |
使用
通过 QueueManager 投递
RabbitmqQueueProvider::boot() 会读取配置,初始化连接池并创建 RabbitmqDriver,注册到 QueueManager 的 rabbitmq 驱动。应用层通过统一队列 API 投递:
// 投递即时消息
queue('rabbitmq')->push('orders', json_encode(['id' => 1, 'action' => 'created']));
// 投递延迟消息(需 enable_delayed=true + broker 插件)
queue('rabbitmq')->later(60, 'orders', json_encode(['id' => 2, 'action' => 'delayed']));
RabbitmqDriver::push() 会:
- 即时消息:声明 exchange + queue + 死信绑定(拓扑缓存,仅首次)
- 延迟消息:声明
x-delayed-message交换机,附带x-delay头 - 通过
ConnectionsManagement::publish()借连接 → 借通道 → 发布 → 归还
消费
queue('rabbitmq')->consume('orders', function (MessageInterface $message): ConsumeResult {
$data = json_decode($message->getBody(), true);
// 处理消息...
return ConsumeResult::Ack; // Ack | Nack | Reject
});
RabbitmqDriver::consume() 在 consumer 专用连接上借通道,注册 basic.consume 回调;handler 抛异常自动 nack 重新入队。nack(requeue=false) 或 reject 会将消息转死信队列 orders.failed。
直接使用连接管理
高频发布场景可绕过 Builder,直接借连接复用 channel:
use FiberPHP\QueueRabbitmq\ConnectionsManagement;
// 便捷发布(单次消息)
ConnectionsManagement::publish(
body: '{"id":1}',
exchange: 'queue.orders',
routingKey: 'orders',
headers: ['delivery-mode' => 2],
connection: 'default',
);
// 借→用→还模式(闭包内复用 channel)
ConnectionsManagement::connection(function (ClientInterface $client) {
$channel = $client->channel(false);
$channel->publish(/* ... */);
}, 'default');
Context 复用模式
同一协程内多次操作可让 channel 自动复用、defer 归还(适合消费者):
$channel = $client->channel(); // 不传 closure → Context 复用
$channel->publish(/* ... */);
$channel->ack($message);
// 协程结束时自动归还 channel
拓扑约定
| 资源 | 命名 | 说明 |
|---|---|---|
| 主交换机 | queue.{queue} | direct 类型,持久化 |
| 主队列 | {queue} | 绑定死信交换机,持久化 |
| 死信交换机 | queue.dlx.{queue} | direct 类型 |
| 死信队列 | {queue}.failed | 持久化 |
| 延迟交换机 | queue.delayed | x-delayed-message 类型 |
License
MIT