EventEmitterPubSub
EventEmitterPubSub はデフォルトの PubSub 実装です。Node.js の EventEmitter を使用してインプロセスでイベントを配信するため、外部サービスなしで動作します。
単一プロセスのアプリケーションで使用します。1台のホスト上でプロセスをまたいで配信する場合は UnixSocketPubSub を参照してください。分散配信には RedisStreamsPubSub または GoogleCloudPubSub を使用します。
インプロセスで動作するため、イベントは永続化されず、他のプロセスとも共有されません。再開可能なストリームにリプレイが必要な場合は、CachingPubSub でラップしてください。
使用例使用例への直接リンク
EventEmitterPubSub は、pubsub オプションを設定しなければ自動的に使用されるため、ほとんどのアプリケーションでは直接インスタンス化しません。設定または共有したい場合にのみ、明示的に作成します。
import { Mastra } from '@mastra/core'
import { EventEmitterPubSub } from '@mastra/core/events'
export const mastra = new Mastra({
pubsub: new EventEmitterPubSub(),
})
アプリケーションの他の部分と emitter を共有するには、既存の EventEmitter を渡します。
import EventEmitter from 'node:events'
import { EventEmitterPubSub } from '@mastra/core/events'
const emitter = new EventEmitter()
const pubsub = new EventEmitterPubSub(emitter)
バッチ配信のエラーを通知するには、logger を渡します。
import { EventEmitterPubSub } from '@mastra/core/events'
const pubsub = new EventEmitterPubSub(undefined, { logger })
コンストラクターのパラメーターコンストラクターのパラメーターへの直接リンク
existingEmitter?:
options?:
プロパティプロパティへの直接リンク
supportedModes:
["pull", "push"] を返します。emitter は pull 形式の worker を提供することも、listener にイベントを直接 push することもできます。supportsNativeBatching:
true を返します。subscriber は options.batch を使用してバッチ配信を選択できます。メソッドメソッドへの直接リンク
EventEmitterPubSub は PubSub の規約を実装します。以下のメソッドには、この実装に固有の動作があります。
subscribe(topic, cb, options?)subscribetopic-cb-optionsへの直接リンク
トピックの callback を登録します。options.group を指定しない場合、すべての subscriber がすべてのイベントを受信します。グループを指定すると、そのグループのメンバーにイベントがラウンドロビン方式で分配されます。
バッチ配信を選択するには options.batch を渡します。以下のバッチ処理を参照してください。
await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})
flush()flushへの直接リンク
nack による保留中の再配信が実行されるまで待機してから完了します。
await pubsub.flush()
close()closeへの直接リンク
すべての listener を削除し、保留中の再配信をキャンセルします。正常終了時にこのメソッドを呼び出してください。
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 を参照してください。
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 のバッファを空にしてから完了します。