본문으로 건너뛰기

RedisStreamsPubSub

RedisStreamsPubSubPubSub구현 지원레디스 스트림. 지속성, 소비자 그룹 및 실패 시 재전송을 통해 프로세스와 호스트 전반에 걸쳐 이벤트를 전달합니다. 또한 구현LeaseProvider, 따라서 신호 계층은 인스턴스 전체에서 리소스당 단일 소유자를 선택할 수 있으며, 이를 통해 신호는 분산 및 서버리스 배포에서 실행을 조정할 수 있습니다.

여러 서비스가 이벤트 스트림을 공유하는 분산 배포에 사용하세요. 단일 프로세스 전송에는 EventEmitterPubSub를 사용하세요. Google Cloud에는 GoogleCloudPubSub을 사용하세요. 각 주제는 Redis 스트림 키에 매핑됩니다. 그룹 구독은 Redis 소비자 그룹을 사용하므로 구성원은 라운드 로빈 방식으로 작업을 공유합니다. 그룹이 없는 구독은 비공개 소비자 그룹을 생성하므로 모든 구독자가 모든 이벤트를 수신합니다.

RedisStreamsPubSub는 풀 방식 전송입니다. 소비자가 XREADGROUP을 사용하여 이벤트를 읽으므로 Mastra가 소비자를 대신해 오케스트레이션 워커를 실행합니다.

설치
설치에 대한 직접 링크

npm install @mastra/redis-streams

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

Redis 연결 URL을 제공하세요.

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',
}),
})

생성자 매개변수
생성자 매개변수에 대한 직접 링크

url?:

string
= redis://localhost:6379
Redis 연결 URL입니다. 생략하면 redisOptions.url을 사용합니다.

keyPrefix?:

string
= mastra:topic
스트림 키의 접두사입니다. 각 주제는 <keyPrefix>:<topic>에 매핑됩니다.

blockMs?:

number
= 1000
각 읽기 작업이 새 이벤트를 기다리며 차단되는 밀리초 단위 시간입니다.

redisOptions?:

RedisClientOptions
고급 구성을 위해 기본 redis 클라이언트에 전달하는 옵션입니다.

maxStreamLength?:

number
= 10000
스트림별로 유지할 대략적인 최대 항목 수입니다. 잘라내기를 비활성화하려면 0으로 설정하세요.

streamIdleTtlMs?:

number
= 0
밀리초 단위 유휴 만료 시간으로, 쓰기 작업(게시, nack 재시도, 그룹 재생성)마다 갱신되는 슬라이딩 TTL입니다. 쓰기 작업마다 재설정되므로 활발하게 쓰는 스트림이 처리 도중 만료되지 않으며, 전체 기간 동안 유휴 상태인 스트림은 Redis가 자동으로 삭제합니다. 쓰기 작업만 TTL을 갱신한다는 점에 유의하세요. 소비자가 백로그를 천천히 소진하는 작업은 TTL을 갱신하지 않으므로 활성 주제에서 예상되는 가장 긴 쓰기 간격보다 훨씬 크게 설정하세요. 이는 기본 정리 수단이 아닌 안전장치입니다. 정상적인 수명 주기 종료 시에는 clearTopic이 삭제를 처리하며, 이 옵션은 clearTopic 호출에 도달하지 못하는 스트림(예: 중단된 실행)의 Memory만 제한합니다. 음이 아닌 정수여야 합니다. 기본값은 0(비활성화)입니다.

reclaimIntervalMs?:

number
= 30000
이전 소비자가 읽었지만 승인하지 않은 이벤트를 구독이 회수하는 밀리초 단위 주기입니다. 비활성화하려면 0으로 설정하세요.

reclaimIdleMs?:

number
= 60000
보류 중인 이벤트를 회수할 수 있게 되기까지 필요한 최소 유휴 시간(밀리초)입니다. 중복 전송을 방지하려면 일반적인 처리 시간보다 훨씬 크게 설정하세요.

maxDeliveryAttempts?:

number
= 5
이벤트가 삭제되기 전에 nack를 통해 재전송되는 최대 횟수입니다. 제한을 비활성화하려면 Infinity를 전달하세요.

logger?:

