fiberphp/queue-rabbitmq

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

Maintainers

Package info

gitee.com/FiberPHP/queue-rabbitmq

Issues

pkg:composer/fiberphp/queue-rabbitmq

Transparency log

Statistics

Installs: 0

Dependents: 1

Suggesters: 0

dev-master 2026-08-23 15:25 UTC

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-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

安装后 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 / 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\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.delayedx-delayed-message 类型

License

MIT