본문으로 건너뛰기

PubSub

Mastra는 게시/구독(pub/sub) 시스템을 내부 이벤트 버스로 사용합니다. 구성 요소는 주제에 이벤트를 게시하고 다른 구성 요소는 해당 주제를 구독하여 반응합니다. 구성하는 백엔드는 해당 이벤트가 한 프로세스 내에서, 한 호스트의 프로세스 간에, 또는 별도의 인스턴스 간에 이동하는 거리를 결정합니다.

백엔드는 Mastra 인스턴스에 한 번 구성하면 나머지 시스템에서 변경 없이 사용합니다. 기본적으로 Mastra는 별도의 설정이 필요 없는 프로세스 내 백엔드를 사용합니다.

Mastra가 PubSub를 사용하는 방법
Mastra가 PubSub를 사용하는 방법에 대한 직접 링크

여러 내장 시스템이 동일한 게시/구독 버스를 통해 이벤트를 게시하고 구독합니다.

  • Workflow 실행: 스케줄러가 workflow.start 이벤트를 게시하면 장기 실행 워커가 이를 소비하여 Workflow를 실행합니다. 실행 중에는 단계 및 수명 주기 이벤트가 pub/sub을 통해 전달됩니다.
  • 예약된 Workflow: 스케줄러는 Workflow 토픽에 게시하여 실행 시점이 된 작업을 전달합니다. 예약된 Workflow를 참조하세요.
  • 백그라운드 작업: 작업 관리자는 작업을 워커 그룹에 배정하고 작업 수명 주기 업데이트를 구독자에게 전달합니다. 이를 통해 백그라운드 작업 스트림이 실시간으로 유지됩니다. 백그라운드 작업 스트리밍을 참조하세요.
  • Agent 신호: 활성 Agent 실행에 신호를 보내면 스레드 토픽에 이벤트가 게시되므로 다른 프로세스에서 실행 중인 작업도 신호를 받습니다.
  • 재개 가능한 스트림: 스트림 청크가 실행별로 게시되므로 다시 연결한 클라이언트가 놓친 내용을 재생할 수 있습니다. 이러한 시스템은 하나의 버스에서 실행되기 때문에 선택한 백엔드는 모든 시스템에 동시에 적용됩니다.

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

백엔드는 다음과 같이 정의된 두 가지 모드 중 하나로 이벤트를 전달합니다.PubSub contract:

  • : 소비자가 백엔드에서 직접 읽으며, Mastra는 장기 실행 워커 루프를 사용하여 이를 수행합니다. RedisStreamsPubSub 같은 분산 백엔드는 이 모드를 사용합니다.
  • 푸시: 소비자가 요청하지 않아도 프로세스 내부 또는 HTTP를 통해 이벤트가 도착합니다. 기본 EventEmitterPubSub은 프로세스 내부에서 이 방식으로 이벤트를 전달합니다. 구독자는 소비자 그룹을 통해 작업을 배포할 수도 있습니다. 같은 그룹의 구성원은 이벤트를 분할하여 각 이벤트가 한 번 처리됩니다. 그룹이 없는 구독자는 모든 이벤트를 수신하며, 이 이벤트는 그룹화되지 않은 모든 구독자에게 스트림을 전달합니다.

기본 백엔드
기본 백엔드에 대한 직접 링크

pubsub 옵션을 설정하지 않으면 Mastra는 EventEmitterPubSub을 사용합니다. Node.js EventEmitter를 사용하여 프로세스 내부에서 이벤트를 전달하므로 외부 서비스 없이 작동합니다.

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

// No pubsub option: Mastra uses EventEmitterPubSub
export const mastra = new Mastra({})

프로세스 내에서 실행되기 때문에 이벤트는 지속되지 않으며 다른 프로세스에 도달하지 않습니다. 기본값은 대부분의 애플리케이션을 포괄하는 단일 인스턴스에 적합합니다.

백엔드 선택
백엔드 선택에 대한 직접 링크

