> Discover all available pages from the documentation index: https://mastra.zisheng.pro/zh-TW/llms.txt # PubSub `PubSub` 是 Mastra event 系統的抽象基底類別。它定義每個 pub/sub 後端都必須實作的 contract,因此 Mastra 的其他部分不需要知道所使用的 transport,也能發布及訂閱 event。 Mastra 在內部使用 pub/sub 處理 Workflow event、串流與跨元件通訊。大多數應用程式使用預設的 [`EventEmitterPubSub`](https://mastra.zisheng.pro/zh-TW/reference/pubsub/event-emitter),不會直接建構 `PubSub`。只有需要自訂 transport 時才實作此類別。 內建實作請參閱 [`EventEmitterPubSub`](https://mastra.zisheng.pro/zh-TW/reference/pubsub/event-emitter)、[`UnixSocketPubSub`](https://mastra.zisheng.pro/zh-TW/reference/pubsub/unix-socket-pubsub)、[`CachingPubSub`](https://mastra.zisheng.pro/zh-TW/reference/pubsub/caching-pubsub)、[`RedisStreamsPubSub`](https://mastra.zisheng.pro/zh-TW/reference/pubsub/redis-streams)和 [`GoogleCloudPubSub`](https://mastra.zisheng.pro/zh-TW/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. } } ``` 將 instance 傳給 [Mastra](https://mastra.zisheng.pro/zh-TW/reference/core/mastra-class) constructor: ```typescript 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)` 將 event 發布至 topic。`id` 與 `createdAt` 欄位由實作指派。 ```typescript await pubsub.publish('my-topic', { type: 'example', data: { value: 1 }, runId: 'run-123', }) ``` #### `subscribe(topic, 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`](#properties) 為 `true` 時才會採用批次處理。其他後端會忽略此選項,逐一傳遞 event。 ```typescript await pubsub.subscribe('my-topic', (event, ack, nack) => { console.log(event) }) ``` #### `unsubscribe(topic, cb)` 從 topic 移除先前註冊的 callback。 ```typescript await pubsub.unsubscribe('my-topic', callback) ``` #### `flush()` 等待所有進行中的傳遞完成。請在關閉前呼叫,以免 event 遺失。 ```typescript await pubsub.flush() ``` #### `clearTopic(topic)` 當 topic 不再發布 event 時,刪除其所有保留 state(快取歷史記錄、持久化 stream entry 與 consumer group)。Mastra 的 run 生命週期(durable Agent 與 event 型 Workflow engine)會在 run 進入終止 state 時自動呼叫此方法,因此每個 run 的 topic 不會在保留訊息的 transport 上持續累積。 預設實作不會執行任何操作:不會按 topic 保留任何內容的 transport(例如 `EventEmitterPubSub`)沒有需要清除的資料。持久化訊息的後端(例如 [`RedisStreamsPubSub`](https://mastra.zisheng.pro/zh-TW/reference/pubsub/redis-streams))會覆寫此方法。此 contract 採盡力而為:實作會記錄失敗,而不會擲回例外,因為 caller 會在清理邊界以 fire-and-forget 方式呼叫。 ```typescript await pubsub.clearTopic('workflow.events.v2.run-123') ``` ### 重播方法 這些方法支援斷線後恢復串流。預設實作會改用一般 `subscribe`,因此不支援歷史記錄的後端只會提供即時內容。[`CachingPubSub`](https://mastra.zisheng.pro/zh-TW/reference/pubsub/caching-pubsub) 會覆寫這些方法,以重播快取 event。 #### `getHistory(topic, offset?)` 傳回 topic 從 `offset` 開始的快取 event。後端沒有歷史記錄時傳回空陣列。 ```typescript const events = await pubsub.getHistory('my-topic', 0) ``` 傳回:`Promise` #### `subscribeWithReplay(topic, cb)` 重播快取 event,再訂閱即時 event。 ```typescript await pubsub.subscribeWithReplay('my-topic', event => { console.log(event) }) ``` #### `subscribeFromOffset(topic, offset, cb)` 從已知位置開始重播快取 event,再訂閱即時 event。client 知道最後位置時,這比完整重播更有效率。 ```typescript await pubsub.subscribeFromOffset('my-topic', 42, event => { console.log(event) }) ``` ## 屬性 **supportedModes** (`ReadonlyArray<"pull" | "push">`): 實作支援的傳遞模式。預設為 \["pull"]。 **supportsNativeBatching** (`boolean`): 實作是否採用 subscribe() 上的 options.batch。預設為 false。內部整合批次處理的後端會覆寫此屬性並傳回 true。 ## 型別 ### `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` **group** (`string`): 設定後,同一 group 中的 subscriber 會競爭訊息,每個 event 只會傳遞給一個成員。省略時,每個 subscriber 都會收到每個 event。 **batch** (`SubscribeBatchOptions`): 為此 subscription 選用批次傳遞。省略時會逐一傳遞 event。只有 supportsNativeBatching 為 true 的後端才會採用。 ### `SubscribeBatchOptions` 各 subscription 專用的批次政策。callback signature 不會改變。一批 N 個 event 會依發布順序轉換為 N 次連續 callback 呼叫。 **maxSize** (`number`): 強制 flush 前最多保留的 event 數量。 **maxWaitMs** (`number`): 最舊 event 可停留在 buffer 中的最長毫秒數。buffer 從空白變為非空白時,timer 便會啟動。 **minIntervalMs** (`number`): 連續批次傳遞之間的最短毫秒數。即使 maxSize 或 maxWaitMs 已觸發,buffer 仍會保留至上次傳遞後已經過此間隔。 **isImmediate** (`(event: Event) => boolean`): 對 event 傳回 true 時,buffer 會在發布時立即 flush,但仍受 minIntervalMs 限制。這是各 event 專用的例外處理機制。 **coalesce** (`(events: Event[]) => Event[]`): 在傳遞前套用至佇列中的批次,以捨棄已被取代的 event。必須依參照識別傳回輸入的子集合;傳回新建構的 Event 物件會違反 contract,並捨棄整個批次。保留的 event 順序不變。 **maxBufferSize** (`number`): 觸發 overflow 處理前,buffer 最多可容納的 event 數量。標示為立即處理的 event 絕不會因 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` subscriber 的 callback signature:`(event: Event, ack?: () => Promise, nack?: () => Promise) => void`。 **event** (`Event`): 已傳遞的 event。 **ack** (`() => Promise`): 確認處理成功。event 會從 queue 中移除。 **nack** (`() => Promise`): 否定確認。event 會重新排入 queue,並在延遲後重新傳遞。