跳至主要內容

PubSub

Mastra 使用發布/訂閱(pub/sub)系統作為內部事件匯流排。元件會將事件發布至主題,其他元件則訂閱這些主題以作出回應。你設定的後端會決定事件的傳遞範圍:僅限單一行程、同一主機上的多個行程,或不同執行個體之間。

你只需在 Mastra 執行個體上設定一次後端,系統其餘部分不必變更即可使用。Mastra 預設使用不需設定的行程內後端。

Mastra 如何使用 PubSub
「Mastra 如何使用 PubSub」的直接連結

數個內建系統會透過相同的 pub/sub 匯流排發布及訂閱事件:

  • Workflow 執行:排程器會發布 workflow.start 事件,常駐 worker 會取用該事件以執行 Workflow。在執行期間,步驟及生命週期事件會透過 pub/sub 傳遞。
  • 排程 Workflow:排程器會發布至 Workflow 主題,以分派到期的執行。請參閱排程 Workflow
  • 背景工作:工作管理員會將工作分派給 worker 群組,並將工作生命週期更新分送給訂閱者,背景工作串流因此能持續即時更新。請參閱背景工作串流
  • Agent 信號:向執行中的 Agent 執行作業傳送信號時,系統會在執行緒主題上發布事件,因此在另一個行程中執行的作業也能收到信號。
  • 可續傳串流:串流區塊會依每次執行發布,因此用戶端重新連線時可以重播錯過的內容。

由於這些系統都在同一條匯流排上執行,你選擇的後端會同時套用至所有系統。

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

後端會使用 PubSub 契約定義的兩種模式之一傳遞事件:

  • 提取:取用者自行從後端讀取,Mastra 會透過常駐 worker 迴圈執行此操作。RedisStreamsPubSub 等分散式後端會使用此模式。
  • 推送:事件無須取用者要求便會抵達,可能在行程內傳遞,也可能透過 HTTP 傳遞。預設的 EventEmitterPubSub 會在行程內以此方式傳遞。

訂閱者也可以透過取用者群組分配工作。同一群組的成員會分攤事件,讓每個事件只處理一次。沒有群組的訂閱者會收到每個事件,將串流分送給所有未分組的訂閱者。

預設後端
「預設後端」的直接連結

未設定 pubsub 選項時,Mastra 會使用 EventEmitterPubSub。它會透過 Node.js EventEmitter 在行程內傳遞事件,因此不需任何外部服務即可運作。

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

// No pubsub option: Mastra uses EventEmitterPubSub
export const mastra = new Mastra({})

由於它在行程內執行,事件不會保存,也不會傳至其他行程。預設後端適合單一執行個體,可涵蓋大多數應用程式。

選擇後端
「選擇後端」的直接連結

Mastra 執行個體上設定 pubsub 選項,即可選擇後端。每個後端都實作相同的 PubSub 契約,因此應用程式其餘部分不需變更。

關鍵在於訂閱者於何處執行。

後端範圍模式套件
EventEmitterPubSub單一行程提取與推送@mastra/core
UnixSocketPubSub同一主機上的多個行程推送@mastra/core
RedisStreamsPubSub分散式、多部主機提取@mastra/redis-streams
GoogleCloudPubSub分散式、多部主機提取@mastra/google-cloud-pubsub

同一主機上的多個行程
「同一主機上的多個行程」的直接連結

同一部機器上的多個行程需要共用串流時,請使用 UnixSocketPubSub。它會透過 Unix 網域通訊端傳遞事件,並選出一個行程作為 broker。若 broker 結束,其餘行程會選出新的 broker。

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { UnixSocketPubSub } from '@mastra/core/events'

export const mastra = new Mastra({
pubsub: new UnixSocketPubSub('/tmp/mastra/events.sock'),
})

分散式部署
「分散式部署」的直接連結

執行多個執行個體或主機時,請使用分散式後端,讓每個執行個體都能收到相同事件。只要由一個執行個體處理的請求必須觸及另一個執行個體上執行的工作,這項設定就很重要。

例如,向 Agent 執行作業傳送信號時,信號事件必須跨越行程界線,傳至擁有該執行作業的執行個體。若使用預設的行程內後端,該執行個體永遠不會收到事件。

下列兩種後端都能跨行程與主機傳遞事件,並保存事件以供重新傳遞。

以下範例使用 Redis Streams

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'

export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})

以下範例使用 Google Cloud Pub/Sub

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { GoogleCloudPubSub } from '@mastra/google-cloud-pubsub'

export const mastra = new Mastra({
pubsub: new GoogleCloudPubSub({
projectId: 'my-project',
}),
})

可續傳串流
「可續傳串流」的直接連結

可續傳串流讓用戶端能重新連線,並重播錯過的事件,因此後端必須保留近期歷程。RedisStreamsPubSub 等分散式後端會保存事件,所以本身就支援重播。

行程內傳遞不會保留歷程。如要在 EventEmitterPubSub 上加入重播功能,請以 CachingPubSub 包裝。它會依主題記錄已發布的事件,讓較晚加入或重新連線的訂閱者先補上進度,再繼續接收即時事件。

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { CachingPubSub, EventEmitterPubSub } from '@mastra/core/events'
import { InMemoryServerCache } from '@mastra/core/cache'

const cache = new InMemoryServerCache()

export const mastra = new Mastra({
pubsub: new CachingPubSub(new EventEmitterPubSub(), cache),
})

如需完整的傳遞契約及各後端的設定選項,請參閱 PubSub 參考