PubSub
PubSub 是 Mastra 事件系統的抽象基礎類別。它定義每個發佈/訂閱後端都會實作的協定,讓 Mastra 的其他部分毋須知道正在使用哪種傳輸方式,亦能發佈及訂閱事件。
Mastra 在內部使用發佈/訂閱機制來處理 Workflow 事件、串流及跨元件通訊。大多數應用程式會使用預設的 EventEmitterPubSub,而不會直接建立 PubSub。只有在需要自訂傳輸方式時,才應實作此類別。
如要了解內置實作,請參閱 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.
}
}
將實例傳入 Mastra 建構函式:
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 的直接連結
將事件發佈至主題。id 及 createdAt 欄位由實作指派。
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) 呼叫傳送。只有後端的 supportsNativeBatching 為 true 時,才會採用批次處理。其他後端會忽略此選項,並逐一傳送事件。
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:
["pull"]。supportsNativeBatching:
options.batch 進行 subscribe()。預設為 false。內部整合批次處理的後端會覆寫此屬性並傳回 true。類型類型 的直接連結
Eventevent 的直接連結
type:
id:
data:
runId:
createdAt:
index?:
deliveryAttempt?:
SubscribeOptionssubscribeoptions 的直接連結
group?:
batch?:
supportsNativeBatching 為 true 的後端才會採用此設定。SubscribeBatchOptionssubscribebatchoptions 的直接連結
每個訂閱的批次處理政策。callback signature 不會改變。包含 N 個事件的批次會按發佈次序轉化為 N 次連續的 callback 呼叫。
maxSize?:
maxWaitMs?:
minIntervalMs?:
maxSize 或 maxWaitMs 會觸發傳送,buffer 亦會保留內容,直至距離上次傳送已經過此段時間。isImmediate?:
true,buffer 會在發佈時立即 flush,但仍受 minIntervalMs 限制。這是按事件使用的例外機制。coalesce?:
Event 物件會違反協定,並令整個批次被捨棄。保留事件的次序不變。maxBufferSize?:
overflow?:
maxBufferSize 時的 overflow 策略。coalesce-or-drop-oldest 會先執行 coalesce,如仍然超出上限,便捨棄最舊的事件。EventCallbackeventcallback 的直接連結
訂閱者的 callback signature:(event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void。