PubSub
PubSub は Mastra のイベントシステムの抽象基底クラスです。すべての pub/sub backend が実装する規約を定義するため、Mastra の他の部分は、使用中の transport を意識せずにイベントを publish および subscribe できます。
Mastra は Workflow のイベント処理、ストリーミング、コンポーネント間通信に pub/sub を内部で使用します。ほとんどのアプリケーションはデフォルトの EventEmitterPubSub を使用し、PubSub を直接構築することはありません。カスタム transport が必要な場合にのみ、このクラスを実装してください。
組み込み実装については、EventEmitterPubSub、UnixSocketPubSub、CachingPubSub、RedisStreamsPubSub、GoogleCloudPubSub を参照してください。
使用例使用例への直接リンク
カスタム backend を追加するには、PubSub を拡張して4つの抽象メソッドを実装します。
import { PubSub } from '@mastra/core/events'
import type { Event, EventCallback, SubscribeOptions } from '@mastra/core/events'
export class CustomPubSub extends PubSub {
async publish(topic: string, event: Omit<Event, 'id' | 'createdAt'>): Promise<void> {
// Deliver the event to subscribers of `topic`.
}
async subscribe(topic: string, cb: EventCallback, options?: SubscribeOptions): Promise<void> {
// Register `cb` to receive events published to `topic`.
}
async unsubscribe(topic: string, cb: EventCallback): Promise<void> {
// Remove a previously registered callback.
}
async flush(): Promise<void> {
// Wait for any in-flight deliveries to settle.
}
}
インスタンスを Mastra コンストラクターに渡します。
import { Mastra } from '@mastra/core'
import { CustomPubSub } from './pubsub'
export const mastra = new Mastra({
pubsub: new CustomPubSub(),
})
配信モード配信モードへの直接リンク
PubSub は、対応する配信モードを supportedModes プロパティで宣言します。Mastra はこの値を読み取り、イベントを pull する長時間稼働 worker を実行するかどうかを判断します。
| モード | 説明 |
|---|---|
pull | consumer が broker から能動的に読み取ります。たとえば Redis Streams の XREADGROUP が該当します。Mastra は読み取り用の orchestration worker を実行します。 |
push | consumer が要求しなくても、インプロセスまたは HTTP endpoint 経由でイベントが到着します。読み取り loop は不要です。 |
カスタム実装が明示的に push 配信を選択しない限り現在の動作を維持できるよう、デフォルトは ['pull'] です。
メソッドメソッドへの直接リンク
コアメソッドコアメソッドへの直接リンク
publish(topic, event)publishtopic-eventへの直接リンク
イベントをトピックに公開します。id フィールドと createdAt フィールドは実装によって割り当てられます。
await pubsub.publish('my-topic', {
type: 'example',
data: { value: 1 },
runId: 'run-123',
})
subscribe(topic, cb, options?)subscribetopic-cb-optionsへの直接リンク
トピックに公開されたイベントを受信する callback を登録します。options.group を設定すると、同じグループの subscriber がメッセージの取得を競合し、各イベントは1つのメンバーに配信されます。グループを指定しない場合は、すべての subscriber がすべてのイベントを受信します。
バッチ配信を選択するには options.batch を渡します。callback のシグネチャは変わりません。N 件のイベントからなるバッチは、publish 順に N 回連続する cb(event, ack, nack) 呼び出しとして配信されます。バッチ処理は、backend の supportsNativeBatching が true の場合にのみ適用されます。それ以外の backend はこのオプションを無視し、イベントを1件ずつ配信します。
await pubsub.subscribe('my-topic', (event, ack, nack) => {
console.log(event)
})
unsubscribe(topic, cb)unsubscribetopic-cbへの直接リンク
以前に登録した callback をトピックから削除します。
await pubsub.unsubscribe('my-topic', callback)
flush()flushへの直接リンク
処理中の配信が完了するまで待機します。イベントの欠落を防ぐため、シャットダウン前に呼び出してください。
await pubsub.flush()
clearTopic(topic)cleartopictopicへの直接リンク
トピックへのイベント公開がすべて終わった後、保持されているトピックの状態(キャッシュ済み履歴、永続 stream entry、consumer group)をすべて削除します。Mastra の run lifecycle(durable Agent とイベント駆動 Workflow engine)は、run が終了状態に達したときにこのメソッドを自動的に呼び出します。そのため、メッセージを保持する transport に run ごとのトピックが蓄積しません。
デフォルト実装は何もしません。EventEmitterPubSub のようにトピックごとの状態を何も保持しない transport には、削除するものがないためです。RedisStreamsPubSub のようにメッセージを永続化する backend は、このメソッドをオーバーライドします。この規約はベストエフォートです。caller はクリーンアップ境界で fire-and-forget として呼び出すため、実装は失敗をスローせずログに記録します。
await pubsub.clearTopic('workflow.events.v2.run-123')
リプレイメソッドリプレイメソッドへの直接リンク
これらのメソッドは、切断後のストリーム再開に対応します。デフォルト実装は通常の subscribe にフォールバックするため、履歴に対応していない backend は live イベントのみを扱います。CachingPubSub はこれらをオーバーライドし、キャッシュ済みイベントをリプレイします。
getHistory(topic, offset?)gethistorytopic-offsetへの直接リンク
offset から始まるトピックのキャッシュ済みイベントを返します。backend に履歴がない場合は空の配列を返します。
const events = await pubsub.getHistory('my-topic', 0)
戻り値:Promise<Event[]>
subscribeWithReplay(topic, cb)subscribewithreplaytopic-cbへの直接リンク
キャッシュ済みイベントをリプレイしてから、live イベントを購読します。
await pubsub.subscribeWithReplay('my-topic', event => {
console.log(event)
})
subscribeFromOffset(topic, offset, cb)subscribefromoffsettopic-offset-cbへの直接リンク
既知の位置からキャッシュ済みイベントをリプレイし、その後 live イベントを購読します。クライアントが最後の位置を把握している場合、完全なリプレイより効率的です。
await pubsub.subscribeFromOffset('my-topic', 42, event => {
console.log(event)
})
プロパティプロパティへの直接リンク
supportedModes:
["pull"] です。supportsNativeBatching:
options.batch を subscribe() で処理するかどうか。デフォルトは false です。バッチ処理を内部に統合する backend はこの値をオーバーライドし、true を返します。型型への直接リンク
Eventeventへの直接リンク
type:
id:
data:
runId:
createdAt:
index?:
deliveryAttempt?:
SubscribeOptionssubscribeoptionsへの直接リンク
group?:
batch?:
supportsNativeBatching が true の backend でのみ処理されます。SubscribeBatchOptionssubscribebatchoptionsへの直接リンク
subscription ごとのバッチ処理ポリシーです。callback のシグネチャは変わりません。N 件のイベントからなるバッチは、publish 順に N 回連続する callback 呼び出しになります。
maxSize?:
maxWaitMs?:
minIntervalMs?:
maxSize または maxWaitMs の条件を満たした場合でも、前回の配信からこの間隔が経過するまでバッファに保持します。isImmediate?:
true を返すと、minIntervalMs に従い、publish 時にバッファを即座に flush します。イベントごとの escape hatch です。coalesce?:
Event オブジェクトを返すと規約違反になり、バッチ全体が破棄されます。保持するイベントの順序は維持されます。maxBufferSize?:
overflow?:
maxBufferSize を超えた場合の overflow 戦略。coalesce-or-drop-oldest は、最初に coalesce を実行し、それでも上限を超えている場合は最も古いイベントを破棄します。EventCallbackeventcallbackへの直接リンク
subscriber 用 callback のシグネチャ:(event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void。