一、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)
发表评论