跳到主要内容

RedisStreamsPubSub

RedisStreamsPubSub 是由 Redis Streams 支持的 PubSub 实现。它可跨进程和主机投递事件,并支持持久化、消费者组及失败后重新投递。它还实现了 LeaseProvider,因此 signals 层可以跨实例为每个资源选出单一所有者,从而让 signals 能够在分布式和 serverless 部署中协调运行。

适用于多个服务共享事件流的分布式部署。单进程投递请使用 EventEmitterPubSub。Google Cloud 请使用 GoogleCloudPubSub

每个主题映射到一个 Redis 流键。带有组的订阅使用 Redis 消费者组,因此成员以轮询方式分担工作。不带组的订阅会创建私有消费者组,因此每个订阅者都会收到每个事件。

RedisStreamsPubSub 是拉取传输:消费者使用 XREADGROUP 读取事件,因此 Mastra 会运行编排 worker 代为读取。

安装
安装的直接链接

npm install @mastra/redis-streams

使用示例
使用示例的直接链接

提供 Redis 连接 URL。

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'

export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})

构造函数参数
构造函数参数的直接链接

url?:

string
= redis://localhost:6379
Redis 连接 URL。回退使用 redisOptions.url

keyPrefix?:

string
= mastra:topic
流键的前缀。每个主题映射到 <keyPrefix>:<topic>

blockMs?:

number
= 1000
每次读取在等待新事件时阻塞的时间(毫秒)。

redisOptions?:

RedisClientOptions
传递给底层 redis 客户端以进行高级配置的选项。

maxStreamLength?:

number
= 10000
每个流保留条目的近似最大数量。设为 0 可禁用修剪。

streamIdleTtlMs?:

number
= 0
以毫秒计的空闲过期时间:每次写入时都会刷新的滑动 TTL(发布、nack 重试、重新创建组)。每次写入都会重置它,因此持续写入的流不会在传输过程中到期;完整时长内保持空闲的流会由 Redis 自动删除。请注意,只有写入会刷新 TTL,缓慢清空积压的消费者不会刷新,因此应将其设为远高于活动主题两次写入之间的最长预期间隔。这是兜底机制,而不是主要清理方式:clearTopic 会处理正常生命周期结束时的删除;此选项仅限制永远不会调用 clearTopic 的流(例如崩溃的运行)的内存。必须为非负整数。默认为 0(禁用)。

reclaimIntervalMs?:

number
= 30000
订阅重新认领先前消费者读取但从未确认的事件的间隔时间(毫秒)。设为 0 可禁用。

reclaimIdleMs?:

number
= 60000
待处理事件可被重新认领前的最短空闲时间(毫秒)。应显著高于典型处理时间,以避免重复投递。

maxDeliveryAttempts?:

number
= 5
事件通过 nack 重新投递的最大次数,超过后将被丢弃。传入 Infinity 可禁用上限。

logger?:

{ debug?: Function; warn?: Function }
用于诊断的可选 logger。省略时,被抑制的错误不会输出。

属性
属性的直接链接

supportedModes:

ReadonlyArray<"pull" | "push">
返回 ["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 后将被丢弃而非重新投递。此外,每个订阅会定期重新认领组中先前消费者读取但从未确认的事件,由 reclaimIntervalMsreclaimIdleMs 控制。

分布式租约
分布式租约的直接链接

RedisStreamsPubSub 在同一 Redis 连接上实现 LeaseProvider 契约。signals 运行时 使用它选出单一所有者(通常按线程键),因此跨实例时只有一个进程会唤醒并运行 agent,其他进程会将后续工作路由给持有者。这使 signals 可在 serverless 和多实例部署中工作;没有共享租约时,每个实例都会启动各自竞争的运行。

租约键与主题使用相同的 keyPrefix 命名空间,格式为 <keyPrefix>:lease:<key>。所有操作都是原子的:acquireLease 使用 SET NX PX 并以幂等方式刷新自身 TTL,而 releaseLeaserenewLeasetransferLease 使用 Lua 脚本,在变更前检查所有权,因此不会覆盖其他所有者的并发续约。

无需直接调用这些方法。将 RedisStreamsPubSub 配置为 pubsub 后端即可让运行时检测并使用该能力。完整的方法契约请参阅 LeaseProvider