fiberphp / queue
📨 FiberPHP 队列核心 —— 契约与管理器,驱动策略模式,支持生产者/消费者模型、重试/死信策略、链路追踪传播。
dev-master
2026-08-23 15:43 UTC
Requires
- php: >=8.3
- fiberphp/event: dev-master
- fiberphp/framework: dev-master
- fiberphp/helper: dev-master
- fiberphp/log: dev-master
Requires (Dev)
- fiberphp/queue-rabbitmq: dev-master
- fiberphp/queue-redis: dev-master
- fiberphp/redis: dev-master
- phpunit/phpunit: ^11.0
- workerman/workerman: ^5.0
This package is auto-updated.
Last update: 2026-08-23 15:43:47 UTC
README
FiberPHP 框架的队列核心子包。提供驱动策略模式的队列管理器、通用生产者/消费者模型、JobHandler 抽象基类与自动扫描注册、驱动无关的重试/死信策略、以及 trace 上下文自动延续。
驱动实现由独立子包提供(fiberphp/queue-redis、fiberphp/queue-rabbitmq),本包仅定义核心契约与管理器,应用层只需依赖本包。
特性
- 驱动策略:
QueueManager按config('queue.default')或queue.map选择驱动并代理调用,上层只依赖DriverInterface - 多驱动挂载:一个项目可同时挂 Redis / RabbitMQ 驱动,按队列维度路由
- 通用 Producer:
Producer::send / sendAsync / sendBatch / sendUsing,自动 JSON 序列化 + trace 注入 - 通用 Consumer:
Consumer::dispatch / subscribe / subscribeMany,统一包装MessageInterface - JobHandler 抽象基类:用户继承后放到
app/Queue/目录即被Queue进程自动扫描注册 - 重试/死信:消费失败后按
max_attempts× 线性退避重试,超限自动转入{queue}:failed死信队列,逻辑驱动无关 - Trace 延续:
Tracer在生产端写入 trace_id,消费端自动恢复,串联跨进程调用链 - 事件分发:消费成功触发
queue.consumed,消费失败触发queue.failed
环境要求
- PHP >= 8.3
fiberphp/frameworkdev-masterfiberphp/logdev-masterfiberphp/eventdev-master- 至少一个队列驱动子包(
fiberphp/queue-redis和/或fiberphp/queue-rabbitmq)
安装
composer require fiberphp/queue
安装后 PackageInstaller::discover 自动把 config/queue.php 与 config/process/queue.php 拷贝到主项目(幂等不覆盖),并通过 PackageManifest 注册 Queue Worker 进程。
配置
config/queue.php
return [
// 默认驱动:'redis' | 'rabbitmq'
'default' => 'redis',
// 队列 → 驱动映射(可选,不配则全部走 default 驱动)
'map' => [
// 'order.pay' => 'rabbitmq',
// 'notify.sms' => 'redis',
],
// 消费侧通用策略
'max_attempts' => 5, // 最大重试次数(超限转死信)
'retry_seconds' => 5, // 重试间隔(第 N 次 = retry_seconds × N,线性退避)
];
config/process/queue.php
return [
'Queue:Worker' => [
'handler' => \FiberPHP\Queue\Queue::class,
'count' => 1,
],
];
使用
生产消息
use FiberPHP\Queue\Queue\Producer;
// 即时消息(走 default 驱动或 queue.map 指定的驱动)
Producer::send('order.pay', ['order_id' => 123, 'amount' => 99.9]);
// 延迟消息(60 秒后投递)
Producer::sendAsync('email.send', ['to' => 'x@y.com'], 60);
// 显式指定驱动(绕开 queue.map 默认映射)
Producer::sendUsing('rabbitmq', 'order.pay', $data);
// 批量发布
Producer::sendBatch('order.pay', [$data1, $data2, $data3]);
消费消息(继承 JobHandler)
namespace App\Queue;
use FiberPHP\Queue\Adapter\ConsumeResult;
use FiberPHP\Queue\Worker\JobHandler;
class OrderPayJob extends JobHandler
{
protected string $queue = 'order.pay'; // 也可留空,用类名 OrderPay → order_pay 推断
public function consume(array $data, array $context = []): ConsumeResult
{
$orderId = $data['order_id'];
// 业务逻辑...
return ConsumeResult::Ack;
}
// 可选:消费失败时的业务扩展点
public function failed(\Throwable $e, array $package): ?array
{
// 告警 / 补偿 / 修改 package 字段
return null;
}
}
将 OrderPayJob.php 放到主项目 app/Queue/ 目录,启动 Queue:Worker 进程后会自动扫描注册。
消费消息(闭包形式)
use FiberPHP\Queue\Queue\QueueManager;
use FiberPHP\Queue\Worker\Consumer;
use FiberPHP\Queue\Adapter\ConsumeResult;
$consumer = new Consumer();
$consumer->subscribe('order.pay', function (array $data, array $context): ConsumeResult {
// 处理消息...
return ConsumeResult::Ack;
});
多驱动路由
// 显式映射特定队列走指定驱动
QueueManager::setQueueDriver('order.pay', 'rabbitmq');
// 或通过 config/queue.php 的 map 字段批量配置
重试与死信
- 消费失败时
Consumer自动按max_attempts上限重试,退避间隔 =retry_seconds × 第 N 次 - 重试通过
Producer重发延迟消息,原消息 ack 后立即从驱动层移除 - 超出
max_attempts后转入同名{queue}:failed死信队列 - 业务可通过
JobHandler::failed()自定义补偿逻辑或修改 package 字段
Trace 上下文
Producer 在序列化时通过 Tracer::extractContext() 注入当前 trace_id,Consumer 在反序列化后通过 Tracer::restoreFromContext() 恢复,确保跨进程调用链延续。
License
MIT