본문으로 건너뛰기

PubSub

PubSubMastra 이벤트 시스템의 추상 기본 클래스입니다. 이는 모든 게시/구독 백엔드가 구현하는 계약을 정의하므로 나머지 Mastra는 어떤 전송이 사용 중인지 알지 못한 채 이벤트를 게시하고 구독할 수 있습니다.

Mastra는 Workflow 이벤트 처리, 스트리밍 및 구성 요소 간 통신을 위해 내부적으로 pub/sub를 사용합니다. 대부분의 애플리케이션은 기본 EventEmitterPubSub를 사용하며 PubSub를 직접 생성하지 않습니다. 사용자 지정 전송이 필요한 경우에만 이 클래스를 구현하세요. 내장 구현에 대해서는 다음을 참조하세요.EventEmitterPubSub, UnixSocketPubSub, CachingPubSub, RedisStreamsPubSub, 그리고GoogleCloudPubSub.

사용예
사용예에 대한 직접 링크

사용자 지정 백엔드를 추가하려면 PubSub을 확장하고 네 개의 추상 메서드를 구현하세요.

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

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { CustomPubSub } from './pubsub'

export const mastra = new Mastra({
pubsub: new CustomPubSub(),
})

배송 모드
배송 모드에 대한 직접 링크

PubSubsupportedModes 속성을 통해 지원하는 전달 모드를 선언합니다. Mastra는 이 값을 읽어 이벤트를 가져오는 장기 실행 워커를 실행할지 결정합니다.

모드설명
pull소비자가 브로커에서 능동적으로 읽습니다(예: Redis Streams XREADGROUP). Mastra는 읽기 작업을 수행하는 오케스트레이션 워커를 실행합니다.
push소비자가 요청하지 않아도 이벤트가 프로세스 내부 또는 HTTP 엔드포인트를 통해 도착합니다. 읽기 루프가 필요하지 않습니다.
사용자 지정 구현이 푸시 전달을 명시적으로 선택하지 않는 한 현재 동작을 유지하도록 기본값은 ['pull']입니다.

행동 양식
행동 양식에 대한 직접 링크

핵심 방법
핵심 방법에 대한 직접 링크

publish(topic, event)
publishtopic-event에 대한 직접 링크

주제에 이벤트를 게시합니다. idcreatedAt 필드는 구현에서 할당합니다.

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) 호출로 전달됩니다. 백엔드의 supportsNativeBatchingtrue일 때만 배치가 적용됩니다. 다른 백엔드는 이 옵션을 무시하고 이벤트를 한 번에 하나씩 전달합니다.

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:

ReadonlyArray<"pull" | "push">
구현에서 지원하는 전달 모드입니다. 기본값은 ["pull"]입니다.

supportsNativeBatching:

boolean
구현이 subscribe()options.batch를 적용하는지 여부입니다. 기본값은 false입니다. 내부적으로 배치를 통합하는 백엔드는 이를 재정의하고 true를 반환합니다.

유형
유형에 대한 직접 링크

Event
event에 대한 직접 링크

type:

string
이벤트 유형 식별자입니다.

id:

string
게시할 때 구현에서 할당하는 고유 이벤트 ID입니다.

data:

any
이벤트 페이로드입니다.

runId:

string
이벤트가 속한 실행입니다.

createdAt:

Date
게시할 때 구현에서 할당하는 타임스탬프입니다.

index?:

number
특정 오프셋부터 재개하는 데 사용하는 순차적 위치입니다.

deliveryAttempt?:

number
이벤트가 전달된 횟수입니다. 1부터 시작합니다. 백엔드가 재전달을 추적하지 않으면 기본값은 1입니다.

SubscribeOptions
subscribeoptions에 대한 직접 링크

group?:

string
설정하면 같은 그룹의 구독자가 메시지를 두고 경쟁하며 각 이벤트는 한 구성원에게 전달됩니다. 생략하면 모든 구독자가 모든 이벤트를 받습니다.

batch?:

SubscribeBatchOptions
이 구독에서 배치 전달을 사용하도록 설정합니다. 생략하면 이벤트가 한 번에 하나씩 전달됩니다. supportsNativeBatchingtrue인 백엔드에서만 적용됩니다.

SubscribeBatchOptions
subscribebatchoptions에 대한 직접 링크

구독별 일괄 처리 정책. 콜백 서명은 변경되지 않습니다. N개 이벤트 배치는 게시 순서에 따라 N개의 연속 콜백 호출이 됩니다.

maxSize?:

number
강제로 플러시하기 전에 보관할 수 있는 최대 이벤트 수입니다.

maxWaitMs?:

number
가장 오래된 이벤트가 버퍼에 머물 수 있는 최대 시간(밀리초)입니다. 타이머는 버퍼가 빈 상태에서 비어 있지 않은 상태로 전환될 때 시작됩니다.

minIntervalMs?:

number
연속된 배치 전달 사이의 최소 시간(밀리초)입니다. maxSize 또는 maxWaitMs 조건이 충족되더라도 마지막 전달 후 이 시간이 지날 때까지 버퍼를 유지합니다.

isImmediate?:

(event: Event) => boolean
이벤트에 대해 true를 반환하면 minIntervalMs 조건에 따라 게시 즉시 버퍼를 플러시합니다. 이벤트별 예외 처리 수단입니다.

coalesce?:

(events: Event[]) => Event[]
전달 전에 대기 중인 배치에 적용하여 대체된 이벤트를 제거합니다. 참조 동일성을 기준으로 입력의 부분집합을 반환해야 합니다. 새로 생성한 Event 객체를 반환하면 계약 위반으로 간주되어 전체 배치가 폐기됩니다. 유지된 이벤트의 순서는 보존됩니다.

maxBufferSize?:

number
= 256
오버플로 처리가 시작되기 전에 버퍼가 보관할 수 있는 최대 이벤트 수입니다. 즉시 처리로 표시된 이벤트는 오버플로 시에도 삭제되지 않습니다.

overflow?:

"drop-oldest" | "drop-newest" | "coalesce-or-drop-oldest"
= coalesce-or-drop-oldest
버퍼가 maxBufferSize를 초과할 때 사용하는 오버플로 전략입니다. coalesce-or-drop-oldest는 먼저 coalesce를 실행한 다음, 여전히 한도를 초과하면 가장 오래된 이벤트를 삭제합니다.

EventCallback
eventcallback에 대한 직접 링크

구독자를 위한 콜백 서명:(event: Event, ack?: () => Promise<void>, nack?: () => Promise<void>) => void.

event:

Event
전달된 이벤트입니다.

ack?:

() => Promise<void>
처리가 성공했음을 확인합니다. 이벤트가 큐에서 제거됩니다.

nack?:

() => Promise<void>
처리 실패를 확인합니다. 일정 시간이 지난 후 재전달하도록 이벤트가 큐에 다시 추가됩니다.