Redis Streams 详解

一、XADD — 添加消息

向 Stream 添加一条新消息,自动生成唯一 ID。

bash
XADD mystream * field1 "value1" field2 "value2"
# 返回生成的消息 ID: "1732867445000-0"

参数说明:

  • mystream — Stream 键名,不存在自动创建
  • * — 自动生成 ID,格式为 timestamp-sequenceNumber
  • 也可手动指定 ID:XADD mystream "1732867445000-0" field "value"

ID 必须严格递增,比已有最大 ID 更大。

常用选项:

bash
# MAXLEN ~ 1000:近似裁剪,保留约 1000 条,高效
XADD mystream MAXLEN ~ 1000 * field "value"

# MINID:删除早于指定 ID 的消息
XADD mystream MINID "1732867445000-0" * field "value"

二、XREAD — 读取新消息(阻塞/非阻塞)

从 Stream 读取消息,支持阻塞等待。

基本语法

bash
# 非阻塞读
XREAD STREAMS mystream "0"        # 从头开始读
XREAD STREAMS mystream "$"        # 只读新消息(写完为止)

# 阻塞读,最长等待 3 秒
XREAD BLOCK 3000 STREAMS mystream "$"

# 阻塞读,永久等待直到有新消息
XREAD BLOCK 0 STREAMS mystream "$"

多 Stream 同时读

bash
XREAD BLOCK 5000 STREAMS s1 s2 s3 "$" "$" "$"

关键特性

  • XREAD 没有消费者组概念,所有客户端看到相同的消息
  • 适合一对多广播场景(类似 Redis Pub/Sub,但消息可持久化)
  • 0 表示从头读,$ 表示只读新消息

三、XREADGROUP — 消费组模式读取(核心)

前提:必须先创建消费者组。 消息只分发给组内一个消费者

创建消费者组

bash
# "0" 表示从头消费,"$" 表示只消费新消息
XGROUP CREATE mystream mygroup "0"
XGROUP CREATE mystream mygroup "$" MKSTREAM    # 如果 Stream 不存在,自动创建

消费者读取消息

bash
XREADGROUP GROUP mygroup consumer1 STREAMS mystream ">"

参数解释:

参数 含义
mygroup 消费者组名
consumer1 当前消费者实例名(用于追踪归属)
mystream Stream 键名
> 只读从未发送给任何消费者的新消息

返回格式示例:

1) 1) "mystream"
   2) 1) 1) "1732867445000-0"
         2) 1) "field"
            2) "value"

读取 Pending 消息(故障恢复)

bash
# "0" 读取 pending 队列中的消息(已分发给消费者但未 ACK)
XREADGROUP GROUP mygroup consumer1 STREAMS mystream "0"

> 读新消息,用 0 读 pending 消息 — 两者结合实现完整的消息追踪。


四、XACK — 确认消息已处理

消费者处理完成后发送 ACK,从 pending 列表中移除。

bash
# 单条确认
XACK mystream mygroup "1732867445000-0"

# 批量确认
XACK mystream mygroup "1732867445000-0" "1732867445000-1" "1732867445000-2"

返回值为确认成功的消息数量。


五、工作流程图

Producer                        Consumer Group
   │                                │
   ├── XADD (写消息)                  │
   │    │                            │
   │    └───► Stream ────────────────►│
   │         pending list             │
   │              │                  │
   │              │  XREADGROUP       │
   │              ▼                  │
   │         分发给一个 consumer       │
   │              │                  │
   │              │  处理中...         │
   │              │                  │
   │              │  XACK ───────────►│ 确认完成
   │              ▼                  │
   │         从 pending 删除           │

六、场景对比

场景 推荐命令 原因
异步任务队列(多 worker 竞争) XADD + XREADGROUP + XACK 消息只被一个 worker 抢到处理
聊天消息广播(多人独立订阅) XADD + XREAD 每个客户端独立读取,无需 ACK
日志收集上报 XADD + XREADGROUP 确认处理防丢失
订单超时检测 XADD + XREADGROUP + XACK pending 中监控超时,ACK 后脱管

七、辅助命令

bash
# 查看消费者组信息
XINFO GROUPS mystream

# 查看 pending 消息(谁在处理、处理多久)
XPENDING mystream mygroup

# 查看 Stream 信息(长度、首个/最后消息 ID)
XINFO STREAM mystream

# 手动删除消息
XDEL mystream "message-id"

# 裁剪 Stream 长度
XTRIM mystream MAXLEN 1000

# 列出消费者
XINFO CONSUMERS mystream mygroup

八、PHP (Webman) 队列示例

