EventEmitterPubSub
EventEmitterPubSub 是默认的 PubSub 实现。它使用 Node.js EventEmitter 在进程内投递事件,因此无需任何外部服务即可工作。
将它用于单进程应用。有关单主机跨进程投递,请参阅 UnixSocketPubSub。有关分布式投递,请参阅 RedisStreamsPubSub 或 GoogleCloudPubSub。
由于它在进程内运行,事件不会持久化,也不会与其他进程共享。需要为可恢复流提供 replay 时,请使用 CachingPubSub 包装它。
使用示例使用示例的直接链接
未配置 pubsub 选项时会自动使用 EventEmitterPubSub,因此大多数应用不需要直接构造它。仅当需要配置或共享它时才显式创建。
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?:
options?:
属性属性的直接链接
supportedModes:
["pull", "push"]。emitter 可服务于 pull 风格 Worker,或将事件直接推送至 listener。supportsNativeBatching:
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 会收到无操作的 ack 和 nack 函数,因为每个事件只到达每个订阅者一次。
批处理批处理的直接链接
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。