fiberphp/queue-redis

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

Maintainers

Package info

gitee.com/FiberPHP/queue-redis

Issues

pkg:composer/fiberphp/queue-redis

Transparency log

Statistics

Installs: 0

Dependents: 1

Suggesters: 0

dev-master 2026-08-24 00:59 UTC

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-redis
  • fiberphp/queue dev-master
  • fiberphp/redis dev-master
  • fiberphp/framework dev-master

安装

composer require fiberphp/queue-redis

安装后 PackageInstaller::discover 自动把 config/queue-redis.php 拷贝到应用 config/queue-redis.php(幂等不覆盖),并通过 PackageManifest 注册 RedisQueueProviderProvider::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,
];
字段默认值说明
connectiondefaultRedis 连接名(对应 config/redis.php 中的键)
prefetch_count1每次拉取消息数上限(0=不限)
timer_interval0.1消费轮询间隔(秒,支持毫秒精度)
pending_timeout30000pending 超时毫秒数,超时后 xAutoClaim 重新分配
delay_scan_interval0.5延迟队列扫描间隔(秒)
fail_fasttrueboot 时 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