PubSub
PubSubMastra 이벤트 시스템의 추상 기본 클래스입니다. 이는 모든 게시/구독 백엔드가 구현하는 계약을 정의하므로 나머지 Mastra는 어떤 전송이 사용 중인지 알지 못한 채 이벤트를 게시하고 구독할 수 있습니다.
Mastra는 Workflow 이벤트 처리, 스트리밍 및 구성 요소 간 통신을 위해 내부적으로 pub/sub를 사용합니다. 대부분의 애플리케이션은 기본 EventEmitterPubSub를 사용하며 PubSub를 직접 생성하지 않습니다. 사용자 지정 전송이 필요한 경우에만 이 클래스를 구현하세요.
내장 구현에 대해서는 다음을 참조하세요.EventEmitterPubSub, UnixSocketPubSub, CachingPubSub, RedisStreamsPubSub, 그리고GoogleCloudPubSub.
사용예사용예에 대한 직접 링크
사용자 지정 백엔드를 추가하려면 PubSub을 확장하고 네 개의 추상 메서드를 구현하세요.
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 constructor:
import { Mastra } from '@mastra/core'
import { CustomPubSub } from './pubsub'
export const mastra = new Mastra({
pubsub: new CustomPubSub(),
})
배송 모드배송 모드에 대한 직접 링크
PubSub은 supportedModes 속성을 통해 지원하는 전달 모드를 선언합니다. Mastra는 이 값을 읽어 이벤트를 가져오는 장기 실행 워커를 실행할지 결정합니다.
| 모드 | 설명 |
|---|---|
pull | 소비자가 브로커에서 능동적으로 읽습니다(예: Redis Streams XREADGROUP). Mastra는 읽기 작업을 수행하는 오케스트레이션 워커를 실행합니다. |
push | 소비자가 요청하지 않아도 이벤트가 프로세스 내부 또는 HTTP 엔드포인트를 통해 도착합니다. 읽기 루프가 필요하지 않습니다. |
사용자 지정 구현이 푸시 전달을 명시적으로 선택하지 않는 한 현재 동작을 유지하도록 기본값은 ['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에 대한 직접 링크
주제에 게시된 이벤트를 수신할 콜백을 등록합니다. options.group이 설정되면 같은 그룹의 구독자가 메시지를 두고 경쟁하며 각 이벤트는 한 구성원에게 전달됩니다. 그룹이 없으면 모든 구독자가 모든 이벤트를 받습니다.
배치 전달을 사용하려면 options.batch를 전달하세요. 콜백 시그니처는 변경되지 않습니다. N개 이벤트의 배치는 게시 순서대로 N번의 연속된 cb(event, ack, nack) 호출로 전달됩니다. 백엔드의 supportsNativeBatching이 true일 때만 배치가 적용됩니다. 다른 백엔드는 이 옵션을 무시하고 이벤트를 한 번에 하나씩 전달합니다.
await pubsub.subscribe('my-topic', (event, ack, nack) => {
console.log(event)
})
unsubscribe(topic, cb)unsubscribetopic-cb에 대한 직접 링크
주제에서 이전에 등록된 콜백을 제거합니다.
await pubsub.unsubscribe('my-topic', callback)
flush()flush에 대한 직접 링크
기내 배송이 완료될 때까지 기다립니다. 이벤트 삭제를 방지하려면 종료하기 전에 이 호출을 호출하세요.
await pubsub.flush()
clearTopic(topic)cleartopictopic에 대한 직접 링크
더 이상 이벤트가 게시되지 않으면 주제에 대해 유지된 모든 상태(캐시된 기록, 영구 스트림 항목 및 소비자 그룹)를 삭제합니다. Mastra의 실행 라이프사이클(지속성 Agent 및 이벤트 Workflow 엔진)은 실행이 최종 상태에 도달할 때 이를 자동으로 호출하므로 실행별 항목이 메시지를 보관하는 전송에 누적되지 않습니다.
기본 구현은 no-op입니다. 즉, 주제별로 아무것도 보관하지 않는 전송 방식(예: EventEmitterPubSub)에는 지울 항목이 없습니다. RedisStreamsPubSub처럼 메시지를 영속화하는 백엔드는 이를 재정의합니다. 이 계약은 최선형(best-effort) 방식입니다. 호출자가 정리 경계에서 이를 fire-and-forget 방식으로 호출하므로, 구현은 실패를 throw하는 대신 로그에 기록합니다.
await pubsub.clearTopic('workflow.events.v2.run-123')
재생 방법재생 방법에 대한 직접 링크
이러한 메서드는 연결이 끊긴 후 스트림 재개를 지원합니다. 기본 구현은 일반 subscribe로 대체되므로 기록을 지원하지 않는 백엔드는 실시간 전용으로 작동합니다. CachingPubSub은 캐시된 이벤트를 재생하도록 이 메서드를 재정의합니다.
getHistory(topic, offset?)gethistorytopic-offset에 대한 직접 링크
offset부터 시작하여 주제에 대해 캐시된 이벤트를 반환합니다. 백엔드에 기록이 없으면 빈 배열을 반환합니다.
const events = await pubsub.getHistory('my-topic', 0)
보고:Promise<Event[]>
subscribeWithReplay(topic, cb)subscribewithreplaytopic-cb에 대한 직접 링크
캐시된 이벤트를 재생한 다음 라이브 이벤트를 구독합니다.
await pubsub.subscribeWithReplay('my-topic', event => {
console.log(event)
})
subscribeFromOffset(topic, offset, cb)subscribefromoffsettopic-offset-cb에 대한 직접 링크
알려진 위치에서 시작하여 캐시된 이벤트를 재생한 다음 라이브 이벤트를 구독합니다. 이는 클라이언트가 마지막 위치를 알고 있는 경우 전체 재생보다 더 효율적입니다.
await pubsub.subscribeFromOffset('my-topic', 42, event => {
console.log(event)
})
속성속성에 대한 직접 링크
supportedModes:
["pull"]입니다.supportsNativeBatching:
subscribe()의 options.batch를 적용하는지 여부입니다. 기본값은 false입니다. 내부적으로 배치를 통합하는 백엔드는 이를 재정의하고 true를 반환합니다.유형유형에 대한 직접 링크
Eventevent에 대한 직접 링크
type:
id:
data:
runId:
createdAt:
index?:
deliveryAttempt?:
SubscribeOptionssubscribeoptions에 대한 직접 링크
group?:
batch?:
supportsNativeBatching이 true인 백엔드에서만 적용됩니다.SubscribeBatchOptionssubscribebatchoptions에 대한 직접 링크
구독별 일괄 처리 정책. 콜백 서명은 변경되지 않습니다. N개 이벤트 배치는 게시 순서에 따라 N개의 연속 콜백 호출이 됩니다.
maxSize?:
maxWaitMs?:
minIntervalMs?:
maxSize 또는 maxWaitMs 조건이 충족되더라도 마지막 전달 후 이 시간이 지날 때까지 버퍼를 유지합니다.isImmediate?:
true를 반환하면 minIntervalMs 조건에 따라 게시 즉시 버퍼를 플러시합니다. 이벤트별 예외 처리 수단입니다.coalesce?:
Event 객체를 반환하면 계약 위반으로 간주되어 전체 배치가 폐기됩니다. 유지된 이벤트의 순서는 보존됩니다.maxBufferSize?:
overflow?:
maxBufferSize를 초과할 때 사용하는 오버플로 전략입니다. coalesce-or-drop-oldest는 먼저 coalesce를 실행한 다음, 여전히 한도를 초과하면 가장 오래된 이벤트를 삭제합니다.EventCallbackeventcallback에 대한 직접 링크
구독자를 위한 콜백 서명:(event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void.