Search by

fiberphp / queue

fiberphp

📨 FiberPHP 队列核心 —— 契约与管理器,驱动策略模式,支持生产者/消费者模型、重试/死信策略、链路追踪传播。

Package info

gitee.com/fiberphp/queue.git

Issues

pkg:composer/fiberphp/queue

Statistics

Installs: 0

Dependents: 2

Suggesters: 0

dev-master 2026-09-10 13:19 UTC

This package is auto-updated.

Last update: 2026-09-10 13:30:55 UTC


README

FiberPHP 框架的队列核心子包。提供驱动策略模式的队列管理器、通用生产者/消费者模型、JobHandler 抽象基类与自动扫描注册、驱动无关的重试/死信策略、以及 trace 上下文自动延续。

驱动实现由独立子包提供(fiberphp/queue-redisfiberphp/queue-rabbitmq),本包仅定义核心契约与管理器,应用层只需依赖本包。

特性

  • 驱动策略QueueManagerconfig('queue.default')queue.map 选择驱动并代理调用,上层只依赖 DriverInterface
  • 多驱动挂载:一个项目可同时挂 Redis / RabbitMQ 驱动,按队列维度路由
  • 通用 ProducerProducer::send / sendAsync / sendBatch / sendUsing,自动 JSON 序列化 + trace 注入
  • 通用 ConsumerConsumer::dispatch / subscribe / subscribeMany,统一包装 MessageInterface
  • JobHandler 抽象基类:用户继承后放到 app/Queue/ 目录即被 Queue 进程自动扫描注册
  • 重试/死信:消费失败后按 max_attempts × 线性退避重试,超限自动转入 {queue}:failed 死信队列,逻辑驱动无关
  • Trace 延续Tracer 在生产端写入 trace_id,消费端自动恢复,串联跨进程调用链
  • 事件分发:消费生命周期内自动派发强类型 Event(需 fiberphp/eventProducedEvent / ConsumedEvent / FailedEvent / RetriedEvent / TimeoutEvent / DeadLetteredEvent
  • 异步事件重派发Consumer 自动挂载 EventDispatchMiddleware,识别 Event::ShouldQueue 监听器推入的 {event_class, payload} 格式消息,反射重建 Event 对象后调用 Event::dispatchQueuedListeners() 只执行异步监听器

环境要求

  • PHP >= 8.3
  • fiberphp/framework dev-master
  • fiberphp/log dev-master
  • fiberphp/event dev-master
  • 至少一个队列驱动子包(fiberphp/queue-redis 和/或 fiberphp/queue-rabbitmq

安装

composer require fiberphp/queue

安装后 PackageInstaller::discover 自动把 config/queue.phpconfig/process/queue.php 拷贝到主项目(幂等不覆盖),并通过 PackageManifest 注册 Queue Worker 进程。

配置

config/queue.php

所有驱动配置集中在这一份文件(随包自动合并,无需为驱动子包维护独立配置),Provider 按 connections.{name}.driver 自动注册驱动:

return [
    // 默认驱动:'redis' | 'rabbitmq'
    'default' => env('QUEUE_CONNECTION', 'redis'),

    // 队列 → 驱动映射(可选,不配则全部走 default 驱动)
    'map' => [
        // 'order.pay'    => 'rabbitmq',
        // 'notify.sms'   => 'redis',
    ],

    // 消费侧通用策略(Job 属性可覆盖)
    'max_attempts'  => 5,    // 最大重试次数(超限转死信)
    'retry_seconds' => 5,    // 重试间隔(第 N 次 = retry_seconds × N,线性退避)
    'timeout'       => 0,    // 消费超时秒数(0=不限制)

    // 驱动连接配置(redis 驱动默认内建;rabbitmq 段随 queue-rabbitmq 包合并)
    'connections' => [
        'redis' => [
            'driver'              => 'redis',
            'connection'          => 'default', // config/redis.php 中的连接名
            'prefetch_count'      => 1,         // 每次拉取消息数上限
            'timer_interval'      => 0.1,       // 消费轮询间隔(秒)
            'pending_timeout'     => 30000,     // pending 超时毫秒数
            'delay_scan_interval' => 0.5,       // 延迟队列扫描间隔(秒)
            'fail_fast'           => true,      // boot 时 ping Redis
            'priority_levels'     => 3,         // 优先级队列分层数(0=禁用优先级)
        ],
        'rabbitmq' => [
            'driver'         => 'rabbitmq',
            'connection'     => 'default',     // broker 连接名(queue-rabbitmq 包合并 connections 池)
            'prefetch_count' => 1,
            'enable_delayed' => false,         // 是否启用 x-delayed 插件
            'fail_fast'      => true,
        ],
    ],
];

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