PubSub
PubSub 是 Mastra event 系統的抽象基底類別。它定義每個 pub/sub 後端都必須實作的 contract,因此 Mastra 的其他部分不需要知道所使用的 transport,也能發布及訂閱 event。
Mastra 在內部使用 pub/sub 處理 Workflow event、串流與跨元件通訊。大多數應用程式使用預設的 EventEmitterPubSub,不會直接建構 PubSub。只有需要自訂 transport 時才實作此類別。
內建實作請參閱 EventEmitterPubSub、UnixSocketPubSub、CachingPubSub、RedisStreamsPubSub和 GoogleCloudPubSub。
使用範例「使用範例」的直接連結
擴充 PubSub 並實作四個抽象方法,以新增自訂後端。
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:
import { Mastra } from '@mastra/core'
import { CustomPubSub } from './pubsub'
export const mastra = new Mastra({
pubsub: new CustomPubSub(),
})
傳遞模式「傳遞模式」的直接連結
PubSub 會透過 supportedModes 屬性宣告支援的傳遞模式。Mastra 會讀取此屬性,以決定是否執行長時間存活並拉取 event 的 worker。
| 模式 | 說明 |
|---|---|
pull | consumer 主動從 broker 讀取,例如 Redis Streams XREADGROUP。Mastra 會執行 orchestration worker 進行讀取。 |
push | event 不需 consumer 要求便會送達,可在 process 內傳送或透過 HTTP endpoint 傳送。不需要 read loop。 |
預設為 ['pull'],除非自訂實作選擇使用 push 傳遞,否則會維持目前的行為。
方法「方法」的直接連結
核心方法「核心方法」的直接連結
publish(topic, event)「publishtopic-event」的直接連結
將 event 發布至 topic。id 與 createdAt 欄位由實作指派。
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) 呼叫傳遞。只有後端的 supportsNativeBatching 為 true 時才會採用批次處理。其他後端會忽略此選項,逐一傳遞 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:
["pull"]。supportsNativeBatching:
subscribe() 上的 options.batch。預設為 false。內部整合批次處理的後端會覆寫此屬性並傳回 true。型別「型別」的直接連結
Event「event」的直接連結
type:
id:
data:
runId:
createdAt:
index?:
deliveryAttempt?:
SubscribeOptions「subscribeoptions」的直接連結
group?:
batch?:
supportsNativeBatching 為 true 的後端才會採用。SubscribeBatchOptions「subscribebatchoptions」的直接連結
各 subscription 專用的批次政策。callback signature 不會改變。一批 N 個 event 會依發布順序轉換為 N 次連續 callback 呼叫。
maxSize?:
maxWaitMs?:
minIntervalMs?:
maxSize 或 maxWaitMs 已觸發,buffer 仍會保留至上次傳遞後已經過此間隔。isImmediate?:
true 時,buffer 會在發布時立即 flush,但仍受 minIntervalMs 限制。這是各 event 專用的例外處理機制。coalesce?:
Event 物件會違反 contract,並捨棄整個批次。保留的 event 順序不變。maxBufferSize?:
overflow?:
maxBufferSize 時的 overflow 策略。coalesce-or-drop-oldest 會先執行 coalesce,若仍超出限制,再捨棄最舊項目。EventCallback「eventcallback」的直接連結
subscriber 的 callback signature:(event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void。