メインコンテンツへ移動

RedisStreamsPubSub

RedisStreamsPubSubPubSub の実装で、Redis Streams を基盤とします。永続化、コンシューマーグループ、障害時の再配信に対応し、プロセスやホストをまたいでイベントを配信します。また、LeaseProvider も実装しているため、signals レイヤーはインスタンス間でリソースごとに単一の所有者を選出できます。これにより、分散環境やサーバーレス環境で signals が実行を調整できます。

複数のサービスがイベントストリームを共有する分散デプロイで使用します。単一プロセスでの配信には EventEmitterPubSub を使用してください。Google Cloud では GoogleCloudPubSub を使用してください。

各トピックは Redis のストリームキーに対応します。グループを指定したサブスクリプションでは Redis のコンシューマーグループを使用するため、メンバー間でラウンドロビン方式により処理を分担します。グループを指定しないサブスクリプションではプライベートなコンシューマーグループが作成されるため、すべてのサブスクライバーがすべてのイベントを受信します。

RedisStreamsPubSub は pull トランスポートです。コンシューマーは 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
ストリームキーのプレフィックス。各トピックは <keyPrefix>:<topic> に対応します。

blockMs?:

number
= 1000
新しいイベントを待つ間、各読み取りをブロックする時間(ミリ秒)。

redisOptions?:

RedisClientOptions
高度な設定のために、基盤となる redis クライアントへ渡すオプション。

maxStreamLength?:

number
= 10000
ストリームごとに保持するエントリのおおよその最大数。トリミングを無効にするには 0 を指定します。

streamIdleTtlMs?:

number
= 0
アイドル状態で期限切れになるまでの時間(ミリ秒)。書き込み(publish、nack による再試行、グループの再作成)のたびに更新されるスライディング TTL です。書き込みのたびにリセットされるため、活発に書き込まれているストリームが処理中に期限切れになることはありません。所定の期間にわたりアイドル状態が続いたストリームは Redis によって自動的に削除されます。TTL が更新されるのは書き込み時だけであり、コンシューマーがバックログをゆっくり処理しても更新されない点に注意してください。そのため、稼働中のトピックで想定される書き込み間隔の最大値を十分に上回る値に設定してください。これは補助的な安全策であり、主要なクリーンアップ手段ではありません。通常のライフサイクル終了時の削除は clearTopic が処理します。この設定は、clearTopic の呼び出しに到達しないストリーム(クラッシュした実行など)のメモリ使用量だけを制限します。0 以上の整数である必要があります。デフォルトは 0(無効)です。

reclaimIntervalMs?:

number
= 30000
以前のコンシューマーが読み取ったものの確認応答しなかったイベントを、サブスクリプションが再取得する間隔(ミリ秒)。無効にするには 0 を指定します。

reclaimIdleMs?:

number
= 60000
保留中のイベントが再取得の対象になるまでに必要な最小アイドル時間(ミリ秒)。二重配信を避けるため、通常の処理時間を十分に上回る値にしてください。

maxDeliveryAttempts?:

number
= 5
イベントが破棄されるまでに nack によって再配信される最大回数。上限を無効にするには Infinity を渡します。

logger?:

{ debug?: Function; warn?: Function }
診断用の任意の logger。省略した場合、抑制されたエラーは出力されません。

プロパティ
プロパティへの直接リンク

supportedModes:

ReadonlyArray<"pull" | "push">
["pull"] を返します。

メソッド
メソッドへの直接リンク

RedisStreamsPubSubPubSub 契約を実装します。以下のメソッドには、この実装固有の動作があります。

subscribe(topic, cb, options?)
subscribetopic-cb-optionsへの直接リンク

トピックを購読します。options.group を指定すると、グループのメンバーは Redis のコンシューマーグループを通じてイベントを分担します。グループを指定しない場合、サブスクライバーはプライベートなコンシューマーグループを通じてすべてのイベントを受信します。

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

flush()
flushへの直接リンク

処理中の publish が完了するまで待機します。

await pubsub.flush()

clearTopic(topic)
cleartopictopicへの直接リンク

トピックのストリームと、そのストリーム上のすべてのコンシューマーグループを削除し、完了したトピックが保持し続けるメモリを解放します。Mastra の実行ライフサイクル(durable Agent とイベント駆動型 Workflow エンジン)は、実行が終了状態に達するとこれを自動的に呼び出します。トピックが今後一切読み取られないことを確認してから、自分で呼び出してください。これはベストエフォートで実行され、例外をスローしません。失敗は warn レベルでログに記録されます。ストリームの削除時に接続中のサブスクライバーは自動的に復旧しますが、削除されたエントリは受信できません。

自動クリーンアップには、clearTopic をサポートするバージョンの @mastra/core@mastra/redis-streams が両方とも必要です。ランタイムはキャッシュレイヤーを介して呼び出しをルーティングするため、実行終了時のストリーム削除を有効にするには、2 つのパッケージを同時にアップグレードしてください。

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

close()
closeへの直接リンク

Redis 接続を閉じ、すべてのサブスクリプションを停止します。正常なシャットダウン時に呼び出してください。

await pubsub.close()

再配信と再取得
再配信と再取得への直接リンク

サブスクライバーが nack を呼び出すと、deliveryAttempt をインクリメントしてイベントが再度 publish され、元のイベントには確認応答が行われます。イベントが maxDeliveryAttempts に達すると、再配信されずに破棄されます。これとは別に、各サブスクリプションは、グループ内の以前のコンシューマーが読み取ったものの確認応答しなかったイベントを定期的に再取得します。この動作は reclaimIntervalMsreclaimIdleMs で制御されます。

分散リース
分散リースへの直接リンク

RedisStreamsPubSub は、同じ Redis 接続上に LeaseProvider 契約を実装します。signals ランタイムはこれを使用して単一の所有者(通常はスレッドキーごと)を選出します。これにより、インスタンスをまたいで 1 つのプロセスだけが Agent を起動して実行し、ほかのプロセスは後続の処理を所有者へルーティングします。この仕組みによって、signals はサーバーレス環境やマルチインスタンス環境で動作します。共有リースがなければ、各インスタンスがそれぞれ競合する実行を開始します。

リースキーは、トピックと同じ keyPrefix の名前空間に <keyPrefix>:lease:<key> として配置されます。すべての操作はアトミックです。acquireLeaseSET NX PX を使用し、自身の TTL を冪等に更新します。releaseLeaserenewLeasetransferLease は、変更前に所有権を確認する Lua スクリプトを使用するため、別の所有者による同時更新が上書きされることはありません。

これらのメソッドを直接呼び出す必要はありません。RedisStreamsPubSubpubsub バックエンドとして設定するだけで、ランタイムがこの機能を検出して使用します。メソッド契約の詳細については、LeaseProvider を参照してください。