# queue **Repository Path**: fiberphp/queue ## Basic Information - **Project Name**: queue - **Description**: 📨 FiberPHP 队列核心抽象 —— 队列管理器、任务处理器、消息接口,不包含具体驱动实现,供驱动包继承扩展。 - **Primary Language**: PHP - **License**: MIT - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-08-22 - **Last Updated**: 2026-09-10 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # FiberPHP Queue 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,消费端自动恢复,串联跨进程调用链 - **事件分发**:消费生命周期内自动派发强类型 Event(需 `fiberphp/event`) `ProducedEvent` / `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`) ## 安装 ```bash composer require fiberphp/queue ``` 安装后 `PackageInstaller::discover` 自动把 `config/queue.php` 与 `config/process/queue.php` 拷贝到主项目(幂等不覆盖),并通过 `PackageManifest` 注册 `Queue` Worker 进程。 ## 配置 ### `config/queue.php` 所有驱动配置集中在这一份文件(随包自动合并,无需为驱动子包维护独立配置),Provider 按 `connections.{name}.driver` 自动注册驱动: ```php 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` ```php return [ 'Queue:Worker' => [ 'handler' => \FiberPHP\Queue\Queue::class, 'count' => 1, ], ]; ``` ## 使用 ### 生产消息 ```php 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) ```php 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` 进程后会自动扫描注册。 ### 消费消息(闭包形式) ```php 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; }); ``` ### 多驱动路由 ```php // 显式映射特定队列走指定驱动 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