EventEmitterPubSub
EventEmitterPubSub 是預設的 PubSub 實作。它使用 Node.js EventEmitter 在進程內傳送事件,因此無須任何外部服務即可運作。
適用於單一進程應用程式。在單一主機上跨進程傳送,請參閱 UnixSocketPubSub。分散式傳送請參閱 RedisStreamsPubSub 或 GoogleCloudPubSub。
由於它在進程內運作,事件不會持久保存,也不會與其他進程共用。若可恢復串流需要重播,請使用 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,或直接將事件 push 至 listener。supportsNativeBatching:
true。訂閱者可透過 options.batch 選用批次傳送。方法方法 的直接連結
EventEmitterPubSub 實作 PubSub 合約。以下方法具有此實作特有的行為。
subscribe(topic, cb, options?)subscribetopic-cb-options 的直接連結
為主題註冊 callback。沒有 options.group 時,每個訂閱者都會收到所有事件。設有群組時,事件會以 round-robin 方式分配給該群組的成員。
傳入 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()
重新傳送重新傳送 的直接連結
群組訂閱者呼叫 nack 時,事件會在短暫延遲後重新傳送至群組,並增加其 deliveryAttempt 計數。呼叫 ack 會清除此事件的追蹤資料。Fan-out 訂閱者收到的 ack 和 nack 函數不會執行任何操作,因為每個事件只會傳送給每個訂閱者一次。
批次處理批次處理 的直接連結
EventEmitterPubSub 原生支援 options.batch。訂閱者選用後,事件會保留在每個訂閱者各自的記憶體緩衝區,並在符合 flush 條件時以連續 callback 呼叫傳送。Fan-out 和群組訂閱者都可使用批次處理。完整政策請參閱 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
},
},
)
緩衝區位於每個進程的記憶體內,因此批次狀態不會持久保存,也無法在重新啟動後保留。由 maxWaitMs 觸發的 flush 屬於盡力而為:如果某個步驟失敗,例如 coalesce 拋出錯誤,錯誤會透過已設定的 logger 顯示,而不會向外拋出。flush() 會在完成前清空每個批次訂閱者的緩衝區。