> Discover all available pages from the documentation index: https://mastra.zisheng.pro/ja/llms.txt # RedisStreamsPubSub `RedisStreamsPubSub` は [`PubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/base) の実装で、[Redis Streams](https://redis.io/docs/latest/develop/data-types/streams/) を基盤とします。永続化、コンシューマーグループ、障害時の再配信に対応し、プロセスやホストをまたいでイベントを配信します。また、[`LeaseProvider`](https://mastra.zisheng.pro/ja/reference/pubsub/lease-provider) も実装しているため、signals レイヤーはインスタンス間でリソースごとに単一の所有者を選出できます。これにより、分散環境やサーバーレス環境で signals が実行を調整できます。 複数のサービスがイベントストリームを共有する分散デプロイで使用します。単一プロセスでの配信には [`EventEmitterPubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/event-emitter) を使用してください。Google Cloud では [`GoogleCloudPubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/google-cloud-pubsub) を使用してください。 各トピックは Redis のストリームキーに対応します。グループを指定したサブスクリプションでは Redis のコンシューマーグループを使用するため、メンバー間でラウンドロビン方式により処理を分担します。グループを指定しないサブスクリプションではプライベートなコンシューマーグループが作成されるため、すべてのサブスクライバーがすべてのイベントを受信します。 `RedisStreamsPubSub` は pull トランスポートです。コンシューマーは `XREADGROUP` でイベントを読み取るため、Mastra が代わりに読み取るオーケストレーション Worker を実行します。 ## インストール **npm**: ```bash npm install @mastra/redis-streams ``` **pnpm**: ```bash pnpm add @mastra/redis-streams ``` **Yarn**: ```bash yarn add @mastra/redis-streams ``` **Bun**: ```bash bun add @mastra/redis-streams ``` ## 使用例 Redis の接続 URL を指定します。 ```typescript 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 の接続 URL。指定しない場合は redisOptions.url を使用します。 (Default: `redis://localhost:6379`) **keyPrefix** (`string`): ストリームキーのプレフィックス。各トピックは \:\ に対応します。 (Default: `mastra:topic`) **blockMs** (`number`): 新しいイベントを待つ間、各読み取りをブロックする時間(ミリ秒)。 (Default: `1000`) **redisOptions** (`RedisClientOptions`): 高度な設定のために、基盤となる redis クライアントへ渡すオプション。 **maxStreamLength** (`number`): ストリームごとに保持するエントリのおおよその最大数。トリミングを無効にするには 0 を指定します。 (Default: `10000`) **streamIdleTtlMs** (`number`): アイドル状態で期限切れになるまでの時間(ミリ秒)。書き込み(publish、nack による再試行、グループの再作成)のたびに更新されるスライディング TTL です。書き込みのたびにリセットされるため、活発に書き込まれているストリームが処理中に期限切れになることはありません。所定の期間にわたりアイドル状態が続いたストリームは Redis によって自動的に削除されます。TTL が更新されるのは書き込み時だけであり、コンシューマーがバックログをゆっくり処理しても更新されない点に注意してください。そのため、稼働中のトピックで想定される書き込み間隔の最大値を十分に上回る値に設定してください。これは補助的な安全策であり、主要なクリーンアップ手段ではありません。通常のライフサイクル終了時の削除は clearTopic が処理します。この設定は、clearTopic の呼び出しに到達しないストリーム(クラッシュした実行など)のメモリ使用量だけを制限します。0 以上の整数である必要があります。デフォルトは 0(無効)です。 (Default: `0`) **reclaimIntervalMs** (`number`): 以前のコンシューマーが読み取ったものの確認応答しなかったイベントを、サブスクリプションが再取得する間隔(ミリ秒)。無効にするには 0 を指定します。 (Default: `30000`) **reclaimIdleMs** (`number`): 保留中のイベントが再取得の対象になるまでに必要な最小アイドル時間(ミリ秒)。二重配信を避けるため、通常の処理時間を十分に上回る値にしてください。 (Default: `60000`) **maxDeliveryAttempts** (`number`): イベントが破棄されるまでに nack によって再配信される最大回数。上限を無効にするには Infinity を渡します。 (Default: `5`) **logger** (`{ debug?: Function; warn?: Function }`): 診断用の任意の logger。省略した場合、抑制されたエラーは出力されません。 ## プロパティ **supportedModes** (`ReadonlyArray<"pull" | "push">`): \["pull"] を返します。 ## メソッド `RedisStreamsPubSub` は [`PubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/base) 契約を実装します。以下のメソッドには、この実装固有の動作があります。 ### `subscribe(topic, cb, options?)` トピックを購読します。`options.group` を指定すると、グループのメンバーは Redis のコンシューマーグループを通じてイベントを分担します。グループを指定しない場合、サブスクライバーはプライベートなコンシューマーグループを通じてすべてのイベントを受信します。 ```typescript await pubsub.subscribe('workflow.events', (event, ack, nack) => { console.log(event) }) ``` ### `flush()` 処理中の publish が完了するまで待機します。 ```typescript await pubsub.flush() ``` ### `clearTopic(topic)` トピックのストリームと、そのストリーム上のすべてのコンシューマーグループを削除し、完了したトピックが保持し続けるメモリを解放します。Mastra の実行ライフサイクル(durable Agent とイベント駆動型 Workflow エンジン)は、実行が終了状態に達するとこれを自動的に呼び出します。トピックが今後一切読み取られないことを確認してから、自分で呼び出してください。これはベストエフォートで実行され、例外をスローしません。失敗は warn レベルでログに記録されます。ストリームの削除時に接続中のサブスクライバーは自動的に復旧しますが、削除されたエントリは受信できません。 自動クリーンアップには、`clearTopic` をサポートするバージョンの `@mastra/core` と `@mastra/redis-streams` が両方とも必要です。ランタイムはキャッシュレイヤーを介して呼び出しをルーティングするため、実行終了時のストリーム削除を有効にするには、2 つのパッケージを同時にアップグレードしてください。 ```typescript await pubsub.clearTopic('workflow.events.run-123') ``` ### `close()` Redis 接続を閉じ、すべてのサブスクリプションを停止します。正常なシャットダウン時に呼び出してください。 ```typescript await pubsub.close() ``` ## 再配信と再取得 サブスクライバーが `nack` を呼び出すと、`deliveryAttempt` をインクリメントしてイベントが再度 publish され、元のイベントには確認応答が行われます。イベントが `maxDeliveryAttempts` に達すると、再配信されずに破棄されます。これとは別に、各サブスクリプションは、グループ内の以前のコンシューマーが読み取ったものの確認応答しなかったイベントを定期的に再取得します。この動作は `reclaimIntervalMs` と `reclaimIdleMs` で制御されます。 ## 分散リース `RedisStreamsPubSub` は、同じ Redis 接続上に [`LeaseProvider`](https://mastra.zisheng.pro/ja/reference/pubsub/lease-provider) 契約を実装します。[signals ランタイム](https://mastra.zisheng.pro/ja/docs/long-running-agents/signals)はこれを使用して単一の所有者(通常はスレッドキーごと)を選出します。これにより、インスタンスをまたいで 1 つのプロセスだけが Agent を起動して実行し、ほかのプロセスは後続の処理を所有者へルーティングします。この仕組みによって、signals はサーバーレス環境やマルチインスタンス環境で動作します。共有リースがなければ、各インスタンスがそれぞれ競合する実行を開始します。 リースキーは、トピックと同じ `keyPrefix` の名前空間に `:lease:` として配置されます。すべての操作はアトミックです。`acquireLease` は `SET NX PX` を使用し、自身の TTL を冪等に更新します。`releaseLease`、`renewLease`、`transferLease` は、変更前に所有権を確認する Lua スクリプトを使用するため、別の所有者による同時更新が上書きされることはありません。 これらのメソッドを直接呼び出す必要はありません。`RedisStreamsPubSub` を `pubsub` バックエンドとして設定するだけで、ランタイムがこの機能を検出して使用します。メソッド契約の詳細については、[`LeaseProvider`](https://mastra.zisheng.pro/ja/reference/pubsub/lease-provider) を参照してください。