# queue-redis **Repository Path**: fiberphp/queue-redis ## Basic Information - **Project Name**: queue-redis - **Description**: 🔴 FiberPHP Redis 队列驱动 —— 基于 `fiberphp/queue` 核心,提供 Redis 队列实现,支持延迟队列、优先级、协程非阻塞 - **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-09 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # FiberPHP Queue Redis FiberPHP 框架的 Redis Stream 队列驱动子包。基于 Redis Stream 实现消息队列,支持延迟消息(ZSET + Timer 扫描)、消费组(Consumer Group)、pending 超时自动 claim 重新分配,通过 `config/queue-redis.php` 声明式配置,由 `RedisQueueProvider` 注册到 `QueueManager`。 ## 特性 - **Redis Stream**:基于 `xAdd / xReadGroup / xAck` 实现可靠消息队列 - **消费组**:多 consumer 共享同一 group,消息均衡分配;`xGroup CREATE` 幂等(忽略 BUSYGROUP) - **延迟消息**:写入 ZSET(score=到期时间戳),Timer 按 `delay_scan_interval` 扫描到期后转移到 Stream - **pending claim**:消费超时(`pending_timeout`)的消息由 `xAutoClaim` 重新分配给其他 consumer - **fail_fast 探活**:boot 时 ping Redis 连接,提前暴露连接问题 - **Timer 轮询**:基于 Workerman Timer,主消费 / 延迟扫描 / pending claim 三组独立 Timer ## 环境要求 - PHP >= 8.3 - `ext-redis` - `fiberphp/queue` dev-master - `fiberphp/redis` dev-master - `fiberphp/framework` dev-master ## 安装 ```bash composer require fiberphp/queue-redis ``` 安装后通过 `PackageManifest` 自动注册 `RedisQueueProvider`,无需任何独立配置文件:`Provider::boot()` 读取 `queue.connections.redis` (默认配置随 `fiberphp/queue` 包自动合并,装包即用),创建 `RedisStreamDriver` 并以 `redis` 名称注册到 `QueueManager`。 ## 配置 驱动参数位于 `config/queue.php` 的 `connections.redis` 节点(queue 包已带默认值);如需调整,在应用 `config/queue.php` 覆盖同名键: ```php // config/queue.php return [ 'connections' => [ 'redis' => [ 'driver' => 'redis', // Redis 连接名(对应 config/redis.php 中的键) 'connection' => 'default', // 每次拉取消息数上限(0=不限) 'prefetch_count' => 1, // 消费轮询间隔(秒,支持毫秒精度如 0.1) 'timer_interval' => 0.1, // pending 超时毫秒数,超时后 xAutoClaim 重新分配给其他 consumer 'pending_timeout' => 30000, // 延迟队列扫描间隔(秒) 'delay_scan_interval' => 0.5, // fail_fast:boot 时 ping Redis 连接,提前暴露问题 'fail_fast' => true, // 优先级队列分层数(0=禁用优先级) 'priority_levels' => 3, ], ], ]; ``` | 字段 | 默认值 | 说明 | |-----------------------|-----------|--------------------------------------------------| | `connection` | `default` | Redis 连接名(对应 `config/redis.php` 中的键) | | `prefetch_count` | `1` | 每次拉取消息数上限(0=不限) | | `timer_interval` | `0.1` | 消费轮询间隔(秒,支持毫秒精度) | | `pending_timeout` | `30000` | pending 超时毫秒数,超时后 `xAutoClaim` 重新分配 | | `delay_scan_interval` | `0.5` | 延迟队列扫描间隔(秒) | | `fail_fast` | `true` | boot 时 ping Redis 连接,提前暴露问题 | | `priority_levels` | `3` | 优先级队列分层数(0=禁用) | ## 使用 ### 发布消息 ```php // 即时消息 queue('redis')->push('email', json_encode(['to' => 'foo@bar', 'subject' => 'hi'])); // 延迟消息(60 秒后投递) queue('redis')->push('reminder', json_encode(['msg' => '...']), 60); ``` ### 消费消息 ```php use FiberPHP\Queue\Adapter\ConsumeResult; use FiberPHP\Queue\Adapter\MessageInterface; queue('redis')->consume('email', function (MessageInterface $message): ConsumeResult { $body = json_decode($message->getBody(), true); // 处理逻辑... return ConsumeResult::Ack; // 成功,确认消费 // return ConsumeResult::Nack; // 失败,留 pending 等待 claim 重新分配 // return ConsumeResult::Reject; // 拒绝,直接 ack 丢弃(死信由上层处理) }); ``` ### Key 命名约定 | 类型 | 规则 | 示例 | |-----------|--------------------------|-------------------------| | Stream | `{queue:}` | `{queue:email}` | | 消费组 | `{queue:}:group` | `{queue:email}:group` | | 延迟 ZSET | `{queue:}:delayed` | `{queue:email}:delayed` | ## 内部机制 ### 延迟消息 `push($queue, $body, $delay)` 当 `$delay > 0` 时,消息不直接进 Stream,而是写入延迟 ZSET(key = `{queue:}:delayed` ,member = 序列化的 entry,score = 到期时间戳)。Timer 按 `delay_scan_interval` 间隔扫描,将 `score <= now` 的成员 `xAdd` 到 Stream 后从 ZSET 删除。Entry 内的 `id` 重置为 `*` 让 Stream 自动生成。 ### 消费组 首次 `consume()` 时通过 `xGroup CREATE ... MKSTREAM` 创建 Stream 与消费组(幂等,已存在时忽略 BUSYGROUP)。主 Timer 按 `timer_interval` 调用 `xReadGroup` 拉取新消息,每条消息回调 handler,根据返回的 `ConsumeResult` 走 ack / nack / reject 分支。 ### pending claim 消费组模式下,消息被 consumer 拉取后进入 pending list 直到 `xAck`。若 consumer 崩溃导致消息长时间未 ack,Timer 按 `pending_timeout` 间隔调用 `xAutoClaim`,将空闲时间超过 `pending_timeout` 毫秒的 pending 消息转移到当前 consumer 重新消费。 ### nack 策略 Redis Stream 无原生 nack,驱动实现如下: - `nack(requeue=true)`:不调用 `xAck`,消息留 pending,由 `pending_timeout` 后的 claim 机制重新分配 - `nack(requeue=false)`:直接 `xAck` 丢弃,死信处理由上层负责 ## License MIT