跳至主要內容

RedisStreamsPubSub

RedisStreamsPubSub 是由 Redis Streams 支援的 PubSub 實作。它可跨程序及主機傳送事件,並支援持久化、使用者群組及失敗後重新傳送。它亦實作 LeaseProvider,因此 signals 層可以跨實例為每項資源選出一名擁有者,讓 signals 能夠在分散式及 serverless 部署中協調執行。

適用於多個服務共用事件串流的分散式部署。如要在單一程序內傳送,請使用 EventEmitterPubSub。如使用 Google Cloud,請使用 GoogleCloudPubSub

每個主題都會對應至 Redis 串流 key。設有群組的訂閱會使用 Redis 使用者群組,讓成員以 round-robin 方式分擔工作。沒有群組的訂閱則會建立私人使用者群組,讓每名訂閱者都收到每個事件。

RedisStreamsPubSub 是拉取式傳輸:使用者透過 XREADGROUP 讀取事件,因此 Mastra 會代為執行協調 worker 以讀取事件。

安裝
安裝 的直接連結

npm install @mastra/redis-streams

使用範例
使用範例 的直接連結

提供 Redis 連線 URL。

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',
}),
})

建構函式參數
建構函式參數 的直接連結

url?:

string
= redis://localhost:6379
Redis 連線 URL。如未設定,會使用 redisOptions.url

keyPrefix?:

string
= mastra:topic
串流 key 的前綴。每個主題都會對應至 <keyPrefix>:<topic>

blockMs?:

number
= 1000
每次讀取在等待新事件時阻塞的時間(毫秒)。

redisOptions?:

RedisClientOptions
傳入底層 redis client、用於進階設定的選項。

maxStreamLength?:

number
= 10000
每個串流大約保留的項目數目上限。設為 0 可停用修剪。

streamIdleTtlMs?:

number
= 0
閒置到期時間(毫秒):每次寫入(發佈、nack 重試、重新建立群組)時都會重新整理的滑動 TTL。每次寫入都會重設 TTL,因此持續有內容寫入的串流不會在處理期間到期;串流閒置整段指定時間後,Redis 會自動將其刪除。請注意,只有寫入操作會重新整理 TTL,緩慢清理 backlog 的使用者並不會重新整理;因此,此值應遠高於即時主題兩次寫入之間預期的最長間隔。這是最後保障,而非主要清理機制;clearTopic 會處理生命週期正常結束時的刪除,而此設定只會限制從未呼叫 clearTopic 的串流(例如執行崩潰)的記憶體用量。必須是非負整數。預設為 0(停用)。

reclaimIntervalMs?:

number
= 30000
訂閱收回先前使用者已讀取但從未確認之事件的頻率(毫秒)。設為 0 可停用。

reclaimIdleMs?:

number
= 60000
待處理事件符合收回資格前的最短閒置時間(毫秒)。此值應遠高於一般處理時間,以免重複傳送。

maxDeliveryAttempts?:

number
= 5
事件透過 nack 重新傳送並最終被捨棄前的最大次數。傳入 Infinity 可停用此上限。

logger?:

{ debug?: Function; warn?: Function }
用於診斷的可選 logger。如省略此項,被抑制的錯誤不會產生任何輸出。

屬性
屬性 的直接連結

supportedModes:

ReadonlyArray<"pull" | "push">
傳回 ["pull"]

方法
方法 的直接連結

RedisStreamsPubSub 會實作 PubSub 協定。以下方法具有此實作特有的行為。

subscribe(topic, cb, options?)
subscribetopic-cb-options 的直接連結

訂閱主題。設有 options.group 時,群組成員會透過 Redis 使用者群組共用事件。沒有群組時,訂閱者會透過私人使用者群組收到每個事件。

await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})

flush()
flush 的直接連結

等待進行中的發佈完成。

await pubsub.flush()

clearTopic(topic)
cleartopictopic 的直接連結

刪除主題的串流及其中所有使用者群組,釋放已完成主題原本會佔用的記憶體。當一次執行到達終止狀態時,Mastra 的執行生命週期(durable Agent 及事件驅動的 Workflow 引擎)會自動呼叫此方法。只有在不會再有任何內容讀取該主題後,才應自行呼叫此方法。此方法會盡力完成操作,絕不會拋出錯誤。失敗會記錄在 warn level。刪除串流時仍保持連接的訂閱者會自行恢復,但會錯過已刪除的項目。

自動清理需要 @mastra/core@mastra/redis-streams 的版本都支援 clearTopic:runtime 會透過快取層轉送呼叫,因此請同時升級兩個依賴套件,以便在執行結束時刪除串流。

await pubsub.clearTopic('workflow.events.run-123')

close()
close 的直接連結

關閉 Redis 連線並停止所有訂閱。請在正常關閉期間呼叫此方法。

await pubsub.close()

重新傳送及收回
重新傳送及收回 的直接連結

訂閱者呼叫 nack 時,事件會以遞增的 deliveryAttempt 重新發佈,而原始事件會被確認。事件到達 maxDeliveryAttempts 後,會被捨棄而不再傳送。此外,每個訂閱都會定期收回群組中先前使用者已讀取但從未確認的事件;此行為由 reclaimIntervalMsreclaimIdleMs 控制。

分散式租約
分散式租約 的直接連結

RedisStreamsPubSub 會在同一 Redis 連線上實作 LeaseProvider 協定。signals runtime 會使用它選出單一擁有者(通常按 thread key 區分),確保跨實例只有一個程序會喚醒及執行 Agent,其他程序則將後續工作轉送給租約持有人。這正是 signals 能夠在 serverless 及多實例部署中運作的原因;如沒有共用租約,每個實例都會開始各自互相競逐的執行。

租約 key 會使用與主題相同的 keyPrefix namespace,格式為 <keyPrefix>:lease:<key>。所有操作都是原子操作:acquireLease 使用 SET NX PX,並以 idempotent 方式重新整理自身的 TTL;releaseLeaserenewLeasetransferLease 則使用 Lua scripts,在修改前檢查擁有權,因此不會覆蓋另一名擁有者同時進行的續期操作。

你毋須直接呼叫這些方法。將 RedisStreamsPubSub 設定為 pubsub 後端,runtime 便會自動偵測並使用此功能。如要了解完整的方法協定,請參閱 LeaseProvider