Search by

fiberphp / queue-rabbitmq

fiberphp

🐰 FiberPHP RabbitMQ 队列驱动 —— AMQP 协议实现,支持连接池、延迟消息、消息确认。

Package info

gitee.com/fiberphp/queue-rabbitmq.git

Issues

pkg:composer/fiberphp/queue-rabbitmq

Statistics

Installs: 0

Dependents: 1

Suggesters: 1

dev-master 2026-09-09 05:55 UTC

This package is auto-updated.

Last update: 2026-09-09 05:55:13 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-pcntlext-posix
  • bunny/bunny ^0.5
  • workerman/workerman ^5.0
  • fiberphp/framework dev-master
  • fiberphp/queue dev-master
  • RabbitMQ Server(延迟消息需安装 rabbitmq_delayed_message_exchange 插件)

安装

composer require fiberphp/queue-rabbitmq

安装后通过 PackageManifest 自动注册 RabbitmqQueueProvider,无需拷贝任何配置文件:驱动开关(driver / connection / prefetch_count / enable_delayed / fail_fast)的默认值随 fiberphp/queue 包合并,broker 连接池(connections.default 发布池、connections.consumer 消费长连接)随本包 config/queue/connections/rabbitmq.php 按目录约定合并,全部落在 queue.connections.rabbitmq 节点,装包即用。

配置

统一位于 config/queue.phpconnections.rabbitmq 节点;如需调整(如 broker 地址账号),在应用 config/queue.php 覆盖同名键:

// config/queue.php
return [
    'connections' => [
        'rabbitmq' => [
            // ── 驱动开关(默认值由 queue 包提供)──
            'driver'         => 'rabbitmq',
            'connection'     => 'default', // broker 连接名(对应下方 connections 键)
            'prefetch_count' => 1,         // 每次推送消息数上限
            'enable_delayed' => false,     // 延迟消息(需 broker 装 rabbitmq_delayed_message_exchange 插件)
            'fail_fast'      => true,      // boot 时初始化连接池并 ping

            // ── broker 连接池(默认值由 queue-rabbitmq 包提供,按目录约定合并)──
            'connections' => [
                'default'  => [ /* 发布连接:启用连接池 + 通道池,支撑影子模式 */ ],
                'consumer' => [ /* 消费专用连接:pool-less 长连接,订阅不归还 */ ],
            ],
        ],
    ],
];

单连接字段(connections.*.config

字段说明
host / port / vhostbroker 地址、端口、虚拟主机
username / password认证账号密码
mechanismSASL 机制(PLAIN / AMQPLAIN
timeout连接超时(秒)
restart_interval重启间隔
debug是否开启协议二进制 dump
channels_pool通道池(max_connections 限制单连接最大通道数)
client_properties客户端标识(name/version)
heartbeat_callback心跳回调 callable

连接池(connections_pool

字段说明
enabletrue=连接池;false=pool-less 专用长连接
min_connections / max_connections最小/最大连接数
idle_timeout / wait_timeout空闲超时 / 借用等待超时

使用

通过 QueueManager 投递

RabbitmqQueueProvider::boot() 会读取配置,初始化连接池并创建 RabbitmqDriver,注册到 QueueManagerrabbitmq 驱动。应用层通过统一队列 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() 会:

  1. 即时消息:声明 exchange + queue + 死信绑定(拓扑缓存,仅首次)
  2. 延迟消息:声明 x-delayed-message 交换机,附带 x-delay
  3. 通过 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\Queue\Rabbitmq\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.delayedx-delayed-message 类型

License

MIT