跳到主要内容

EventEmitterPubSub

EventEmitterPubSub 是默认的 PubSub 实现。它使用 Node.js EventEmitter 在进程内投递事件,因此无需任何外部服务即可工作。

将它用于单进程应用。有关单主机跨进程投递,请参阅 UnixSocketPubSub。有关分布式投递,请参阅 RedisStreamsPubSubGoogleCloudPubSub

由于它在进程内运行,事件不会持久化,也不会与其他进程共享。需要为可恢复流提供 replay 时,请使用 CachingPubSub 包装它。

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

未配置 pubsub 选项时会自动使用 EventEmitterPubSub,因此大多数应用不需要直接构造它。仅当需要配置或共享它时才显式创建。

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { EventEmitterPubSub } from '@mastra/core/events'

export const mastra = new Mastra({
pubsub: new EventEmitterPubSub(),
})

要与应用的其他部分共享 emitter,请传入现有的 EventEmitter

import EventEmitter from 'node:events'
import { EventEmitterPubSub } from '@mastra/core/events'

const emitter = new EventEmitter()
const pubsub = new EventEmitterPubSub(emitter)

要显示批量投递错误,请传入 logger:

import { EventEmitterPubSub } from '@mastra/core/events'

const pubsub = new EventEmitterPubSub(undefined, { logger })

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

existingEmitter?:

EventEmitter
用于投递的现有 Node.js EventEmitter。省略时将创建新的 EventEmitter。

options?:

EventEmitterPubSubOptions
可选配置。
IMastraLogger

属性
属性的直接链接

supportedModes:

ReadonlyArray<"pull" | "push">
返回 ["pull", "push"]。emitter 可服务于 pull 风格 Worker,或将事件直接推送至 listener。

supportsNativeBatching:

boolean
返回 true。订阅者可以通过 options.batch 选择批量投递。

方法
方法的直接链接

EventEmitterPubSub 实现 PubSub contract。以下方法具有此实现特有的行为。

subscribe(topic, cb, options?)
subscribetopic-cb-options的直接链接

为 topic 注册 callback。没有 options.group 时,每个订阅者都会收到每个事件。使用 group 时,事件会在该组成员间轮询分发。

传入 options.batch 可选择批量投递。请参阅下面的批处理

await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})

flush()
flush的直接链接

等待由 nack 引发的所有待处理重新投递触发后再返回。

await pubsub.flush()

close()
close的直接链接

移除所有 listener 并取消待处理的重新投递。在优雅关闭期间调用此方法。

await pubsub.close()

重新投递
重新投递的直接链接

当 grouped subscriber 调用 nack 时,事件会在短暂延迟后重新投递给该组,其 deliveryAttempt 计数会增加。调用 ack 会清除该事件的跟踪。fan-out subscriber 会收到无操作的 acknack 函数,因为每个事件只到达每个订阅者一次。

批处理
批处理的直接链接

EventEmitterPubSub 原生支持 options.batch。当订阅者选择加入时,事件会保存在每个订阅者的内存 buffer 中;满足 flush 条件后,以连续 callback 调用的形式投递。fan-out 和 group subscriber 都可批处理。完整策略请参阅 SubscribeBatchOptions

await pubsub.subscribe(
'workflow.events',
event => {
console.log(event)
},
{
batch: {
maxSize: 10, // flush once 10 events have queued
maxWaitMs: 500, // ...or after 500ms, whichever comes first
},
},
)

buffer 位于内存中且按进程隔离,因此批处理状态不会持久化,也无法在重启后保留。由 maxWaitMs 触发的 flush 是尽力而为的:如果抛出异常的 coalesce 等步骤失败,错误会通过配置的 logger 显示,而不会抛出。flush() 会在返回前清空每个批处理订阅者的 buffer。