> Discover all available pages from the documentation index: https://mastra.zisheng.pro/ko/llms.txt # PubSub `PubSub`Mastra 이벤트 시스템의 추상 기본 클래스입니다. 이는 모든 게시/구독 백엔드가 구현하는 계약을 정의하므로 나머지 Mastra는 어떤 전송이 사용 중인지 알지 못한 채 이벤트를 게시하고 구독할 수 있습니다. Mastra는 Workflow 이벤트 처리, 스트리밍 및 구성 요소 간 통신을 위해 내부적으로 pub/sub를 사용합니다. 대부분의 애플리케이션은 기본 [`EventEmitterPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/event-emitter)를 사용하며 `PubSub`를 직접 생성하지 않습니다. 사용자 지정 전송이 필요한 경우에만 이 클래스를 구현하세요. 내장 구현에 대해서는 다음을 참조하세요.[`EventEmitterPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/event-emitter), [`UnixSocketPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/unix-socket-pubsub), [`CachingPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/caching-pubsub), [`RedisStreamsPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/redis-streams), 그리고[`GoogleCloudPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/google-cloud-pubsub). ## 사용예 사용자 지정 백엔드를 추가하려면 `PubSub`을 확장하고 네 개의 추상 메서드를 구현하세요. ```typescript 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): Promise { // Deliver the event to subscribers of `topic`. } async subscribe(topic: string, cb: EventCallback, options?: SubscribeOptions): Promise { // Register `cb` to receive events published to `topic`. } async unsubscribe(topic: string, cb: EventCallback): Promise { // Remove a previously registered callback. } async flush(): Promise { // Wait for any in-flight deliveries to settle. } } ``` 인스턴스를[Mastra](https://mastra.zisheng.pro/ko/reference/core/mastra-class) constructor: ```typescript 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)` 주제에 이벤트를 게시합니다. `id` 및 `createdAt` 필드는 구현에서 할당합니다. ```typescript await pubsub.publish('my-topic', { type: 'example', data: { value: 1 }, runId: 'run-123', }) ``` #### `subscribe(topic, cb, options?)` 주제에 게시된 이벤트를 수신할 콜백을 등록합니다. `options.group`이 설정되면 같은 그룹의 구독자가 메시지를 두고 경쟁하며 각 이벤트는 한 구성원에게 전달됩니다. 그룹이 없으면 모든 구독자가 모든 이벤트를 받습니다. 배치 전달을 사용하려면 `options.batch`를 전달하세요. 콜백 시그니처는 변경되지 않습니다. N개 이벤트의 배치는 게시 순서대로 N번의 연속된 `cb(event, ack, nack)` 호출로 전달됩니다. 백엔드의 [`supportsNativeBatching`](#properties)이 `true`일 때만 배치가 적용됩니다. 다른 백엔드는 이 옵션을 무시하고 이벤트를 한 번에 하나씩 전달합니다. ```typescript await pubsub.subscribe('my-topic', (event, ack, nack) => { console.log(event) }) ``` #### `unsubscribe(topic, cb)` 주제에서 이전에 등록된 콜백을 제거합니다. ```typescript await pubsub.unsubscribe('my-topic', callback) ``` #### `flush()` 기내 배송이 완료될 때까지 기다립니다. 이벤트 삭제를 방지하려면 종료하기 전에 이 호출을 호출하세요. ```typescript await pubsub.flush() ``` #### `clearTopic(topic)` 더 이상 이벤트가 게시되지 않으면 주제에 대해 유지된 모든 상태(캐시된 기록, 영구 스트림 항목 및 소비자 그룹)를 삭제합니다. Mastra의 실행 라이프사이클(지속성 Agent 및 이벤트 Workflow 엔진)은 실행이 최종 상태에 도달할 때 이를 자동으로 호출하므로 실행별 항목이 메시지를 보관하는 전송에 누적되지 않습니다. 기본 구현은 no-op입니다. 즉, 주제별로 아무것도 보관하지 않는 전송 방식(예: `EventEmitterPubSub`)에는 지울 항목이 없습니다. [`RedisStreamsPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/redis-streams)처럼 메시지를 영속화하는 백엔드는 이를 재정의합니다. 이 계약은 최선형(best-effort) 방식입니다. 호출자가 정리 경계에서 이를 fire-and-forget 방식으로 호출하므로, 구현은 실패를 throw하는 대신 로그에 기록합니다. ```typescript await pubsub.clearTopic('workflow.events.v2.run-123') ``` ### 재생 방법 이러한 메서드는 연결이 끊긴 후 스트림 재개를 지원합니다. 기본 구현은 일반 `subscribe`로 대체되므로 기록을 지원하지 않는 백엔드는 실시간 전용으로 작동합니다. [`CachingPubSub`](https://mastra.zisheng.pro/ko/reference/pubsub/caching-pubsub)은 캐시된 이벤트를 재생하도록 이 메서드를 재정의합니다. #### `getHistory(topic, offset?)` `offset`부터 시작하여 주제에 대해 캐시된 이벤트를 반환합니다. 백엔드에 기록이 없으면 빈 배열을 반환합니다. ```typescript const events = await pubsub.getHistory('my-topic', 0) ``` 보고:`Promise` #### `subscribeWithReplay(topic, cb)` 캐시된 이벤트를 재생한 다음 라이브 이벤트를 구독합니다. ```typescript await pubsub.subscribeWithReplay('my-topic', event => { console.log(event) }) ``` #### `subscribeFromOffset(topic, offset, cb)` 알려진 위치에서 시작하여 캐시된 이벤트를 재생한 다음 라이브 이벤트를 구독합니다. 이는 클라이언트가 마지막 위치를 알고 있는 경우 전체 재생보다 더 효율적입니다. ```typescript await pubsub.subscribeFromOffset('my-topic', 42, event => { console.log(event) }) ``` ## 속성 **supportedModes** (`ReadonlyArray<"pull" | "push">`): 구현에서 지원하는 전달 모드입니다. 기본값은 \["pull"]입니다. **supportsNativeBatching** (`boolean`): 구현이 subscribe()의 options.batch를 적용하는지 여부입니다. 기본값은 false입니다. 내부적으로 배치를 통합하는 백엔드는 이를 재정의하고 true를 반환합니다. ## 유형 ### `Event` **type** (`string`): 이벤트 유형 식별자입니다. **id** (`string`): 게시할 때 구현에서 할당하는 고유 이벤트 ID입니다. **data** (`any`): 이벤트 페이로드입니다. **runId** (`string`): 이벤트가 속한 실행입니다. **createdAt** (`Date`): 게시할 때 구현에서 할당하는 타임스탬프입니다. **index** (`number`): 특정 오프셋부터 재개하는 데 사용하는 순차적 위치입니다. **deliveryAttempt** (`number`): 이벤트가 전달된 횟수입니다. 1부터 시작합니다. 백엔드가 재전달을 추적하지 않으면 기본값은 1입니다. ### `SubscribeOptions` **group** (`string`): 설정하면 같은 그룹의 구독자가 메시지를 두고 경쟁하며 각 이벤트는 한 구성원에게 전달됩니다. 생략하면 모든 구독자가 모든 이벤트를 받습니다. **batch** (`SubscribeBatchOptions`): 이 구독에서 배치 전달을 사용하도록 설정합니다. 생략하면 이벤트가 한 번에 하나씩 전달됩니다. supportsNativeBatching이 true인 백엔드에서만 적용됩니다. ### `SubscribeBatchOptions` 구독별 일괄 처리 정책. 콜백 서명은 변경되지 않습니다. N개 이벤트 배치는 게시 순서에 따라 N개의 연속 콜백 호출이 됩니다. **maxSize** (`number`): 강제로 플러시하기 전에 보관할 수 있는 최대 이벤트 수입니다. **maxWaitMs** (`number`): 가장 오래된 이벤트가 버퍼에 머물 수 있는 최대 시간(밀리초)입니다. 타이머는 버퍼가 빈 상태에서 비어 있지 않은 상태로 전환될 때 시작됩니다. **minIntervalMs** (`number`): 연속된 배치 전달 사이의 최소 시간(밀리초)입니다. maxSize 또는 maxWaitMs 조건이 충족되더라도 마지막 전달 후 이 시간이 지날 때까지 버퍼를 유지합니다. **isImmediate** (`(event: Event) => boolean`): 이벤트에 대해 true를 반환하면 minIntervalMs 조건에 따라 게시 즉시 버퍼를 플러시합니다. 이벤트별 예외 처리 수단입니다. **coalesce** (`(events: Event[]) => Event[]`): 전달 전에 대기 중인 배치에 적용하여 대체된 이벤트를 제거합니다. 참조 동일성을 기준으로 입력의 부분집합을 반환해야 합니다. 새로 생성한 Event 객체를 반환하면 계약 위반으로 간주되어 전체 배치가 폐기됩니다. 유지된 이벤트의 순서는 보존됩니다. **maxBufferSize** (`number`): 오버플로 처리가 시작되기 전에 버퍼가 보관할 수 있는 최대 이벤트 수입니다. 즉시 처리로 표시된 이벤트는 오버플로 시에도 삭제되지 않습니다. (Default: `256`) **overflow** (`"drop-oldest" | "drop-newest" | "coalesce-or-drop-oldest"`): 버퍼가 maxBufferSize를 초과할 때 사용하는 오버플로 전략입니다. coalesce-or-drop-oldest는 먼저 coalesce를 실행한 다음, 여전히 한도를 초과하면 가장 오래된 이벤트를 삭제합니다. (Default: `coalesce-or-drop-oldest`) ### `EventCallback` 구독자를 위한 콜백 서명:`(event: Event, ack?: () => Promise, nack?: () => Promise) => void`. **event** (`Event`): 전달된 이벤트입니다. **ack** (`() => Promise`): 처리가 성공했음을 확인합니다. 이벤트가 큐에서 제거됩니다. **nack** (`() => Promise`): 처리 실패를 확인합니다. 일정 시간이 지난 후 재전달하도록 이벤트가 큐에 다시 추가됩니다.