跳至主要內容

EventEmitterPubSub

EventEmitterPubSub 是預設的 PubSub 實作。它使用 Node.js EventEmitter 在進程內傳送事件,因此無須任何外部服務即可運作。

適用於單一進程應用程式。在單一主機上跨進程傳送,請參閱 UnixSocketPubSub。分散式傳送請參閱 RedisStreamsPubSubGoogleCloudPubSub

由於它在進程內運作,事件不會持久保存,也不會與其他進程共用。若可恢復串流需要重播,請使用 CachingPubSub 包裝它。

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

如果你沒有設定 pubsub 選項,系統會自動使用 EventEmitterPubSub,因此大部分應用程式都無須直接建構它。只有在需要設定或共用它時,才明確建立執行個體。

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

建構函數參數
建構函數參數 的直接連結

existingEmitter?:

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

options?:

EventEmitterPubSubOptions
可選設定。
IMastraLogger

屬性
屬性 的直接連結

supportedModes:

ReadonlyArray<"pull" | "push">
傳回 ["pull", "push"]。emitter 可以服務 pull 式 worker,或直接將事件 push 至 listener。

supportsNativeBatching:

boolean
傳回 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 訂閱者收到的 acknack 函數不會執行任何操作,因為每個事件只會傳送給每個訂閱者一次。

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

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() 會在完成前清空每個批次訂閱者的緩衝區。