メインコンテンツへ移動

EventEmitterPubSub

EventEmitterPubSub はデフォルトの PubSub 実装です。Node.js の EventEmitter を使用してインプロセスでイベントを配信するため、外部サービスなしで動作します。

単一プロセスのアプリケーションで使用します。1台のホスト上でプロセスをまたいで配信する場合は UnixSocketPubSub を参照してください。分散配信には RedisStreamsPubSub または GoogleCloudPubSub を使用します。

インプロセスで動作するため、イベントは永続化されず、他のプロセスとも共有されません。再開可能なストリームにリプレイが必要な場合は、CachingPubSub でラップしてください。

使用例
使用例への直接リンク

EventEmitterPubSub は、pubsub オプションを設定しなければ自動的に使用されるため、ほとんどのアプリケーションでは直接インスタンス化しません。設定または共有したい場合にのみ、明示的に作成します。

src/mastra/index.ts
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?:

EventEmitter
配信に使用する既存の Node.js EventEmitter。省略すると、新しい EventEmitter が作成されます。

options?:

EventEmitterPubSubOptions
任意の設定。
IMastraLogger

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

supportedModes:

ReadonlyArray<"pull" | "push">
["pull", "push"] を返します。emitter は pull 形式の worker を提供することも、listener にイベントを直接 push することもできます。

supportsNativeBatching:

boolean
true を返します。subscriber は options.batch を使用してバッチ配信を選択できます。

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

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

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 関数を受け取ります。

バッチ処理
バッチ処理への直接リンク

EventEmitterPubSuboptions.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 のバッファを空にしてから完了します。