Search by

fiberphp / queue-redis

fiberphp

📨 FiberPHP Redis 队列驱动 —— 基于 Redis Stream,支持延迟消息、消费组、pending 认领。

Package info

gitee.com/fiberphp/queue-redis.git

Issues

pkg:composer/fiberphp/queue-redis

Statistics

Installs: 0

Dependents: 1

Suggesters: 1

dev-master 2026-09-09 05:55 UTC

This package is auto-updated.

Last update: 2026-09-09 05:55:16 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-redis
  • fiberphp/queue dev-master
  • fiberphp/redis dev-master
  • fiberphp/framework dev-master

安装

composer require fiberphp/queue-redis

安装后通过 PackageManifest 自动注册 RedisQueueProvider,无需任何独立配置文件:Provider::boot() 读取 queue.connections.redis (默认配置随 fiberphp/queue 包自动合并,装包即用),创建 RedisStreamDriver 并以 redis 名称注册到 QueueManager

配置

驱动参数位于 config/queue.phpconnections.redis 节点(queue 包已带默认值);如需调整,在应用 config/queue.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,
        ],
    ],
];
字段默认值说明
connectiondefaultRedis 连接名(对应 config/redis.php 中的键)
prefetch_count1每次拉取消息数上限(0=不限)
timer_interval0.1消费轮询间隔(秒,支持毫秒精度)
pending_timeout30000pending 超时毫秒数,超时后 xAutoClaim 重新分配
delay_scan_interval0.5延迟队列扫描间隔(秒)
fail_fasttrueboot 时 ping Redis 连接,提前暴露问题
priority_levels3优先级队列分层数(0=禁用)

使用

发布消息

// 即时消息
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