RedisStreamsPubSub
RedisStreamsPubSub 是由 Redis Streams 支援的 PubSub 實作。它可跨程序及主機傳送事件,並支援持久化、使用者群組及失敗後重新傳送。它亦實作 LeaseProvider,因此 signals 層可以跨實例為每項資源選出一名擁有者,讓 signals 能夠在分散式及 serverless 部署中協調執行。
適用於多個服務共用事件串流的分散式部署。如要在單一程序內傳送,請使用 EventEmitterPubSub。如使用 Google Cloud,請使用 GoogleCloudPubSub。
每個主題都會對應至 Redis 串流 key。設有群組的訂閱會使用 Redis 使用者群組,讓成員以 round-robin 方式分擔工作。沒有群組的訂閱則會建立私人使用者群組,讓每名訂閱者都收到每個事件。
RedisStreamsPubSub 是拉取式傳輸:使用者透過 XREADGROUP 讀取事件,因此 Mastra 會代為執行協調 worker 以讀取事件。
安裝安裝 的直接連結
- npm
- pnpm
- Yarn
- Bun
npm install @mastra/redis-streams
pnpm add @mastra/redis-streams
yarn add @mastra/redis-streams
bun add @mastra/redis-streams
使用範例使用範例 的直接連結
提供 Redis 連線 URL。
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'
export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})
建構函式參數建構函式參數 的直接連結
url?:
redisOptions.url。keyPrefix?:
<keyPrefix>:<topic>。blockMs?:
redisOptions?:
redis client、用於進階設定的選項。maxStreamLength?:
streamIdleTtlMs?:
clearTopic 會處理生命週期正常結束時的刪除,而此設定只會限制從未呼叫 clearTopic 的串流(例如執行崩潰)的記憶體用量。必須是非負整數。預設為 0(停用)。reclaimIntervalMs?:
reclaimIdleMs?:
maxDeliveryAttempts?:
nack 重新傳送並最終被捨棄前的最大次數。傳入 Infinity 可停用此上限。logger?:
屬性屬性 的直接連結
supportedModes:
["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 後,會被捨棄而不再傳送。此外,每個訂閱都會定期收回群組中先前使用者已讀取但從未確認的事件;此行為由 reclaimIntervalMs 及 reclaimIdleMs 控制。
分散式租約分散式租約 的直接連結
RedisStreamsPubSub 會在同一 Redis 連線上實作 LeaseProvider 協定。signals runtime 會使用它選出單一擁有者(通常按 thread key 區分),確保跨實例只有一個程序會喚醒及執行 Agent,其他程序則將後續工作轉送給租約持有人。這正是 signals 能夠在 serverless 及多實例部署中運作的原因;如沒有共用租約,每個實例都會開始各自互相競逐的執行。
租約 key 會使用與主題相同的 keyPrefix namespace,格式為 <keyPrefix>:lease:<key>。所有操作都是原子操作:acquireLease 使用 SET NX PX,並以 idempotent 方式重新整理自身的 TTL;releaseLease、renewLease 及 transferLease 則使用 Lua scripts,在修改前檢查擁有權,因此不會覆蓋另一名擁有者同時進行的續期操作。
你毋須直接呼叫這些方法。將 RedisStreamsPubSub 設定為 pubsub 後端,runtime 便會自動偵測並使用此功能。如要了解完整的方法協定,請參閱 LeaseProvider。