生产者(入队)

php
use support\Redis;

// 入队
Redis::xadd('task_queue', '*', ['task' => json_encode($data)]);

// 入队并设置上限(近似裁剪 10000 条)
Redis::xadd('task_queue', 'MAXLEN', '~', 10000, '*', ['task' => json_encode($data)]);

消费者(出队)

php
Redis::xReadGroup(
    string $group,         // 消费者组名
    string $consumer,      // 当前消费者名称
    array  $streams,       // Stream 名 => ID 的关联数组
    ?int   $count  = null, // 每次最多取多少条(可选)
    ?int   $block  = null  // 阻塞等待毫秒数(可选)
): array|Redis|false
// 读新消息("$" 不适用,XREADGROUP 用 ">" 代表未分配的新消息)
['stream_key' => '>']

// 读 pending(待确认的旧消息,用于故障恢复)
['stream_key' => '0']

// 从指定 ID 开始读
['stream_key' => '1732867445000-0']
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);

// 读新消息(阻塞 2 秒,最多取 5 条)
$result = $redis->xReadGroup(
    'mygroup',     // $group
    'worker_1',    // $consumer
    ['mystream' => '>'],  // $streams
    5,             // $count,最多取 5 条
    2000           // $block,2 秒
);

// 读 pending 消息(非阻塞)
$pending = $redis->xReadGroup(
    'mygroup',
    'worker_1',
    ['mystream' => '0'],  // "0" 读 pending
    10,             // 最多 10 条
    null            // 非阻塞
);

// 多 Stream 同时读
$result = $redis->xReadGroup(
    'mygroup',
    'worker_1',
    [
        'stream_a' => '>',
        'stream_b' => '>',
    ],
    3,     // 每个 Stream 最多 3 条
    5000   // 5 秒阻塞
);


//简单例子
use support\Redis;

$group = 'task_queue';
$consumer = 'worker_' . gethostname();
$count = 10;   // 每次最多取 10 条
$block = 3000; // 阻塞 3 秒

while (true) {
    $result = Redis::xreadgroup($group, $consumer, [$group, '>'], $count, $block);

    if (empty($result)) {
        continue;
    }

    foreach ($result as $stream => $messages) {
        foreach ($messages as $id => $fields) {
            $task = json_decode($fields['task'] ?? '{}', true);

            try {
                // 处理任务
                $this->process($task);

                // 确认完成
                Redis::xack($group, $consumer, $id);
            } catch (\Throwable $e) {
                // 记录日志,消息仍在 pending 中,后续可重新处理
                log_error('Task failed: ' . $e->getMessage());
            }
        }
    }
}

Pending 消息兜底(故障恢复)

php
// 定期检查 pending 消息,防止消息卡死
$result = Redis::xreadgroup($group, $consumer, [$group, '0'], 100, 0);

foreach ($result as $stream => $messages) {
    foreach ($messages as $id => $fields) {
        // 检查处理时长,超过阈值则重试或转移
        $info = Redis::xpending($group, $consumer, '-', '+', 1, $id);
        // $info: [messageId, consumer, idleTime(ms), deliveryCount]
        if ($info[2] > 60000) { // 超过 60 秒未确认
            // 重新处理或记录告警
        }
    }
}

// 读新消息(阻塞 2 秒,最多取 5 条)
$result = $redis->xReadGroup(
    'mygroup',     // $group
    'worker_1',    // $consumer
    ['mystream' => '>'],  // $streams
    5,             // $count,最多取 5 条
    2000           // $block,2 秒
);

// 读 pending 消息(非阻塞)
$pending = $redis->xReadGroup(
    'mygroup',
    'worker_1',
    ['mystream' => '0'],  // "0" 读 pending
    10,             // 最多 10 条
    null            // 非阻塞
);

// 多 Stream 同时读
$result = $redis->xReadGroup(
    'mygroup',
    'worker_1',
    [
        'stream_a' => '>',
        'stream_b' => '>',
    ],
    3,     // 每个 Stream 最多 3 条
    5000   // 5 秒阻塞
);

九、与 Redis List 的对比

特性 Redis List (LPUSH/BRPOP) Redis Streams
消息持久化 需自行处理 自动支持
多消费者 不支持 XREADGROUP 原生支持
消息追踪 pending 队列完整追踪
消息确认 XACK 确认机制
消息广播 无(轮询分发) XREAD 多客户端独立消费
复杂度 简单 稍高,但功能完整

需要可靠队列、多消费者竞争、消息确认,选 Streams;只是简单的先进先出,选 List 更轻量。

评论 (0)

发表评论