跳至主要內容

PubSub

PubSub 是 Mastra 事件系統的抽象基礎類別。它定義每個發佈/訂閱後端都會實作的協定,讓 Mastra 的其他部分毋須知道正在使用哪種傳輸方式,亦能發佈及訂閱事件。

Mastra 在內部使用發佈/訂閱機制來處理 Workflow 事件、串流及跨元件通訊。大多數應用程式會使用預設的 EventEmitterPubSub,而不會直接建立 PubSub。只有在需要自訂傳輸方式時,才應實作此類別。

如要了解內置實作,請參閱 EventEmitterPubSubUnixSocketPubSubCachingPubSubRedisStreamsPubSubGoogleCloudPubSub

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

擴充 PubSub 並實作四個抽象方法,以加入自訂後端。

src/mastra/pubsub.ts
import { PubSub } from '@mastra/core/events'
import type { Event, EventCallback, SubscribeOptions } from '@mastra/core/events'

export class CustomPubSub extends PubSub {
async publish(topic: string, event: Omit<Event, 'id' | 'createdAt'>): Promise<void> {
// Deliver the event to subscribers of `topic`.
}

async subscribe(topic: string, cb: EventCallback, options?: SubscribeOptions): Promise<void> {
// Register `cb` to receive events published to `topic`.
}

async unsubscribe(topic: string, cb: EventCallback): Promise<void> {
// Remove a previously registered callback.
}

async flush(): Promise<void> {
// Wait for any in-flight deliveries to settle.
}
}

將實例傳入 Mastra 建構函式:

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { CustomPubSub } from './pubsub'

export const mastra = new Mastra({
pubsub: new CustomPubSub(),
})

傳送模式
傳送模式 的直接連結

PubSub 透過 supportedModes 屬性宣告其支援的傳送模式。Mastra 會讀取此屬性,以決定是否執行長時間運作、負責拉取事件的 worker。

模式說明
pull使用者主動從 broker 讀取資料,例如 Redis Streams XREADGROUP。Mastra 會執行協調 worker 來讀取資料。
push事件毋須使用者提出要求便會送達,可在程序內傳送,亦可透過 HTTP endpoint 傳送。不需要讀取迴圈。

預設值為 ['pull'],因此除非自訂實作選擇採用推送傳送,否則會維持現有行為。

方法
方法 的直接連結

核心方法
核心方法 的直接連結

publish(topic, event)
publishtopic-event 的直接連結

將事件發佈至主題。idcreatedAt 欄位由實作指派。

await pubsub.publish('my-topic', {
type: 'example',
data: { value: 1 },
runId: 'run-123',
})

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

註冊 callback,以接收發佈至主題的事件。設定 options.group 後,同一群組的訂閱者會競逐訊息,而每個事件只會傳送給其中一名成員。如沒有群組,每名訂閱者都會收到每個事件。

傳入 options.batch 以選用批次傳送。callback signature 維持不變:包含 N 個事件的批次會按發佈次序,透過 N 次連續的 cb(event, ack, nack) 呼叫傳送。只有後端的 supportsNativeBatchingtrue 時,才會採用批次處理。其他後端會忽略此選項,並逐一傳送事件。

await pubsub.subscribe('my-topic', (event, ack, nack) => {
console.log(event)
})

unsubscribe(topic, cb)
unsubscribetopic-cb 的直接連結

從主題移除先前註冊的 callback。

await pubsub.unsubscribe('my-topic', callback)

flush()
flush 的直接連結

等待所有進行中的傳送完成。關閉前請呼叫此方法,以免遺失事件。

await pubsub.flush()

clearTopic(topic)
cleartopictopic 的直接連結

當不會再有事件發佈至某個主題後,刪除該主題保留的所有狀態(快取記錄、持久化串流項目及使用者群組)。當一次執行到達終止狀態時,Mastra 的執行生命週期(durable Agent 及事件驅動的 Workflow 引擎)會自動呼叫此方法,因此在會保留訊息的傳輸方式上,個別執行的主題不會持續累積。

預設實作不會執行任何操作:不會為個別主題保留資料的傳輸方式(例如 EventEmitterPubSub)沒有需要清除的內容。會持久保存訊息的後端(例如 RedisStreamsPubSub)則會覆寫此方法。此協定採取盡力而為的方式:實作會記錄失敗,而不會拋出錯誤,因為呼叫者會在清理邊界以不等待結果的方式呼叫此方法。

