fiberphp / queue-redis
📨 FiberPHP Redis 队列驱动 —— 基于 Redis Stream,支持延迟消息、消费组、pending 认领。
Requires
- php: >=8.3
- fiberphp/queue: dev-master
- fiberphp/redis: dev-master
Requires (Dev)
- phpunit/phpunit: ^11.0
This package is auto-updated.
Last update: 2026-08-24 00:59:54 UTC
README
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-redisfiberphp/queuedev-masterfiberphp/redisdev-masterfiberphp/frameworkdev-master
安装
composer require fiberphp/queue-redis
安装后 PackageInstaller::discover 自动把 config/queue-redis.php 拷贝到应用 config/queue-redis.php(幂等不覆盖),并通过 PackageManifest 注册 RedisQueueProvider。Provider::boot() 读取配置,创建 RedisStreamDriver 并以 redis 名称注册到 QueueManager。
配置
config/queue-redis.php:
return [
// 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,
];
| 字段 | 默认值 | 说明 |
|---|---|---|
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 连接,提前暴露问题 |
使用
发布消息
// 即时消息
queue('redis')->push('email', json_encode(['to' => 'foo@bar', 'subject' => 'hi']));
// 延迟消息(60 秒后投递)
queue('redis')->push('reminder', json_encode(['msg' => '...']), 60);
消费消息
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:<name>} | {queue:email} |
| 消费组 | {queue:<name>}:group | {queue:email}:group |
| 延迟 ZSET | {queue:<name>}:delayed | {queue:email}:delayed |
内部机制
延迟消息
push($queue, $body, $delay) 当 $delay > 0 时,消息不直接进 Stream,而是写入延迟 ZSET(key = {queue:<name>}: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