> Discover all available pages from the documentation index: https://mastra.zisheng.pro/llms.txt # EventEmitterPubSub `EventEmitterPubSub` 是默认的 [`PubSub`](https://mastra.zisheng.pro/reference/pubsub/base) 实现。它使用 Node.js [`EventEmitter`](https://nodejs.org/api/events.html#class-eventemitter) 在进程内投递事件,因此无需任何外部服务即可工作。 将它用于单进程应用。有关单主机跨进程投递,请参阅 [`UnixSocketPubSub`](https://mastra.zisheng.pro/reference/pubsub/unix-socket-pubsub)。有关分布式投递,请参阅 [`RedisStreamsPubSub`](https://mastra.zisheng.pro/reference/pubsub/redis-streams) 或 [`GoogleCloudPubSub`](https://mastra.zisheng.pro/reference/pubsub/google-cloud-pubsub)。 由于它在进程内运行,事件不会持久化,也不会与其他进程共享。需要为可恢复流提供 replay 时,请使用 [`CachingPubSub`](https://mastra.zisheng.pro/reference/pubsub/caching-pubsub) 包装它。 ## 使用示例 未配置 `pubsub` 选项时会自动使用 `EventEmitterPubSub`,因此大多数应用不需要直接构造它。仅当需要配置或共享它时才显式创建。 ```typescript import { Mastra } from '@mastra/core' import { EventEmitterPubSub } from '@mastra/core/events' export const mastra = new Mastra({ pubsub: new EventEmitterPubSub(), }) ``` 要与应用的其他部分共享 emitter,请传入现有的 `EventEmitter`: ```typescript import EventEmitter from 'node:events' import { EventEmitterPubSub } from '@mastra/core/events' const emitter = new EventEmitter() const pubsub = new EventEmitterPubSub(emitter) ``` 要显示批量投递错误,请传入 logger: ```typescript import { EventEmitterPubSub } from '@mastra/core/events' const pubsub = new EventEmitterPubSub(undefined, { logger }) ``` ## 构造函数参数 **existingEmitter** (`EventEmitter`): 用于投递的现有 Node.js EventEmitter。省略时将创建新的 EventEmitter。 **options** (`EventEmitterPubSubOptions`): 可选配置。 ## 属性 **supportedModes** (`ReadonlyArray<"pull" | "push">`): 返回 \["pull", "push"]。emitter 可服务于 pull 风格 Worker,或将事件直接推送至 listener。 **supportsNativeBatching** (`boolean`): 返回 true。订阅者可以通过 options.batch 选择批量投递。 ## 方法 `EventEmitterPubSub` 实现 [`PubSub`](https://mastra.zisheng.pro/reference/pubsub/base) contract。以下方法具有此实现特有的行为。 ### `subscribe(topic, cb, options?)` 为 topic 注册 callback。没有 `options.group` 时,每个订阅者都会收到每个事件。使用 group 时,事件会在该组成员间轮询分发。 传入 `options.batch` 可选择批量投递。请参阅下面的[批处理](#batching)。 ```typescript await pubsub.subscribe('workflow.events', (event, ack, nack) => { console.log(event) }) ``` ### `flush()` 等待由 `nack` 引发的所有待处理重新投递触发后再返回。 ```typescript await pubsub.flush() ``` ### `close()` 移除所有 listener 并取消待处理的重新投递。在优雅关闭期间调用此方法。 ```typescript await pubsub.close() ``` ## 重新投递 当 grouped subscriber 调用 `nack` 时,事件会在短暂延迟后重新投递给该组,其 `deliveryAttempt` 计数会增加。调用 `ack` 会清除该事件的跟踪。fan-out subscriber 会收到无操作的 `ack` 和 `nack` 函数,因为每个事件只到达每个订阅者一次。 ## 批处理 `EventEmitterPubSub` 原生支持 `options.batch`。当订阅者选择加入时,事件会保存在每个订阅者的内存 buffer 中;满足 flush 条件后,以连续 callback 调用的形式投递。fan-out 和 group subscriber 都可批处理。完整策略请参阅 [`SubscribeBatchOptions`](https://mastra.zisheng.pro/reference/pubsub/base)。 ```typescript 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。