> Discover all available pages from the documentation index: https://mastra.zisheng.pro/zh-HK/llms.txt # PubSub `PubSub` 是 Mastra 事件系統的抽象基礎類別。它定義每個發佈/訂閱後端都會實作的協定,讓 Mastra 的其他部分毋須知道正在使用哪種傳輸方式,亦能發佈及訂閱事件。 Mastra 在內部使用發佈/訂閱機制來處理 Workflow 事件、串流及跨元件通訊。大多數應用程式會使用預設的 [`EventEmitterPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/event-emitter),而不會直接建立 `PubSub`。只有在需要自訂傳輸方式時,才應實作此類別。 如要了解內置實作,請參閱 [`EventEmitterPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/event-emitter)、[`UnixSocketPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/unix-socket-pubsub)、[`CachingPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/caching-pubsub)、[`RedisStreamsPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/redis-streams) 及 [`GoogleCloudPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/google-cloud-pubsub)。 ## 使用範例 擴充 `PubSub` 並實作四個抽象方法,以加入自訂後端。 ```typescript 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): Promise { // Deliver the event to subscribers of `topic`. } async subscribe(topic: string, cb: EventCallback, options?: SubscribeOptions): Promise { // Register `cb` to receive events published to `topic`. } async unsubscribe(topic: string, cb: EventCallback): Promise { // Remove a previously registered callback. } async flush(): Promise { // Wait for any in-flight deliveries to settle. } } ``` 將實例傳入 [Mastra](https://mastra.zisheng.pro/zh-HK/reference/core/mastra-class) 建構函式: ```typescript 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)` 將事件發佈至主題。`id` 及 `createdAt` 欄位由實作指派。 ```typescript await pubsub.publish('my-topic', { type: 'example', data: { value: 1 }, runId: 'run-123', }) ``` #### `subscribe(topic, cb, options?)` 註冊 callback,以接收發佈至主題的事件。設定 `options.group` 後,同一群組的訂閱者會競逐訊息,而每個事件只會傳送給其中一名成員。如沒有群組,每名訂閱者都會收到每個事件。 傳入 `options.batch` 以選用批次傳送。callback signature 維持不變:包含 N 個事件的批次會按發佈次序,透過 N 次連續的 `cb(event, ack, nack)` 呼叫傳送。只有後端的 [`supportsNativeBatching`](#properties) 為 `true` 時,才會採用批次處理。其他後端會忽略此選項,並逐一傳送事件。 ```typescript await pubsub.subscribe('my-topic', (event, ack, nack) => { console.log(event) }) ``` #### `unsubscribe(topic, cb)` 從主題移除先前註冊的 callback。 ```typescript await pubsub.unsubscribe('my-topic', callback) ``` #### `flush()` 等待所有進行中的傳送完成。關閉前請呼叫此方法,以免遺失事件。 ```typescript await pubsub.flush() ``` #### `clearTopic(topic)` 當不會再有事件發佈至某個主題後,刪除該主題保留的所有狀態(快取記錄、持久化串流項目及使用者群組)。當一次執行到達終止狀態時,Mastra 的執行生命週期(durable Agent 及事件驅動的 Workflow 引擎)會自動呼叫此方法,因此在會保留訊息的傳輸方式上,個別執行的主題不會持續累積。 預設實作不會執行任何操作:不會為個別主題保留資料的傳輸方式(例如 `EventEmitterPubSub`)沒有需要清除的內容。會持久保存訊息的後端(例如 [`RedisStreamsPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/redis-streams))則會覆寫此方法。此協定採取盡力而為的方式:實作會記錄失敗,而不會拋出錯誤,因為呼叫者會在清理邊界以不等待結果的方式呼叫此方法。 ```typescript await pubsub.clearTopic('workflow.events.v2.run-123') ``` ### 重播方法 這些方法支援在中斷連線後恢復串流。預設實作會退回一般的 `subscribe`,因此不支援記錄的後端只會處理即時事件。[`CachingPubSub`](https://mastra.zisheng.pro/zh-HK/reference/pubsub/caching-pubsub) 會覆寫這些方法,以重播已快取的事件。 #### `getHistory(topic, offset?)` 傳回主題從 `offset` 開始的已快取事件。如後端沒有記錄,則傳回空陣列。 ```typescript const events = await pubsub.getHistory('my-topic', 0) ``` 傳回:`Promise` #### `subscribeWithReplay(topic, cb)` 重播已快取的事件,然後訂閱即時事件。 ```typescript await pubsub.subscribeWithReplay('my-topic', event => { console.log(event) }) ``` #### `subscribeFromOffset(topic, offset, cb)` 從已知位置開始重播已快取的事件,然後訂閱即時事件。當 client 知道其最後位置時,這比完整重播更有效率。 ```typescript await pubsub.subscribeFromOffset('my-topic', 42, event => { console.log(event) }) ``` ## 屬性 **supportedModes** (`ReadonlyArray<"pull" | "push">`): 實作支援的傳送模式。預設為 \["pull"]。 **supportsNativeBatching** (`boolean`): 實作是否採用 options.batch 進行 subscribe()。預設為 false。內部整合批次處理的後端會覆寫此屬性並傳回 true。 ## 類型 ### `Event` **type** (`string`): 事件類型識別碼。 **id** (`string`): 唯一事件 ID,由實作在發佈時指派。 **data** (`any`): 事件 payload。 **runId** (`string`): 事件所屬的執行。 **createdAt** (`Date`): 由實作在發佈時指派的時間戳記。 **index** (`number`): 用於從特定 offset 恢復的順序位置。 **deliveryAttempt** (`number`): 事件已傳送的次數,由 1 開始。後端不追蹤重新傳送時,預設為 1。 ### `SubscribeOptions` **group** (`string`): 設定後,同一群組的訂閱者會競逐訊息,而每個事件只會傳送給其中一名成員。如省略此項,每名訂閱者都會收到每個事件。 **batch** (`SubscribeBatchOptions`): 為此訂閱選用批次傳送。如省略此項,事件會逐一傳送。只有 supportsNativeBatching 為 true 的後端才會採用此設定。 ### `SubscribeBatchOptions` 每個訂閱的批次處理政策。callback signature 不會改變。包含 N 個事件的批次會按發佈次序轉化為 N 次連續的 callback 呼叫。 **maxSize** (`number`): 強制 flush 前最多保留的事件數目。 **maxWaitMs** (`number`): 最舊事件可在 buffer 中停留的最長時間(毫秒)。buffer 從空白變為非空白時,計時器便會啟動。 **minIntervalMs** (`number`): 連續兩次批次傳送之間的最短時間(毫秒)。即使 maxSize 或 maxWaitMs 會觸發傳送,buffer 亦會保留內容,直至距離上次傳送已經過此段時間。 **isImmediate** (`(event: Event) => boolean`): 如為事件傳回 true,buffer 會在發佈時立即 flush,但仍受 minIntervalMs 限制。這是按事件使用的例外機制。 **coalesce** (`(events: Event[]) => Event[]`): 在傳送前套用至排隊中的批次,以捨棄已被取代的事件。必須按 reference identity 傳回輸入內容的子集;傳回新建的 Event 物件會違反協定,並令整個批次被捨棄。保留事件的次序不變。 **maxBufferSize** (`number`): 觸發 overflow 處理前,buffer 可保留的事件數目上限。標記為立即處理的事件絕不會因 overflow 而被捨棄。 (Default: `256`) **overflow** (`"drop-oldest" | "drop-newest" | "coalesce-or-drop-oldest"`): buffer 超出 maxBufferSize 時的 overflow 策略。coalesce-or-drop-oldest 會先執行 coalesce,如仍然超出上限,便捨棄最舊的事件。 (Default: `coalesce-or-drop-oldest`) ### `EventCallback` 訂閱者的 callback signature:`(event: Event, ack?: () => Promise, nack?: () => Promise) => void`。 **event** (`Event`): 已傳送的事件。 **ack** (`() => Promise`): 確認處理成功。事件會從佇列移除。 **nack** (`() => Promise`): 否定確認。事件會重新加入佇列,並在延遲後再次傳送。