跳至主要內容

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 傳送訊號時,系統會在 thread 主題上發佈事件,讓在另一進程中執行的項目收到訊號。
  • 可恢復串流:串流區塊會按每次執行發佈,讓重新連線的用戶端重播錯過的內容。

由於這些系統使用同一個匯流排,你選擇的後端會同時套用至所有系統。

傳送模式
傳送模式 的直接連結

後端以兩種模式之一傳送事件,模式由 PubSub 合約定義:

  • 拉取(Pull):接收端自行從後端讀取資料,Mastra 會透過長時間運行的 worker 迴圈執行此操作。RedisStreamsPubSub 等分佈式後端使用此模式。
  • 推送(Push):事件無須接收端提出請求便會送達,可在進程內或透過 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({})

由於它在進程內運行,事件不會持久保存,也無法傳送至其他進程。預設後端適合單一執行個體,足以應付大多數應用程式。

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

設定 pubsub 選項於 Mastra 執行個體上,以選擇後端。每個後端都實作相同的 PubSub 合約,因此應用程式其他部分無須更改。

決定因素是訂閱者在哪裏運行。

後端範圍模式依賴套件
EventEmitterPubSub單一進程拉取及推送@mastra/core
UnixSocketPubSub同一主機上的多個進程推送@mastra/core
RedisStreamsPubSub分佈式、多部主機拉取@mastra/redis-streams
GoogleCloudPubSub分佈式、多部主機拉取@mastra/google-cloud-pubsub

同一主機上的多個進程
同一主機上的多個進程 的直接連結

當同一部機器上的多個進程需要共用串流時,請使用 UnixSocketPubSub。它透過 Unix domain socket 傳送事件,並選出一個進程作為 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 參考資料,了解完整的傳送合約及各個後端的配置選項。