EventEmitterPubSub
EventEmitterPubSub 是預設的 PubSub 實作。它使用 Node.js EventEmitter 在 process 內傳遞 event,因此不需要任何外部服務即可運作。
適用於單一 process 應用程式。若要在同一主機的多個 process 之間傳遞,請參閱 UnixSocketPubSub。若要進行分散式傳遞,請參閱 RedisStreamsPubSub 或 GoogleCloudPubSub。
由於它在 process 內運作,event 不會持久化,也不會與其他 process 共用。需要為可恢復串流提供重播功能時,請使用 CachingPubSub 包裝它。
使用範例「使用範例」的直接連結
未設定 pubsub 選項時,系統會自動使用 EventEmitterPubSub,因此大多數應用程式不會直接建構它。只有需要設定或共用時,才明確建立 instance。
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?:
options?:
屬性「屬性」的直接連結
supportedModes:
["pull", "push"]。emitter 可供 pull 型 worker 使用,也能將 event 直接 push 給 listener。supportsNativeBatching:
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 收到的 ack 與 nack 函式不會執行任何操作。
批次處理「批次處理」的直接連結
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,再完成解析。