跳至主要內容

EventEmitterPubSub

EventEmitterPubSub 是預設的 PubSub 實作。它使用 Node.js EventEmitter 在 process 內傳遞 event,因此不需要任何外部服務即可運作。

適用於單一 process 應用程式。若要在同一主機的多個 process 之間傳遞,請參閱 UnixSocketPubSub。若要進行分散式傳遞,請參閱 RedisStreamsPubSubGoogleCloudPubSub

由於它在 process 內運作,event 不會持久化,也不會與其他 process 共用。需要為可恢復串流提供重播功能時,請使用 CachingPubSub 包裝它。

使用範例
「使用範例」的直接連結

未設定 pubsub 選項時,系統會自動使用 EventEmitterPubSub,因此大多數應用程式不會直接建構它。只有需要設定或共用時,才明確建立 instance。

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 })

Constructor 參數
「Constructor 參數」的直接連結

existingEmitter?:

EventEmitter
用於傳遞的現有 Node.js EventEmitter。省略時會建立新的 EventEmitter。

options?:

EventEmitterPubSubOptions
選填設定。
IMastraLogger

屬性
「屬性」的直接連結

supportedModes:

ReadonlyArray<"pull" | "push">
傳回 ["pull", "push"]。emitter 可供 pull 型 worker 使用,也能將 event 直接 push 給 listener。

supportsNativeBatching:

boolean
傳回 true。subscriber 可透過 options.batch 選擇使用批次傳遞。

方法
「方法」的直接連結

EventEmitterPubSub 實作 PubSub contract。以下方法具有此實作專屬的行為。

subscribe(topic, cb, options?)
「subscribetopic-cb-options」的直接連結

為 topic 註冊 callback。沒有 options.group 時,每個 subscriber 都會收到每個 event。設定 group 後,event 會以 round-robin 方式分配給該 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()

重新傳遞
「重新傳遞」的直接連結

分組 subscriber 呼叫 nack 時,event 會在短暫延遲後重新傳遞給 group,且 deliveryAttempt 次數會增加。呼叫 ack 會清除該 event 的追蹤。由於每個 event 都只會送達每個 fan-out subscriber 一次,因此 fan-out subscriber 收到的 acknack 函式不會執行任何操作。

批次處理
「批次處理」的直接連結

EventEmitterPubSub 原生支援 options.batch。subscriber 選用後,event 會保留在各 subscriber 專用的記憶體內 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 位於記憶體內且只屬於單一 process,因此批次 state 不會持久化,也無法在重新啟動後保留。由 maxWaitMs 觸發的 flush 採盡力而為:如果擲回例外的 coalesce 等步驟失敗,錯誤會透過設定的 logger 顯示,而不會擲回。flush() 會先清空每個批次 subscriber buffer,再完成解析。