> Discover all available pages from the documentation index: https://mastra.zisheng.pro/ja/llms.txt # EventEmitterPubSub `EventEmitterPubSub` はデフォルトの [`PubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/base) 実装です。Node.js の [`EventEmitter`](https://nodejs.org/api/events.html#class-eventemitter) を使用してインプロセスでイベントを配信するため、外部サービスなしで動作します。 単一プロセスのアプリケーションで使用します。1台のホスト上でプロセスをまたいで配信する場合は [`UnixSocketPubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/unix-socket-pubsub) を参照してください。分散配信には [`RedisStreamsPubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/redis-streams) または [`GoogleCloudPubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/google-cloud-pubsub) を使用します。 インプロセスで動作するため、イベントは永続化されず、他のプロセスとも共有されません。再開可能なストリームにリプレイが必要な場合は、[`CachingPubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/caching-pubsub) でラップしてください。 ## 使用例 `EventEmitterPubSub` は、`pubsub` オプションを設定しなければ自動的に使用されるため、ほとんどのアプリケーションでは直接インスタンス化しません。設定または共有したい場合にのみ、明示的に作成します。 ```typescript import { Mastra } from '@mastra/core' import { EventEmitterPubSub } from '@mastra/core/events' export const mastra = new Mastra({ pubsub: new EventEmitterPubSub(), }) ``` アプリケーションの他の部分と emitter を共有するには、既存の `EventEmitter` を渡します。 ```typescript import EventEmitter from 'node:events' import { EventEmitterPubSub } from '@mastra/core/events' const emitter = new EventEmitter() const pubsub = new EventEmitterPubSub(emitter) ``` バッチ配信のエラーを通知するには、logger を渡します。 ```typescript import { EventEmitterPubSub } from '@mastra/core/events' const pubsub = new EventEmitterPubSub(undefined, { logger }) ``` ## コンストラクターのパラメーター **existingEmitter** (`EventEmitter`): 配信に使用する既存の Node.js EventEmitter。省略すると、新しい EventEmitter が作成されます。 **options** (`EventEmitterPubSubOptions`): 任意の設定。 ## プロパティ **supportedModes** (`ReadonlyArray<"pull" | "push">`): \["pull", "push"] を返します。emitter は pull 形式の worker を提供することも、listener にイベントを直接 push することもできます。 **supportsNativeBatching** (`boolean`): true を返します。subscriber は options.batch を使用してバッチ配信を選択できます。 ## メソッド `EventEmitterPubSub` は [`PubSub`](https://mastra.zisheng.pro/ja/reference/pubsub/base) の規約を実装します。以下のメソッドには、この実装に固有の動作があります。 ### `subscribe(topic, cb, options?)` トピックの callback を登録します。`options.group` を指定しない場合、すべての subscriber がすべてのイベントを受信します。グループを指定すると、そのグループのメンバーにイベントがラウンドロビン方式で分配されます。 バッチ配信を選択するには `options.batch` を渡します。以下の[バッチ処理](#batching)を参照してください。 ```typescript await pubsub.subscribe('workflow.events', (event, ack, nack) => { console.log(event) }) ``` ### `flush()` `nack` による保留中の再配信が実行されるまで待機してから完了します。 ```typescript await pubsub.flush() ``` ### `close()` すべての listener を削除し、保留中の再配信をキャンセルします。正常終了時にこのメソッドを呼び出してください。 ```typescript await pubsub.close() ``` ## 再配信 グループ化された subscriber が `nack` を呼び出すと、少し遅れてイベントがグループに再配信され、`deliveryAttempt` の回数が増加します。`ack` を呼び出すと、そのイベントの追跡情報がクリアされます。fan-out subscriber では各イベントがすべての subscriber に一度ずつ届くため、何もしない `ack` 関数と `nack` 関数を受け取ります。 ## バッチ処理 `EventEmitterPubSub` は `options.batch` をネイティブに処理します。subscriber がバッチ配信を選択すると、イベントは subscriber ごとのインメモリバッファに保持され、flush 条件を満たした時点で連続する callback 呼び出しとして配信されます。fan-out subscriber とグループ subscriber のどちらもバッチ処理を利用できます。ポリシーの詳細は [`SubscribeBatchOptions`](https://mastra.zisheng.pro/ja/reference/pubsub/base) を参照してください。 ```typescript await pubsub.subscribe( 'workflow.events', event => { console.log(event) }, { batch: { maxSize: 10, // flush once 10 events have queued maxWaitMs: 500, // ...or after 500ms, whichever comes first }, }, ) ``` バッファはインメモリかつプロセス単位であるため、バッチの状態は永続化されず、再起動後には残りません。`maxWaitMs` によってトリガーされる flush はベストエフォートです。例外をスローする `coalesce` などのステップが失敗した場合、エラーはスローされず、設定された `logger` を通じて通知されます。`flush()` は、すべてのバッチ subscriber のバッファを空にしてから完了します。