跳至主要內容

PubSub

PubSub 是 Mastra event 系統的抽象基底類別。它定義每個 pub/sub 後端都必須實作的 contract,因此 Mastra 的其他部分不需要知道所使用的 transport,也能發布及訂閱 event。

Mastra 在內部使用 pub/sub 處理 Workflow event、串流與跨元件通訊。大多數應用程式使用預設的 EventEmitterPubSub,不會直接建構 PubSub。只有需要自訂 transport 時才實作此類別。

內建實作請參閱 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.
}
}

將 instance 傳給 Mastra constructor:

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

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

傳遞模式
「傳遞模式」的直接連結

PubSub 會透過 supportedModes 屬性宣告支援的傳遞模式。Mastra 會讀取此屬性,以決定是否執行長時間存活並拉取 event 的 worker。

模式說明
pullconsumer 主動從 broker 讀取,例如 Redis Streams XREADGROUP。Mastra 會執行 orchestration worker 進行讀取。
pushevent 不需 consumer 要求便會送達,可在 process 內傳送或透過 HTTP endpoint 傳送。不需要 read loop。

預設為 ['pull'],除非自訂實作選擇使用 push 傳遞,否則會維持目前的行為。

方法
「方法」的直接連結

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

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

將 event 發布至 topic。idcreatedAt 欄位由實作指派。

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

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

註冊 callback,以接收發布至 topic 的 event。設定 options.group 後,同一 group 中的 subscriber 會競爭訊息,每個 event 只會傳遞給一個成員。未設定 group 時,每個 subscriber 都會收到每個 event。

傳入 options.batch 以選用批次傳遞。callback signature 不變:一批 N 個 event 會依發布順序,透過 N 次連續的 cb(event, ack, nack) 呼叫傳遞。只有後端的 supportsNativeBatchingtrue 時才會採用批次處理。其他後端會忽略此選項,逐一傳遞 event。

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

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

從 topic 移除先前註冊的 callback。

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

flush()
「flush」的直接連結

等待所有進行中的傳遞完成。請在關閉前呼叫,以免 event 遺失。

await pubsub.flush()

clearTopic(topic)
「cleartopictopic」的直接連結

當 topic 不再發布 event 時,刪除其所有保留 state(快取歷史記錄、持久化 stream entry 與 consumer group)。Mastra 的 run 生命週期(durable Agent 與 event 型 Workflow engine)會在 run 進入終止 state 時自動呼叫此方法,因此每個 run 的 topic 不會在保留訊息的 transport 上持續累積。

預設實作不會執行任何操作:不會按 topic 保留任何內容的 transport(例如 EventEmitterPubSub)沒有需要清除的資料。持久化訊息的後端(例如 RedisStreamsPubSub)會覆寫此方法。此 contract 採盡力而為:實作會記錄失敗,而不會擲回例外,因為 caller 會在清理邊界以 fire-and-forget 方式呼叫。

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

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

這些方法支援斷線後恢復串流。預設實作會改用一般 subscribe,因此不支援歷史記錄的後端只會提供即時內容。CachingPubSub 會覆寫這些方法,以重播快取 event。

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

傳回 topic 從 offset 開始的快取 event。後端沒有歷史記錄時傳回空陣列。

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

傳回:Promise<Event[]>

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

重播快取 event,再訂閱即時 event。

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

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

從已知位置開始重播快取 event,再訂閱即時 event。client 知道最後位置時,這比完整重播更有效率。

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

屬性
「屬性」的直接連結

supportedModes:

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

supportsNativeBatching:

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

型別
「型別」的直接連結

Event
「event」的直接連結

type:

string
Event 型別識別碼。

id:

string
不重複 event ID,由實作在發布時指派。

data:

any
Event payload。

runId:

string
Event 所屬的 run。

createdAt:

Date
由實作在發布時指派的 timestamp。

index?:

number
用於從特定 offset 恢復的連續位置。

deliveryAttempt?:

number
Event 已傳遞的次數。從 1 開始。後端未追蹤重新傳遞時預設為 1。

SubscribeOptions
「subscribeoptions」的直接連結

group?:

string
設定後,同一 group 中的 subscriber 會競爭訊息,每個 event 只會傳遞給一個成員。省略時,每個 subscriber 都會收到每個 event。

batch?:

SubscribeBatchOptions
為此 subscription 選用批次傳遞。省略時會逐一傳遞 event。只有 supportsNativeBatchingtrue 的後端才會採用。

SubscribeBatchOptions
「subscribebatchoptions」的直接連結

各 subscription 專用的批次政策。callback signature 不會改變。一批 N 個 event 會依發布順序轉換為 N 次連續 callback 呼叫。

maxSize?:

number
強制 flush 前最多保留的 event 數量。

maxWaitMs?:

number
最舊 event 可停留在 buffer 中的最長毫秒數。buffer 從空白變為非空白時,timer 便會啟動。

minIntervalMs?:

number
連續批次傳遞之間的最短毫秒數。即使 maxSizemaxWaitMs 已觸發,buffer 仍會保留至上次傳遞後已經過此間隔。

isImmediate?:

(event: Event) => boolean
對 event 傳回 true 時,buffer 會在發布時立即 flush,但仍受 minIntervalMs 限制。這是各 event 專用的例外處理機制。

coalesce?:

(events: Event[]) => Event[]
在傳遞前套用至佇列中的批次,以捨棄已被取代的 event。必須依參照識別傳回輸入的子集合;傳回新建構的 Event 物件會違反 contract,並捨棄整個批次。保留的 event 順序不變。

maxBufferSize?:

number
= 256
觸發 overflow 處理前,buffer 最多可容納的 event 數量。標示為立即處理的 event 絕不會因 overflow 遭到捨棄。

overflow?:

"drop-oldest" | "drop-newest" | "coalesce-or-drop-oldest"
= coalesce-or-drop-oldest
buffer 超過 maxBufferSize 時的 overflow 策略。coalesce-or-drop-oldest 會先執行 coalesce,若仍超出限制,再捨棄最舊項目。

EventCallback
「eventcallback」的直接連結

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

event:

Event
已傳遞的 event。

ack?:

() => Promise<void>
確認處理成功。event 會從 queue 中移除。

nack?:

() => Promise<void>
否定確認。event 會重新排入 queue,並在延遲後重新傳遞。