EventEmitterPubSub
EventEmitterPubSub기본값입니다PubSub구현. Node.js를 사용하여 진행 중인 이벤트를 전달합니다.EventEmitter, 따라서 외부 서비스 없이 작동합니다.
단일 프로세스 애플리케이션에 사용하세요. 한 호스트의 여러 프로세스에 전달하려면 UnixSocketPubSub을 참조하세요. 분산 전송에는 RedisStreamsPubSub 또는 GoogleCloudPubSub을 사용하세요.
프로세스 내부에서 작동하므로 이벤트가 유지되거나 다른 프로세스와 공유되지 않습니다. 재개 가능한 스트림에 재생 기능이 필요하면 CachingPubSub로 래핑하세요.
사용예사용예에 대한 직접 링크
pubsub 옵션을 구성하지 않으면 EventEmitterPubSub가 자동으로 사용되므로 대부분의 애플리케이션에서는 이를 직접 생성할 필요가 없습니다. 구성하거나 공유하려는 경우에만 명시적으로 생성하세요.
import { Mastra } from '@mastra/core'
import { EventEmitterPubSub } from '@mastra/core/events'
export const mastra = new Mastra({
pubsub: new EventEmitterPubSub(),
})
애플리케이션의 다른 부분과 이미터를 공유하려면 기존EventEmitter:
import EventEmitter from 'node:events'
import { EventEmitterPubSub } from '@mastra/core/events'
const emitter = new EventEmitter()
const pubsub = new EventEmitterPubSub(emitter)
일괄 전송 오류를 표시하려면 로거를 전달하세요.
import { EventEmitterPubSub } from '@mastra/core/events'
const pubsub = new EventEmitterPubSub(undefined, { logger })
생성자 매개변수생성자 매개변수에 대한 직접 링크
existingEmitter?:
options?:
속성속성에 대한 직접 링크
supportedModes:
["pull", "push"]를 반환합니다. 이미터는 풀 방식 워커에 이벤트를 제공하거나 리스너에 이벤트를 직접 푸시할 수 있습니다.supportsNativeBatching:
true를 반환합니다. 구독자는 options.batch를 사용하여 일괄 전송을 선택할 수 있습니다.행동 양식행동 양식에 대한 직접 링크
EventEmitterPubSub는 PubSub 계약을 구현합니다. 아래 메서드는 이 구현에 특화된 동작을 제공합니다.
subscribe(topic, cb, options?)subscribetopic-cb-options에 대한 직접 링크
주제에 대한 콜백을 등록합니다. options.group이 없으면 모든 구독자가 모든 이벤트를 수신합니다. 그룹을 지정하면 해당 그룹의 멤버에게 이벤트가 라운드 로빈 방식으로 분배됩니다.
일괄 전송을 사용하려면 options.batch를 전달하세요. 아래의 일괄 처리를 참조하세요.
await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})
flush()flush에 대한 직접 링크
완료되기 전에 nack로 예약된 재전송이 실행될 때까지 기다립니다.
await pubsub.flush()
close()close에 대한 직접 링크
모든 리스너를 제거하고 보류 중인 재전송을 취소합니다. 정상적인 종료 중에 이를 호출하십시오.
await pubsub.close()
재배송재배송에 대한 직접 링크
그룹 구독자가 nack를 호출하면 잠시 후 이벤트가 그룹에 다시 전달되고 deliveryAttempt 횟수가 증가합니다. ack를 호출하면 해당 이벤트의 추적 정보가 지워집니다. 팬아웃 구독자에게는 각 이벤트가 한 번씩 모두 전달되므로 아무 작업도 하지 않는 ack 및 nack 함수가 제공됩니다.
일괄 처리일괄 처리에 대한 직접 링크
EventEmitterPubSub는 options.batch를 기본적으로 지원합니다. 구독자가 이 옵션을 선택하면 이벤트는 구독자별 Memory 버퍼에 보관되며 플러시 조건이 충족될 때 연속된 콜백 호출로 전달됩니다. 팬아웃 구독자와 그룹 구독자 모두 일괄 처리를 사용할 수 있습니다. 전체 정책은 SubscribeBatchOptions를 참조하세요.
await pubsub.subscribe(
'workflow.events',
event => {
console.log(event)
},
{
batch: {
maxSize: 10, // flush once 10 events have queued
maxWaitMs: 500, // ...or after 500ms, whichever comes first
},
},
)
버퍼는 Memory에 프로세스별로 존재하므로 일괄 처리 상태가 유지되거나 재시작 후 보존되지 않습니다. maxWaitMs로 트리거되는 플러시는 최선형 방식으로 동작합니다. 예외가 발생하는 coalesce 같은 단계가 실패하면 오류를 던지는 대신 구성된 logger를 통해 표시합니다. flush()는 완료되기 전에 모든 일괄 구독자 버퍼를 비웁니다.