RedisStreamsPubSub
RedisStreamsPubSub 是由 Redis Streams 支持的 PubSub 实现。它可跨进程和主机投递事件,并支持持久化、消费者组及失败后重新投递。它还实现了 LeaseProvider,因此 signals 层可以跨实例为每个资源选出单一所有者,从而让 signals 能够在分布式和 serverless 部署中协调运行。
适用于多个服务共享事件流的分布式部署。单进程投递请使用 EventEmitterPubSub。Google Cloud 请使用 GoogleCloudPubSub。
每个主题映射到一个 Redis 流键。带有组的订阅使用 Redis 消费者组,因此成员以轮询方式分担工作。不带组的订阅会创建私有消费者组,因此每个订阅者都会收到每个事件。
RedisStreamsPubSub 是拉取传输:消费者使用 XREADGROUP 读取事件,因此 Mastra 会运行编排 worker 代为读取。
安装安装的直接链接
- npm
- pnpm
- Yarn
- Bun
npm install @mastra/redis-streams
pnpm add @mastra/redis-streams
yarn add @mastra/redis-streams
bun add @mastra/redis-streams
使用示例使用示例的直接链接
提供 Redis 连接 URL。
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'
export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})
构造函数参数构造函数参数的直接链接
url?:
redisOptions.url。keyPrefix?:
<keyPrefix>:<topic>。blockMs?:
redisOptions?:
redis 客户端以进行高级配置的选项。maxStreamLength?:
streamIdleTtlMs?:
clearTopic 会处理正常生命周期结束时的删除;此选项仅限制永远不会调用 clearTopic 的流(例如崩溃的运行)的内存。必须为非负整数。默认为 0(禁用)。reclaimIntervalMs?:
reclaimIdleMs?:
maxDeliveryAttempts?:
nack 重新投递的最大次数,超过后将被丢弃。传入 Infinity 可禁用上限。logger?:
属性属性的直接链接
supportedModes:
["pull"]。方法方法的直接链接
RedisStreamsPubSub 实现 PubSub 契约。以下方法的行为特定于此实现。
subscribe(topic, cb, options?)subscribetopic-cb-options的直接链接
订阅主题。使用 options.group 时,组成员通过 Redis 消费者组共享事件。没有组时,订阅者通过私有消费者组接收每个事件。
await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})
flush()flush的直接链接
等待正在进行的发布完成。
await pubsub.flush()
clearTopic(topic)cleartopictopic的直接链接
删除主题的流及其上的所有消费者组,释放已完成主题原本会占用的内存。Mastra 的运行生命周期(持久化 agent 和事件驱动的 workflow 引擎)会在运行达到终止状态时自动调用它。仅在不会再读取该主题时自行调用。它是尽力而为的操作,绝不会抛出错误。失败会以 warn 级别记录。删除流时仍附加的订阅者会自行恢复,但会错过已删除的条目。
自动清理要求 @mastra/core 和 @mastra/redis-streams 版本都支持 clearTopic:运行时通过其缓存层路由调用,因此请同时升级两个包以获得运行结束时的流删除。
await pubsub.clearTopic('workflow.events.run-123')
close()close的直接链接
关闭 Redis 连接并停止所有订阅。请在优雅关闭期间调用此方法。
await pubsub.close()
重新投递与重新认领重新投递与重新认领的直接链接
当订阅者调用 nack 时,事件会以递增的 deliveryAttempt 重新发布,原事件会被确认。事件达到 maxDeliveryAttempts 后将被丢弃而非重新投递。此外,每个订阅会定期重新认领组中先前消费者读取但从未确认的事件,由 reclaimIntervalMs 和 reclaimIdleMs 控制。
分布式租约分布式租约的直接链接
RedisStreamsPubSub 在同一 Redis 连接上实现 LeaseProvider 契约。signals 运行时 使用它选出单一所有者(通常按线程键),因此跨实例时只有一个进程会唤醒并运行 agent,其他进程会将后续工作路由给持有者。这使 signals 可在 serverless 和多实例部署中工作;没有共享租约时,每个实例都会启动各自竞争的运行。
租约键与主题使用相同的 keyPrefix 命名空间,格式为 <keyPrefix>:lease:<key>。所有操作都是原子的:acquireLease 使用 SET NX PX 并以幂等方式刷新自身 TTL,而 releaseLease、renewLease 和 transferLease 使用 Lua 脚本,在变更前检查所有权,因此不会覆盖其他所有者的并发续约。
无需直接调用这些方法。将 RedisStreamsPubSub 配置为 pubsub 后端即可让运行时检测并使用该能力。完整的方法契约请参阅 LeaseProvider。