백엔드를 선택하려면 Mastra 인스턴스에 pubsub 옵션을 설정합니다. 각 백엔드는 동일한 PubSub 규약을 구현하므로 애플리케이션의 나머지 부분은 변경되지 않습니다. 결정적인 질문은 구독자가 실행되는 위치입니다.

백엔드범위모드패키지
EventEmitterPubSub단일 프로세스풀 및 푸시@mastra/core
UnixSocketPubSub단일 호스트의 여러 프로세스푸시@mastra/core
RedisStreamsPubSub여러 호스트에 분산@mastra/redis-streams
GoogleCloudPubSub여러 호스트에 분산@mastra/google-cloud-pubsub

하나의 호스트에 여러 프로세스가 있음
하나의 호스트에 여러 프로세스가 있음에 대한 직접 링크

같은 머신의 여러 프로세스가 스트림을 공유해야 할 때 UnixSocketPubSub을 사용합니다. Unix 도메인 소켓을 통해 이벤트를 전달하고 하나의 프로세스를 브로커로 선출합니다. 브로커가 종료되면 나머지 프로세스가 새 브로커를 선출합니다.

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { UnixSocketPubSub } from '@mastra/core/events'

export const mastra = new Mastra({
pubsub: new UnixSocketPubSub('/tmp/mastra/events.sock'),
})

분산 배포
분산 배포에 대한 직접 링크

둘 이상의 인스턴스나 호스트를 실행할 때 분산 백엔드를 사용하면 모든 인스턴스가 동일한 이벤트를 수신하게 됩니다. 이는 한 인스턴스에서 처리한 요청이 다른 인스턴스에서 실행 중인 작업에 도달해야 할 때마다 중요합니다.

예를 들어 Agent 실행에 신호를 보내려면 신호 이벤트가 프로세스 경계를 넘어 해당 실행을 소유한 인스턴스로 전달되어야 합니다. 프로세스 내 기본값을 사용하면 해당 인스턴스가 이벤트를 수신하지 못합니다. 아래의 두 백엔드는 프로세스와 호스트 전체에 걸쳐 전달되며 재전송을 위해 이벤트를 유지합니다.

다음 예제에서는Redis Streams:

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { RedisStreamsPubSub } from '@mastra/redis-streams'

export const mastra = new Mastra({
pubsub: new RedisStreamsPubSub({
url: 'redis://localhost:6379',
}),
})

다음 예제에서는Google Cloud Pub/Sub:

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { GoogleCloudPubSub } from '@mastra/google-cloud-pubsub'

export const mastra = new Mastra({
pubsub: new GoogleCloudPubSub({
projectId: 'my-project',
}),
})

재개 가능한 스트림
재개 가능한 스트림에 대한 직접 링크

재개 가능한 스트림을 사용하면 클라이언트가 다시 연결하여 놓친 이벤트를 재생할 수 있으며, 그러려면 백엔드가 최근 기록을 유지해야 합니다. RedisStreamsPubSub 같은 분산 백엔드는 이벤트를 영구 저장하므로 자체적으로 재생을 지원합니다. 실시간 전달은 기록을 유지하지 않습니다. EventEmitterPubSub에 재생 기능을 추가하려면 이를 CachingPubSub으로 래핑하세요. 그러면 게시된 이벤트가 토픽별로 기록되어, 늦게 구독하거나 다시 연결한 구독자가 실시간 이벤트를 계속 받기 전에 누락된 이벤트를 따라잡을 수 있습니다.

src/mastra/index.ts
import { Mastra } from '@mastra/core'
import { CachingPubSub, EventEmitterPubSub } from '@mastra/core/events'
import { InMemoryServerCache } from '@mastra/core/cache'

const cache = new InMemoryServerCache()

export const mastra = new Mastra({
pubsub: new CachingPubSub(new EventEmitterPubSub(), cache),
})

전체 전달 규약과 각 백엔드의 구성 옵션은 PubSub 레퍼런스를 참조하세요.