使用 node-redis 构建 Redis 事件流

本指南用 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 独立处理完整事件流。

工作流程

  1. 应用调用 stream.produce(eventType, payload),执行带近似 MAXLEN ~ 上限的 XADD,Redis 分配按时间排序的 ID。
  2. 每个消费者的异步循环用特殊 ID > 调用 XREADGROUP,表示读取该组尚未投递给任何人的条目,并设置较短阻塞超时。
  3. 处理后调用 XACK,让 Redis 将条目移出组的待处理列表。
  4. 消费者在确认前被终止或崩溃,条目留在 PEL;定期 XAUTOCLAIM 扫描将闲置条目转给健康消费者。
  5. 任何代码,包括组外代码,都可用 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 概览。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容