本指南用 node-redis 构建以 Redis 为后端的 Node.js 事件流管线。示例包含一个基于 Node.js 标准 http 模块的小型本地服务器:向单个 Redis Stream 生产事件,观察两个独立消费者组按各自速度读取,并在模拟消费者崩溃后用 XAUTOCLAIM 恢复卡住的投递。
概览
Redis Stream 是只追加的字段/值日志,条目拥有自动生成、按时间排序的 ID。生产者用 XADD 追加;消费者属于消费者组,用 XREADGROUP 读取。每组整体维护一个 last-delivered-id 游标,每个消费者也有已收到但尚未确认的待处理条目列表(PEL)。
处理完成后调用 XACK 清除 PEL 条目;未确认时间超过空闲阈值的条目可通过 XAUTOCLAIM 转交健康消费者。这提供:
- 有序、持久的历史,多个独立组按各自速度读取。
- 至少一次投递、各消费者待处理列表及崩溃恢复。
- 组内水平扩展:增加消费者,Redis 自动分配工作。
- 通过 XRANGE 重放任意范围,不依赖消费者组状态。
- 通过
XADD MAXLEN ~或 XTRIM MINID ~ 限定保留范围,无需单独清理任务。
示例把 order.placed、order.paid、order.shipped、order.cancelled 写入 demo:events:orders。两个组读取相同流:notifications 中 worker-a 与 worker-b 分担工作,模拟工作池;analytics 中 worker-c 独立处理完整事件流。
工作流程
- 应用调用
stream.produce(eventType, payload),执行带近似MAXLEN ~上限的 XADD,Redis 分配按时间排序的 ID。 - 每个消费者的异步循环用特殊 ID
>调用 XREADGROUP,表示读取该组尚未投递给任何人的条目,并设置较短阻塞超时。 - 处理后调用 XACK,让 Redis 将条目移出组的待处理列表。
- 消费者在确认前被终止或崩溃,条目留在 PEL;定期 XAUTOCLAIM 扫描将闲置条目转给健康消费者。
- 任何代码,包括组外代码,都可用 XRANGE 读取历史,不影响组游标。
每组都有自己的游标和待处理列表,因此两个组处理相同事件时不需要彼此协调。
事件流辅助类
EventStream 封装流操作,见完整源码:
const { createClient } = require("redis");
const { EventStream } = require("./eventStream");
const client = createClient({ socket: { host: "localhost", port: 6379 } });
await client.connect();
const stream = new EventStream({
redisClient: client,
streamKey: "demo:events:orders",
maxlenApprox: 2000, // retention guardrail
claimMinIdleMs: 5000, // XAUTOCLAIM threshold
});
// Producer
const streamId = await stream.produce("order.placed", {
order_id: "o-1234",
customer: "alice",
amount: "49.50",
});
// Consumer group + one consumer
await stream.ensureGroup("notifications", "0-0");
const entries = await stream.consume("notifications", "worker-a", 10, 500);
for (const [entryId, fields] of entries) {
handle(fields); // your processing
await stream.ack("notifications", [entryId]); // XACK
}
// Recover stuck PEL entries by reaping them into a healthy consumer.
// The textbook pattern: each consumer periodically calls XAUTOCLAIM
// with itself as the target and processes whatever it claimed.
// `ConsumerWorker.reapIdlePel` wraps that flow; the low-level helper
// `stream.autoclaim(group, targetName)` is also available if you
// want to drive XAUTOCLAIM directly.
const result = await workerB.reapIdlePel();
// result === { claimed: N, processed: M, deletedIds: [...] }
// deletedIds are PEL entries whose payload was already trimmed.
// Redis 7+ has already removed those slots from the PEL, so no XACK
// is needed — log them and route to a dead-letter store for audit.
// Replay history (independent of any group's cursor)
for (const [entryId, fields] of await stream.replay("-", "+", 50)) {
console.log(entryId, fields);
}
数据模型
每个事件是一个流条目,内容为扁平字符串字段/值对象,带自动生成的时间有序 ID:
demo:events:orders
1716998413541-0 type=order.placed order_id=o-1234 customer=alice amount=49.50 ts_ms=...
1716998413542-0 type=order.paid order_id=o-1234 customer=alice amount=49.50 ts_ms=...
1716998413542-1 type=order.shipped order_id=o-1235 customer=bob amount=12.00 ts_ms=...
...
ID 为 {milliseconds}-{sequence},在同一条流内单调递增,因此无需额外索引即可按近似墙上时钟时间查询。排序仅在单流内成立;不同流同一毫秒追加的事件可能得到相同 ID。
实现使用:
XADD ... MAXLEN ~ n,通过管线批量生产并限制保留长度。XREADGROUP与>,投递新条目。XACK,确认每条已处理事件。XAUTOCLAIM,将闲置待处理条目交给健康消费者。XRANGE,重放与审计。- XPENDING,检查组内 PEL。
- XINFO STREAM、XINFO GROUPS 和 XINFO CONSUMERS,提供基本可观测信息。
XTRIM,显式执行保留策略。
生产事件
produceBatch 通过管线在一次往返中发送多次 XADD。每次都带近似 MAXLEN ~ 上限:
async produceBatch(events) {
const list = Array.from(events);
if (list.length === 0) return [];
const pipe = this.redis.multi();
for (const [eventType, payload] of list) {
const fields = EventStream._encodeFields(eventType, payload);
pipe.xAdd(this.streamKey, "*", fields, {
TRIM: {
strategy: "MAXLEN",
strategyModifier: "~",
threshold: this.maxlenApprox,
},
});
}
// execAsPipeline sends the commands in one round trip without
// wrapping them in MULTI/EXEC.
const ids = await pipe.execAsPipeline();
this._producedTotal += ids.length;
return ids.map((id) => String(id));
}
~ 允许 Redis 在宏节点边界裁剪,比精确裁剪便宜得多,适合保留长度的保护阈值而非严格大小约束。生产300条事件、设置 MAXLEN ~ 50 时,可能剩100条:Redis 释放了最旧的整个宏节点后停止,后续 XADD 将继续保持长度稳定。
确实需要精确上限时,把 strategyModifier: "~" 改为 "=";繁忙流中的性能差别可能明显。
通过消费者组读取
组内每个消费者执行相同循环,> 表示该组尚未向任何消费者投递的条目:
async consume(group, consumer, count = 10, blockMs = 500) {
const raw = await this.redis.xReadGroup(
group,
consumer,
[{ key: this.streamKey, id: ">" }],
{ COUNT: count, BLOCK: blockMs },
);
return flattenEntries(raw);
}
blockMs 让空闲流的读取也高效:连接在服务端等待新消息或超时,避免忙循环。XREADGROUP BLOCK 会占用底层 socket 直到返回,因此示例每个消费者都有复制出的独立 Redis 客户端,阻塞读取不会排在其他工作者读取或 HTTP 处理命令之后。
改用 0-0 等显式 ID 时,含义不同:重放已投递给该消费者名称的私有 PEL 条目。同名消费者重启时,通常先补齐自己的待处理条目,再读取新消息。
确认条目
处理完成后,XACK 告诉 Redis 可以从组的 PEL 移除条目:
async ack(group, ids) {
const idList = Array.from(ids);
if (idList.length === 0) return 0;
const n = Number(await this.redis.xAck(this.streamKey, group, idList));
this._ackedTotal += n;
return n;
}
这是至少一次投递的关键:未确认条目一直留在 PEL,直到被认领。消费者在处理与确认之间崩溃,下次扫描会重新取得条目。
但保留策略是例外:即使 ID 仍在 PEL,XADD MAXLEN ~ 与 XTRIM 也可能释放其载荷。下次 XAUTOCLAIM 会在 deletedMessages 中返回这些 ID,并在同一命令内清除 PEL 槽位;内容已不能重试,应记录并送到死信存储以便审计。后文 autoclaim 明确处理这一情况。
这与 pub/sub 的取舍相反:慢或崩溃消费者可以恢复未确认消息,但下游必须幂等。第一次处理产生副作用后、确认前崩溃,第二次再处理相同订单也必须安全。
多个消费者组共享一条流
Redis Streams 与任务队列的重要区别,是任意数量独立组都能读取同一条流。本例创建两组:
await stream.ensureGroup("notifications", "0-0");
await stream.ensureGroup("analytics", "0-0");
每组独立游标。生产五条事件,两组各收到全部五条,无需协调;notifications 组内由两名工作者分担。Redis 为每次 XREADGROUP 分配组内尚未投递的工作,无需应用实现再均衡。原示例描述增加第二名工作者可使吞吐翻倍,实际仍取决于处理能力与瓶颈。
第二参数 "0-0" 表示从流开始投递,适合演示或从历史初始化新组。生产环境中新组读取长期存在的流时,通常从 "$" 开始,只接收此后的事件,需要历史时显式使用 XRANGE。
用 XAUTOCLAIM 恢复崩溃消费者
演示中的“Crash next 3”让所选消费者丢弃接下来三次投递而不确认,模拟处理途中崩溃。条目留在 PEL,投递计数增加。空闲至少 claimMinIdleMs 后,同组健康消费者可用自己作为目标调用 XAUTOCLAIM。ConsumerWorker.reapIdlePel 封装该模式:
async reapIdlePel() {
const { claimed, deletedIds } = await this.stream.autoclaim(
this.group,
this.name,
{ pageCount: 100, maxPages: 10 },
);
let processed = 0;
for (const [entryId, fields] of claimed) {
try {
if (this.processLatencyMs) {
await sleep(this.processLatencyMs);
}
await this._handleEntry(entryId, fields);
processed += 1;
} catch (err) {
console.error(`reap failed on ${entryId}: ${err.message}`);
}
}
this._reaped += processed;
return { claimed: claimed.length, deletedIds, processed };
}
底层 stream.autoclaim 用延续游标分页遍历 PEL:
async autoclaim(group, consumer, options = {}) {
const { pageCount = 100, startId = "0-0", maxPages = 10 } = options;
const claimedAll = [];
const deletedAll = [];
let cursor = startId;
for (let i = 0; i < maxPages; i += 1) {
const reply = await this.redis.xAutoClaim(
this.streamKey, group, consumer,
this.claimMinIdleMs, cursor, { COUNT: pageCount },
);
for (const entry of reply.messages || []) {
const tuple = toTuple(entry);
if (tuple) claimedAll.push(tuple);
}
for (const id of reply.deletedMessages || []) {
deletedAll.push(String(id));
}
const nextId = String(reply.nextId || "0-0");
if (nextId === "0-0") break;
cursor = nextId;
}
this._claimedTotal += claimedAll.length;
return { claimed: claimedAll, deletedIds: deletedAll };
}
单次调用从 startId 开始扫描 PEL,将满足空闲阈值的条目转给指定消费者,并在 reply.nextId 返回后续游标。完整扫描持续到游标回到 "0-0";maxPages 防止一次操作独占庞大 PEL。原文用 pageCount 描述扫描规模,具体扫描与返回数量规则应以 XAUTOCLAIM 命令文档为准。
每次认领增加投递次数。多次循环后可识别让每个消费者都失败的“毒丸”消息,将其送入死信流。新条目仍由 XREADGROUP > 正常投递,但反复认领坏消息会浪费消费者时间并扩大 PEL。
deletedMessages 包含载荷已被裁剪的 PEL ID,常见原因是保留策略快于慢消费者。XAUTOCLAIM 自行删除这些悬空槽位,无需再 XACK,但也不能重试,需记录并送入死信存储供检查。该第三返回值自 Redis 7.0 引入,因此示例要求7.0以上。
reapIdlePel 把认领与处理合在一起:返回条目已属于本消费者 PEL,因此由同一消费者负责处理和确认。生产中各消费者每隔数秒定时执行,避免崩溃同伴留下的条目无人处理;演示采用按钮便于等待阈值后手动触发。
XCLAIM 对已知的具体 ID 列表执行同类操作,适合接管某条卡住消息,或把指定消费者的 PEL 移交同伴。演示“Remove consumer”通过 handoverPending 处理这种情况。XAUTOCLAIM 不能按来源消费者筛选,不能用于定向移交。
用 XRANGE 重放
XRANGE 读取历史片段,与消费者组完全独立,不移动游标、不确认条目,可以从任何进程多次调用:
async replay(startId = "-", endId = "+", count = 100) {
const raw = await this.redis.xRange(this.streamKey, startId, endId, {
COUNT: count,
});
const out = [];
for (const entry of raw || []) {
const tuple = toTuple(entry);
if (tuple) out.push(tuple);
}
return out;
}
- 和 + 表示最开始与最末尾。也可传真实 ID,例如 1716998413541-0,或毫秒部分 1716998413541,后者匹配该时间戳内的条目。
常见用途:
- 初始化新投影:从
-读取全流,在搜索索引、SQL 表或其他缓存建立衍生视图。使用组会消费消息,XRANGE 则不干扰在线消费者。 - 审计近期活动:按 ID 范围读取最近几分钟,不动任何组游标。
- 调试:按 ID 获取单条消息,或读取事故时间附近的小范围,确认生产者到底写了什么。
消费者工作任务
ConsumerWorker 把“XREADGROUP → 处理 → XACK”封装为异步任务,见完整源码:
async _run() {
while (!this._stopped) {
if (this._paused) {
await sleep(50);
continue;
}
let entries;
try {
// Use the dedicated blocking client so the shared client stays
// free for HTTP-handler commands.
entries = await this.blockingStream.consume(
this.group, this.name, 10, 500,
);
} catch (err) {
console.error(`[${this.group}/${this.name}] read failed: ${err.message}`);
await sleep(500);
continue;
}
for (const [entryId, fields] of entries) {
if (this._stopped) break;
await this._dispatch(entryId, fields);
}
}
this._runPromise = null;
}
_dispatch 用 try/catch 包裹逐条处理,避免通常由 XACK 产生的错误终止循环。正常处理会确认;演示要求“崩溃”时则丢弃条目并增加计数,以便 UI 显示当前等待认领的 PEL。
Node.js 的 JavaScript 执行是单线程,因此这里的“工作线程”实际是事件循环上的异步任务。各任务的 XREADGROUP 交错执行,在自己的 Redis 连接上等待消息或超时。由此有两个设计选择:
- 每个工作者拥有专门用于阻塞读取的复制客户端。node-redis 在一个连接上串行发送命令,共享客户端会使读取或 HTTP 命令互相排队。
- 其余命令(XACK、XAUTOCLAIM、XCLAIM、XPENDING、XADD)走共享 EventStream,使生产、确认、认领计数跨工作者汇总。
恢复自己或崩溃同伴的 PEL,走独立 reapIdlePel,而不是主读取循环。它以当前消费者为目标认领,再像新条目一样处理。这是消费者各自定期或按需回收的模式。
演示按钮“XAUTOCLAIM to selected”调用所选消费者的方法;生产中应由每隔数秒的计时器调用。主循环刻意不在每次迭代执行 XREADGROUP 0 清空自己的 PEL,否则会持续重新投递待处理条目,每次重置空闲时间,使崩溃条目永远达不到认领阈值。只认领超时条目的 XAUTOCLAIM(self) 避免这一问题。
暂停和崩溃开关仅用于演示。真实消费者核心就是读取、处理、确认循环,其余是观测和演示逻辑。
前提条件
- Redis 7.0以上。XAUTOCLAIM 在6.2加入,7.0新增已删除 ID 列表,本例依赖该返回结构。
- Node.js 18以上。
- node-redis 5.x;演示 package.json 已声明依赖:
npm install
若 Redis 在其他地址,启动时传 --redis-host、--redis-port。
运行演示
取得源文件
演示由三个 JavaScript 文件与 package.json 组成,从 GitHub nodejs 源目录下载,或使用:
mkdir streaming-demo && cd streaming-demo
BASE=https://raw.githubusercontent.com/redis/docs/main/content/develop/use-cases/streaming/nodejs
curl -O $BASE/package.json
curl -O $BASE/eventStream.js
curl -O $BASE/consumerWorker.js
curl -O $BASE/demoServer.js
启动服务器
在该目录运行:
npm install
node demoServer.js
预期输出:
Deleting any existing data at key 'demo:events:orders' for a clean demo run (pass --no-reset to keep it).
Redis streaming demo server listening on http://127.0.0.1:8083
Using Redis at localhost:6379 with stream key 'demo:events:orders' (MAXLEN ~ 2000)
Seeded 3 consumer(s) across 2 group(s)
浏览器打开 http://127.0.0.1:8083,可进行:
- 生产任意数量指定或随机类型事件,观察长度和尾部更新。
- 查看各组游标、待处理数量及消费者;消费者显示处理计数、待处理数和空闲时间。
- 运行时增删消费者,观察工作重新分配。
- 点击“Crash next 3”,模拟读取后确认前崩溃,观察 XPENDING 面板。
- 等待超过默认5000毫秒阈值,选择健康消费者,点击“XAUTOCLAIM to selected”,观察重新分配和投递计数增加。
- 用 XRANGE 重放范围,确认历史独立于组状态。
- 用近似 MAXLEN 的 XTRIM 限制保留。小流的
MAXLEN ~ 50可能不删除任何内容;300条流通常裁到约100条。 - 点击“Reset demo”删除流并重新创建默认组。
生产环境使用
按长度或最小 ID 保留
演示在每次 XADD 使用 MAXLEN ~。也可:
- 用
MINID ~ <id>保留某 ID 之后的条目。保留最近24小时,可计算时间截止点并传XTRIM MINID ~ <ms>-0。 - XADD 不裁剪,另行周期执行 XTRIM。适合繁忙生产者需要减少单次开销,或复杂保留规则需要独立进程考虑组延迟。
三种方式默认均可近似裁剪。确实需要精确约束时才使用不带 ~ 的 MAXLEN 或 MINID。
避免消费者组延迟悄然增长
XINFO GROUPS 的 lag 是组未读取的条目,pending 是已投递未确认条目。生产中任一超过阈值都应告警:pending 持续增加通常表示崩溃后无人认领,lag 增长则表示消费跟不上生产。XINFO CONSUMERS 提供消费者级待处理数量与空闲时间,可识别持有条目的慢消费者。
保证幂等
XAUTOCLAIM 可在崩溃后把同一条目交给其他消费者。如果处理会发送邮件、扣款或更新下游存储,要保证重复处理安全:使用幂等键、带条件检查的 upsert,或每 ID 只执行一次的保护表。Redis Streams 本身不提供恰好一次语义。
用投递次数识别毒丸
XPENDING 返回投递次数,每次认领都会增加。消息多次投递失败,下一名消费者可能仍失败。达到阈值,例如 deliveries >= 5,将其送死信流、在原组确认并告警。否则反复认领会浪费时间、增大 PEL,并逐步触发保留/延迟告警;新消息虽仍能通过,却不能消除这一浪费。
按租户或实体分区
单个 Stream 是单键,在 Redis Cluster 中位于单个分片。吞吐超过一片能力时可按租户分流,例如 events:orders:{tenant_a} 与 events:orders:{tenant_b},让不同租户落到不同分片。需要多流原子操作时,{tenant_a} 这样的 hash tag 可使相关流位于同一分片。
按实体分区,如 events:order:{order_id},适合把每个实体的流当作其事件溯源日志:订单每次状态变更写入自身流,大小由实体生命周期约束。
每组独立消费者池
演示把全部消费者放在一个 Node.js 进程。生产中每组通常是独立 Pod 或 VM 池,使 analytics 的慢投影不会拖住 notifications。每个 Pod 可每 CPU 核运行一个消费者,或按事件循环和逐条处理吞吐容纳更多异步循环;XAUTOCLAIM 可每 N 次读取后嵌入消费者循环,或由独立回收器执行。
每次阻塞读取独占连接
XREADGROUP BLOCK 会一直占用连接到消息到达或超时。node-redis 按到达顺序发送同一 TCP 连接上的命令,共享连接会使全部命令排在尚未结束的阻塞读取后。演示为每个工作者复制阻塞客户端,其余操作走共享客户端。生产环境也应采用这一模式,或由连接池为每个阻塞调用分配独立连接。
不要用 XREAD 读取后再确认
XREAD 与 XREADGROUP 不同。前者只是跟随日志,不维护组状态、不把条目加入 PEL,因此不能 XACK。需要至少一次投递和崩溃恢复时必须使用消费者组。XREAD 仍适用于只读尾随客户端,例如事件 UI、调试器或 tail -f 风格工具,但不属于至少一次投递路径。
用 redis-cli 直接检查
测试或排障时,可以直接核对流与消费者状态:
# Stream summary
redis-cli XLEN demo:events:orders
redis-cli XINFO STREAM demo:events:orders
# Group cursors and pending counts
redis-cli XINFO GROUPS demo:events:orders
# Consumers within a group
redis-cli XINFO CONSUMERS demo:events:orders notifications
# Pending entries with idle time and delivery count
redis-cli XPENDING demo:events:orders notifications - + 20
# Tail the stream live (no consumer-group state — like tail -f)
redis-cli XREAD BLOCK 0 STREAMS demo:events:orders '$'
# Replay a range
redis-cli XRANGE demo:events:orders - + COUNT 50
组 lag 增长而消费者 idle 较短,通常说明消费者健康但生产更快,应增加消费者。pending 增长而 lag 小,则说明消费者接到消息但没有确认,可能中途崩溃,也可能确认逻辑有错误。
更多资料
本例使用:XADD 追加带近似上限的事件;XREADGROUP 读取组内新条目;XACK 确认;XAUTOCLAIM 重新分配闲置条目;XCLAIM 按已知 ID 手动移交;XRANGE 重放与审计;XPENDING 检查空闲与投递次数;XTRIM 裁剪;XGROUP CREATE 与 XGROUP DELCONSUMER 管理组和消费者;各 XINFO 命令提供可观测信息。
完整客户端参考见 node-redis 文档;消费者组、PEL、认领、限长流及 Kafka 分区的区别见Streams 概览。











暂无评论内容