RedisStreamsPubSub
RedisStreamsPubSub 是以 Redis Streams 為基礎的 PubSub 實作。它可在多個 process 與主機間傳遞 event,並支援持久化、consumer group 與失敗時重新傳遞。它也實作 LeaseProvider,因此 signal layer 可跨 instance 為每項資源選出單一擁有者,讓 signal 能在分散式與 serverless 部署中協調 run。
適用於多個服務共用 event stream 的分散式部署。單一 process 傳遞請使用 EventEmitterPubSub。Google Cloud 請使用 GoogleCloudPubSub。
每個 topic 對應至 Redis stream key。具有 group 的 subscription 會使用 Redis consumer group,因此成員以 round-robin 方式分擔工作。沒有 group 的 subscription 會建立私有 consumer group,因此每個 subscriber 都會收到每個 event。
RedisStreamsPubSub 是 pull transport:consumer 使用 XREADGROUP 讀取 event,因此 Mastra 會代為執行 orchestration 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',
}),
})
Constructor 參數「Constructor 參數」的直接連結
url?:
redisOptions.url。keyPrefix?:
<keyPrefix>:<topic>。blockMs?:
redisOptions?:
redis client 的選項。maxStreamLength?:
streamIdleTtlMs?:
clearTopic 會處理一般生命週期結束時的刪除,此設定只限制從未呼叫 clearTopic 的 stream(例如當機的 run)所占用的記憶體。必須是非負整數。預設為 0(停用)。reclaimIntervalMs?:
reclaimIdleMs?:
maxDeliveryAttempts?:
nack 重新傳遞後,遭到捨棄前的最大次數。傳入 Infinity 可停用上限。logger?:
屬性「屬性」的直接連結
supportedModes:
["pull"]。方法「方法」的直接連結
RedisStreamsPubSub 實作 PubSub contract。以下方法具有此實作專屬的行為。
subscribe(topic, cb, options?)「subscribetopic-cb-options」的直接連結
訂閱 topic。設定 options.group 後,group 成員會透過 Redis consumer group 共用 event。未設定 group 時,subscriber 會透過私有 consumer group 收到每個 event。
await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})
flush()「flush」的直接連結
等待進行中的 publish 完成。
await pubsub.flush()
clearTopic(topic)「cleartopictopic」的直接連結
刪除 topic 的 stream 及其上的每個 consumer group,釋放已完成 topic 原本會持續占用的記憶體。Mastra 的 run 生命週期(durable Agent 與 event 型 Workflow engine)會在 run 進入終止 state 時自動呼叫。只有確認不會再有人讀取 topic 時,才自行呼叫。此方法採盡力而為,絕不會擲回例外;失敗會以 warn 層級記錄。刪除 stream 時仍連接的 subscriber 會自行恢復,但會錯過已刪除的 entry。
自動清理需要 @mastra/core 與 @mastra/redis-streams 版本都支援 clearTopic:runtime 會透過其 caching layer 路由呼叫,因此請同時升級兩個 package,才能在 run 結束時刪除 stream。
await pubsub.clearTopic('workflow.events.run-123')
close()「close」的直接連結
關閉 Redis 連線並停止所有 subscription。請在正常關閉期間呼叫此方法。
await pubsub.close()
重新傳遞與重新取得「重新傳遞與重新取得」的直接連結
subscriber 呼叫 nack 時,event 會以遞增的 deliveryAttempt 重新發布,原始 event 則會獲得確認。event 達到 maxDeliveryAttempts 後會遭到捨棄,不再重新傳遞。另外,每個 subscription 會定期重新取得同一 group 中先前 consumer 已讀取但未確認的 event,此行為由 reclaimIntervalMs 與 reclaimIdleMs 控制。
分散式租約「分散式租約」的直接連結
RedisStreamsPubSub 會在相同 Redis 連線之上實作 LeaseProvider contract。signal runtime 會使用它選出單一擁有者(通常以 thread key 為單位),讓多個 instance 中只會有一個 process 喚醒並執行 Agent,其他 process 則將後續工作路由給持有者。這讓 signal 可在 serverless 與多 instance 部署中運作;如果沒有共用租約,每個 instance 都會啟動自己的競爭 run。
租約 key 與 topic 使用相同的 keyPrefix namespace,格式為 <keyPrefix>:lease:<key>。所有操作皆為原子操作:acquireLease 使用 SET NX PX,並以冪等方式重新整理自己的 TTL;releaseLease、renewLease 與 transferLease 則使用 Lua script,在變更前檢查擁有權,因此絕不會覆蓋其他 owner 的並行續約。
你不需直接呼叫這些方法。將 RedisStreamsPubSub 設定為 pubsub 後端,runtime 就能偵測並使用此能力。完整方法 contract 請參閱 LeaseProvider。