{ debug?: Function; warn?: Function }
진단용 선택적 로거입니다. 생략하면 억제된 오류가 출력되지 않습니다.

속성
속성에 대한 직접 링크

supportedModes:

ReadonlyArray<"pull" | "push">
["pull"]을 반환합니다.

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

RedisStreamsPubSubPubSub 계약을 구현합니다. 아래 메서드는 이 구현에 특화된 동작을 제공합니다.

subscribe(topic, cb, options?)
subscribetopic-cb-options에 대한 직접 링크

주제를 구독합니다. options.group을 지정하면 그룹 멤버가 Redis 소비자 그룹을 통해 이벤트를 공유합니다. 그룹이 없으면 구독자가 비공개 소비자 그룹을 통해 모든 이벤트를 수신합니다.

await pubsub.subscribe('workflow.events', (event, ack, nack) => {
console.log(event)
})

flush()
flush에 대한 직접 링크

진행 중인 게시가 완료될 때까지 기다립니다.

await pubsub.flush()

clearTopic(topic)
cleartopictopic에 대한 직접 링크

주제의 스트림과 그 안에 있는 모든 소비자 그룹을 삭제하여 완성된 주제가 보유할 Memory를 확보합니다. Mastra의 실행 수명 주기(지속성 Agent 및 이벤트 Workflow 엔진)는 실행이 최종 상태에 도달하면 이를 자동으로 호출합니다. 아무 것도 주제를 다시 읽지 않을 때만 직접 호출하십시오. 최선을 다하고 절대 던지지 않습니다. 실패는 경고 수준에서 기록됩니다. 스트림이 삭제될 때 여전히 연결된 구독자는 자체적으로 복구되지만 삭제된 항목을 놓칩니다.

자동 정리를 사용하려면 @mastra/core@mastra/redis-streams 버전이 모두 clearTopic을 지원해야 합니다. 런타임이 캐싱 계층을 통해 호출을 라우팅하므로 실행 종료 시 스트림을 삭제하려면 두 패키지를 함께 업그레이드하세요.

await pubsub.clearTopic('workflow.events.run-123')

close()
close에 대한 직접 링크

Redis 연결을 닫고 모든 구독을 중지합니다. 정상적인 종료 중에 이를 호출하십시오.

await pubsub.close()

재전송 및 회수
재전송 및 회수에 대한 직접 링크

구독자가 nack를 호출하면 deliveryAttempt가 증가한 상태로 이벤트가 다시 게시되고 원본은 승인됩니다. 이벤트가 maxDeliveryAttempts에 도달하면 재전송되지 않고 삭제됩니다. 이와 별도로 각 구독은 그룹의 이전 소비자가 읽었지만 승인하지 않은 이벤트를 주기적으로 회수하며, 이는 reclaimIntervalMsreclaimIdleMs로 제어합니다.

분산리스
분산리스에 대한 직접 링크

RedisStreamsPubSub는 동일한 Redis 연결을 기반으로 LeaseProvider 계약을 구현합니다. 신호 런타임은 이를 사용해 단일 소유자(일반적으로 스레드 키별)를 선출하므로 여러 인스턴스 중 하나의 프로세스만 깨어나 Agent를 실행하고 나머지는 후속 작업을 보유자에게 라우팅합니다. 이를 통해 서버리스 및 다중 인스턴스 배포에서 신호가 작동합니다. 공유 임대가 없으면 각 인스턴스가 서로 경쟁하는 자체 실행을 시작합니다. 임대 키는 주제와 동일한 keyPrefix를 사용하여 <keyPrefix>:lease:<key> 형식으로 네임스페이스가 지정됩니다. 모든 작업은 원자적입니다. acquireLeaseSET NX PX를 사용하고 자체 TTL을 멱등적으로 갱신하며, releaseLease, renewLease, transferLease는 변경 전에 소유권을 확인하는 Lua 스크립트를 사용하므로 다른 소유자의 동시 갱신을 덮어쓰지 않습니다. 이러한 메서드를 직접 호출할 필요는 없습니다. RedisStreamsPubSubpubsub 백엔드로 구성하기만 하면 런타임이 이 기능을 감지하여 사용합니다. 전체 메서드 계약은 LeaseProvider를 참조하세요.