await pubsub.clearTopic('workflow.events.v2.run-123')

重播方法
重播方法 的直接連結

這些方法支援在中斷連線後恢復串流。預設實作會退回一般的 subscribe,因此不支援記錄的後端只會處理即時事件。CachingPubSub 會覆寫這些方法,以重播已快取的事件。

getHistory(topic, offset?)
gethistorytopic-offset 的直接連結

傳回主題從 offset 開始的已快取事件。如後端沒有記錄,則傳回空陣列。

const events = await pubsub.getHistory('my-topic', 0)

傳回:Promise<Event[]>

subscribeWithReplay(topic, cb)
subscribewithreplaytopic-cb 的直接連結

重播已快取的事件,然後訂閱即時事件。

await pubsub.subscribeWithReplay('my-topic', event => {
console.log(event)
})

subscribeFromOffset(topic, offset, cb)
subscribefromoffsettopic-offset-cb 的直接連結

從已知位置開始重播已快取的事件,然後訂閱即時事件。當 client 知道其最後位置時,這比完整重播更有效率。

await pubsub.subscribeFromOffset('my-topic', 42, event => {
console.log(event)
})

屬性
屬性 的直接連結

supportedModes:

ReadonlyArray<"pull" | "push">
實作支援的傳送模式。預設為 ["pull"]

supportsNativeBatching:

boolean
實作是否採用 options.batch 進行 subscribe()。預設為 false。內部整合批次處理的後端會覆寫此屬性並傳回 true

類型
類型 的直接連結

Event
event 的直接連結

type:

string
事件類型識別碼。

id:

string
唯一事件 ID,由實作在發佈時指派。

data:

any
事件 payload。

runId:

string
事件所屬的執行。

createdAt:

Date
由實作在發佈時指派的時間戳記。

index?:

number
用於從特定 offset 恢復的順序位置。

deliveryAttempt?:

number
事件已傳送的次數,由 1 開始。後端不追蹤重新傳送時,預設為 1。

SubscribeOptions
subscribeoptions 的直接連結

group?:

string
設定後,同一群組的訂閱者會競逐訊息,而每個事件只會傳送給其中一名成員。如省略此項,每名訂閱者都會收到每個事件。

batch?:

SubscribeBatchOptions
為此訂閱選用批次傳送。如省略此項,事件會逐一傳送。只有 supportsNativeBatchingtrue 的後端才會採用此設定。

SubscribeBatchOptions
subscribebatchoptions 的直接連結

每個訂閱的批次處理政策。callback signature 不會改變。包含 N 個事件的批次會按發佈次序轉化為 N 次連續的 callback 呼叫。

maxSize?:

number
強制 flush 前最多保留的事件數目。

maxWaitMs?:

number
最舊事件可在 buffer 中停留的最長時間(毫秒)。buffer 從空白變為非空白時,計時器便會啟動。

minIntervalMs?:

number
連續兩次批次傳送之間的最短時間(毫秒)。即使 maxSizemaxWaitMs 會觸發傳送,buffer 亦會保留內容,直至距離上次傳送已經過此段時間。

isImmediate?:

(event: Event) => boolean
如為事件傳回 true,buffer 會在發佈時立即 flush,但仍受 minIntervalMs 限制。這是按事件使用的例外機制。

coalesce?:

(events: Event[]) => Event[]
在傳送前套用至排隊中的批次,以捨棄已被取代的事件。必須按 reference identity 傳回輸入內容的子集;傳回新建的 Event 物件會違反協定,並令整個批次被捨棄。保留事件的次序不變。

maxBufferSize?:

number
= 256
觸發 overflow 處理前,buffer 可保留的事件數目上限。標記為立即處理的事件絕不會因 overflow 而被捨棄。

overflow?:

"drop-oldest" | "drop-newest" | "coalesce-or-drop-oldest"
= coalesce-or-drop-oldest
buffer 超出 maxBufferSize 時的 overflow 策略。coalesce-or-drop-oldest 會先執行 coalesce,如仍然超出上限,便捨棄最舊的事件。

EventCallback
eventcallback 的直接連結

訂閱者的 callback signature:(event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void

event:

Event
已傳送的事件。

ack?:

() => Promise<void>
確認處理成功。事件會從佇列移除。

nack?:

() => Promise<void>
否定確認。事件會重新加入佇列,並在延遲後再次傳送。