PROD주권 유럽 BaaS 플랫폼대시보드 열기 →

공학 · 10분 읽음

PostgreSQL CDC with NATS: 실시간 아키텍처

Affane Daylami · Fondateur · 2026년 7월 24일

블로그로 돌아가기

WebSocket을 통해 연결된 Aurabase 클라이언트는 폴링이나 구성할 웹후크 없이 COMMIT 후 몇 밀리초 후에 베이스에 삽입된 라인이 나타나는 것을 확인합니다. 메커니즘은 세 단계로 구성됩니다. 논리적 복제 플러그인은 Postgres WAL을 JSON으로 디코딩하고 하나의 프로세스만 이를 읽을 수 있는 권한을 가지며 NATS는 각 이벤트를 클러스터의 모든 WebSocket 서버에 전달합니다. aura-realtimeservice의 실제 코드에 연결되는 방법은 다음과 같습니다.

이 영어 텍스트는 프랑스어 원본에서 자동으로 생성되었으며 아직 검토되지 않았습니다.
이 페이지는 자동으로 번역되었습니다. 영어 버전은 권위가 있습니다.

여기에는 이상적인 패턴이 없습니다. 각 추출물은 출판일 현재 aura-realtime상자에서 있는 그대로 제공됩니다. 여기에는 명확한 차이가 있습니다. NATS JetStream은 이 파이프라인에서 CDC 이벤트를 전송하는 것이 아닙니다. 여기에서 JetStream이 실제로 수행하는 작업과 그 자리에서 CDC 트래픽을 전달하는 작업이 아래에 자세히 설명되어 있습니다. 오류를 발생시키지 않고 전체 메커니즘을 자동으로 만들 수 있는 검증된 Kubernetes 트랩은 아래에 자세히 설명되어 있습니다.

필수사항
  • Aurabase의 Postgres CDC는 Debezium/Kafka Connect 브리지가 아닌 pgoutput바이너리 프로토콜이 아닌 wal2json 및 pg_logical_slot_get_changes을 통해 WAL을 읽습니다.
  • 논리 슬롯 aura_cdc_slot은 하나의 리더만 허용합니다. aura-realtime-cdc-worker는 K8을 제외한 "항상 리더" 모드를 사용하여 Kubernetes 임대(coordination.k8s.io/v1)를 통해 선택됩니다.
  • ws-front 복제본에 대한 팬아웃은 영구 JetStream 스트림이 아닌 NATS 코어(aura.realtime.>의 단순 게시/구독)를 통해 진행됩니다. 동일한 서비스에서 JetStream은 인스턴스 간 존재 KV를 독점적으로 제공합니다.
  • 각 이벤트는 캐시된 근사치가 아닌 해당 회선의 실제 SET LOCAL ROLE 요청을 통해 전송 직전에 가입자가 RLS에서 다시 확인합니다.
  • 저장소 기록에서 확인된 함정: Postgres Kubernetes 이미지에 wal2json 플러그인이 없으면 슬롯 생성이 실패하고 CDC가 자동으로 비활성 상태로 유지됩니다.
#
엔진 선택

pgoutput이 아닌 wal2json, Debezium이 아님

대부분의 CDC Postgres 파이프라인은 바이너리 논리 복제 프로토콜인 pgoutput를 거친 다음 이를 Kafka로 변환하는 Debezium과 같은 커넥터를 통과합니다. Aurabase는 이 단계를 건너뜁니다. aura-realtime 서비스는 중간 변환 단계 없이 사용 가능한 JSON을 생성하는 wal2json 플러그인을 사용하여 WAL을 논리적 복제로 직접 디코딩합니다.

cdc/postgres.rs (실제 쿼리)sql
-- 없으면 한 번 생성되고, WAL이 만료되면 다시 생성됩니다.
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- 반복 폴링(배치가 비어 있을 때 100ms 백오프)
SELECT data FROM pg_logical_slot_get_changes(
  'aura_cdc_slot', NULL, 1000,
  'include-xids', '0', 'format-version', '2',
  'include-schemas', '1', 'include-types', '0');

코드에서 확인된 세 가지 Postgres 전제 조건: 애플리케이션 역할은 REPLICATION를 전달해야 하고, wal_level = logical은 활성 상태여야 하며, max_slot_wal_keep_size는 WAL 보존을 제한해야 합니다. 이 마지막 제한이 없으면 느린 소비자로 인해 디스크 볼륨이 무한정 증가합니다. 서비스는 슬롯의 wal_status도 모니터링합니다. lost(WAL이 이 제한을 초과하여 제거됨)로 변경되면 슬롯이 자동으로 다시 생성됩니다. 그 동안 소비되지 않은 이벤트의 손실을 가정하고 기록합니다.

3초마다 새로 고쳐지는 동적 테이블 필터를 사용하여 폴링당 1,000개의 변경 사항을 일괄 처리합니다. 프로젝트에서 실시간을 명시적으로 활성화한 테이블만 add-tables를 입력합니다. 이렇게 하면 구독자가 관심을 두지 않는 테이블의 WAL을 디코딩하는 것을 방지할 수 있습니다.

#
단일 생산자

cdc-worker: 중복되지 않은 Kubernetes 임대에 의해 선택됨

Postgres 논리적 복제 슬롯은 한 번에 하나의 활성 드라이브만 허용합니다. 동일한 aura_cdc_slot의 여러 소비자를 병렬로 실행하면 변경 내용이 단순히 복제되는 것이 아니라 변경 순서가 중단됩니다. Aurabase는 WebSocket의 수평 규모를 희생하지 않고 이 제약 조건을 해결하기 위해 aura-realtime를 두 개의 별도 바이너리로 분할했습니다.

Cargo.tomltoml
# cdc-worker: 선출된 리더, CDC를 시작하고 NATS에 게시합니다. WS 서버가 없습니다.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N개의 복제본(HPA), NATS, RLS 필터를 사용하고 WS/SSE를 제공합니다.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker는 임대 Kubernetes 개체 (coordination.k8s.io/v1)를 보유하여 슬롯을 읽을 수 있는 권한만 얻습니다. 그는 그것을 만들려고 시도하거나 만료된 경우 훔친 다음 수명의 3분의 1 후에 갱신합니다. 임대가 손실되면(갱신 실패, 소유자 변경) CancellationToken는 진행 중인 작업을 중단하고 프로세스는 획득 루프로 돌아갑니다. Kubernetes 외부(예: 로컬 테스트 또는 비클러스터 개발 중에) K8s 클라이언트가 연결에 실패하고 서비스가 "항상 선행" 모드로 전환됩니다. 개발에는 유용하지만 여러 복제본이 있는 프로덕션에서 잊어버린 경우에는 올바르지 않습니다.

cdc/leader.rs (실제 서명)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // 모든 rent_duration_secs / 3을 갱신합니다.
  // K8s 제외: 작업(토큰)은 한 번만, 토큰은 취소되지 않음
}
#
실제 팬아웃

CDC용 NATS 코어, 프레즌스용 JetStream

이는 이러한 종류의 파이프라인에서 가장 자주 오해되는 점입니다. NATS와 JetStream은 서로 다른 두 가지이며 CDC 이벤트는 지속적인 JetStream 스트림을 통과하지 않습니다. cdc-worker는 Client::publish_with_headers — NATS 핵심 게시/구독 API를 사용하여 각 이벤트를 게시합니다. 이 API는 JetStream과 별도로 지속성이나 재생 없이 최대 한 번 전달됩니다.

NATS 하트(pub/sub)CDC 이벤트를 ws-front 복제본으로 팬아웃지속성이나 재생 없이 최대 한 번
NATS 제트스트림(KV)인스턴스 간 존재(aura_presence)재생할 CDC 스트림이 아닌 공유 상태 60초
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker 게시(ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front는 전체 하위 트리(NatsRelay::start_relay)를 수신합니다.
client.subscribe("aura.realtime.>").await

각 복제본 ws-front은 하위 트리 aura.realtime.>을 구독하고 주제를 원래 채널 이름으로 디코딩합니다. 그런 다음 연결된 WebSocket에 대해 이벤트를 로컬 tokio::sync::broadcast에 다시 게시합니다. Aura-Origin 헤더는 보내는 인스턴스의 식별자를 전달합니다. 각 복제본은 자체적으로 게시한 메시지를 무시하므로 중앙 조정 없이 중복을 방지합니다.

이 선택은 직접적인 결과를 가져옵니다. 즉, NATS 코어는 오래 지속되는 대기열 없이 기껏해야 한 번만(최대 한 번) 전달합니다. 이벤트가 통과할 때 ws-front 복제본이 NATS에서 잠시 연결이 끊어지면 따라잡지 못합니다. 채널별 기록(Channel.history, HISTORY_SIZE로 제한되는 메모리 내 버퍼)은 클러스터 수준이 아닌 각 복제본에서 로컬로 유지됩니다. Postgres 논리적 복제 슬롯이 업스트림 내구성 있는 정보 소스로 남아 있기 때문에 이는 허용 가능한 절충안입니다. NATS가 아닌 사용 전에 DB 변경 사항이 손실되지 않도록 보장하는 것은 wal2json입니다.

JetStream은 aura-realtime에 존재하지만 완전히 다른 용도로 사용됩니다. 모든 ws-front복제본 사이에서 누가 어떤 채널에 온라인인지 동기화하는 인스턴스 간 존재 KV 저장소(aura_presence, 메모리 저장소, 60초 max_age)를 제공합니다. 재생할 CDC 이벤트 스트림이 아닌 단기 공유 상태입니다.

#
보안

각 이벤트에서 구독자가 RLS를 다시 확인함

CDC 이벤트는 채널의 모든 구독자에게 그대로 전달되지 않습니다. 폴러는 각 변경 사항을 CDC_NUM_SHARDS 작업자(기본적으로 4개, project_id에 의해 해시됨) 중 하나로 라우팅하고, 먼저 pg_policies를 폴링합니다. RLS 정책이 없는 테이블은 모든 구독자를 허용하며(Supabase 동작), 정책이 있는 테이블은 구독자별 확인을 트리거합니다.

subscriptions.rs::check_rls_visible_batch (실제, 단순화)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- GUC에 주입된 구독자의 JWT 클레임
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

트랜잭션은 구독자의 물리적 Postgres 역할을 맡고 해당 JWT 클레임을 request.jwt.claims(auth.uid() 및 auth.role()이 정책 측에서 읽는 것)에 삽입합니다. 그런 다음 이 역할 아래에 행이 계속 표시되는지 확인한 다음 취소합니다. 쓰기 없음, 실제 RLS 읽기입니다. 이벤트의 allowed_sub_ids에 있는 sub_id 토지만 검증되었습니다. 빈 vec는 페일클로즈 상태에서 아무도 없음을 의미합니다.

아스투스

DELETE는 라인별로 이 검사를 이스케이프합니다. 라인이 사라졌고 RLS에서 다시 읽을 수 없으며 채널의 모든 구독자에게 돌아갑니다. 그 대가로 DELETE의 페이로드는 삭제된 행의 내용이 아닌 ID 열(기본 키)만 노출합니다.

마이그레이션을 방해하는 세부 사항: wal2json는 테이블에 REPLICA IDENTITY FULL이 있는 경우 UPDATE 또는 DELETE의 전체 콘텐츠만 캡처합니다. 이것이 없으면 INSERT만 다시 나타납니다. 이것이 바로 테이블에서 실시간을 활성화하는 엔드포인트가 이 ALTER TABLE를 자동으로 설정하는 이유입니다. 이는 별도의 확인란이 아닌 저장소에 대한 실제 수정 사항입니다.

#
트랩 확인됨

CDC가 Kubernetes에서 침묵을 유지할 수 있는 이유

서류의 이력은 이론적인 사례가 아닌 실제 함정을 기록합니다. Kubernetes 클러스터에서 기본적으로 사용되는 Postgres 이미지(pgvector/pgvector 스톡 커뮤니티 이미지)에는 wal2json플러그인이 포함되어 있지 않습니다. 이것이 없으면 pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json')는 실패하고 클라이언트 측에서는 이를 명시적으로 알리는 신호가 없습니다. WebSocket은 열려 있고 구독이 허용되지만 DB 이벤트는 발생하지 않습니다.

docker/Postgres.Dockerfile (진짜 추출물)dockerfile
RUN apt-get update && \
    apt-get install -y --no-install-recommends \
      postgresql-16-wal2json=2.6-4.pgdg12+1 \
      ...

실제 수정 사항은 이 패키지가 설치된 전용 Postgres 이미지를 빌드 및 게시한 다음 스톡 이미지 대신 Kubernetes StatefulSet를 가리키는 것이었습니다. 또한 wal_level=logical, max_replication_slots 및 max_slot_wal_keep_size을 시작 인수로 명시적으로 설정합니다. 일반 이미지에는 없습니다.

cdc-worker가 전용 /health 엔드포인트에서 cdc_active 플래그(슬롯이 정상으로 확인되지 않은 경우 거짓)를 노출하는 것은 이러한 유형의 오류를 관찰 가능하게 만드는 것입니다. 이는 클라이언트 측 이벤트만 없어 CDC가 작동하지 않는다고 추론하는 것이 아니라 주의 깊게 살펴봐야 할 바이너리 신호입니다.

#
현재 한도

이 파이프라인에서 아직 다루지 않는 것: 전용 클러스터

이 메커니즘은 단일 POSTGRES_REPLICATION_URL변수, 즉 단일 Postgres 서버와 단일 논리 슬롯을 읽습니다. 이는 공유 클러스터의 프로젝트 모델당 다중 테넌트 스키마(project_<uuid>)와 일치합니다. 이 모델에서는 모든 것이 작동합니다. 단일 cdc-worker는 동일한 클러스터에 있는 모든 프로젝트의 변경 사항을 확인하고 다이어그램별로 라우팅합니다.

전용 CNPG 클러스터(자체 격리된 Postgres 인스턴스)의 프로젝트의 경우 오늘 사용된 이미지(CloudNativePG 스톡 이미지에서 구축)는 wal2json가 아닌 pg_graphql만 추가합니다. 어느 것도 전용 클러스터당 cdc-worker를 실행하지 않습니다. 따라서 코드의 현재 상태에서 여기에 설명된 실시간 CDC는 로드맵 기능이 아닌 실제 아키텍처 경계인 공유 계층에 계속 초점을 맞추고 있습니다.

#
클라이언트 측

변경되지 않는 사항: SDK의 postgres_changes API

이 전체 메커니즘은 SDK에서 보이지 않는 상태로 유지됩니다.aurabase-js 체인 가능 API는 이벤트가 NATS를 통해 cdc-worker에서 오는지 아니면 로컬 개발에서 분할되지 않은 aura-realtime 바이너리에서 오는지 여부에 따라 변경되지 않습니다.

app.tstypescript
const channel = aura.realtime.channel('room-1')
  .on('postgres_changes', {
    event: 'INSERT', schema: 'public', table: 'messages'
  }, (payload) => { /* … */ })
  .subscribe()

그러나 단일 이벤트가 발생하기 전에 테이블이 CDC에 등록되어야 합니다. 이는 생성 시 자동으로 수행되지 않습니다. 전용 엔드포인트(또는 Studio의 동등한 토글)에 대한 service_role 호출은 프로젝트 플랫폼 스키마의 realtime_tables에 테이블을 작성하고 REPLICA IDENTITY FULL를 설정합니다.

terminalbash
curl -X PUT \
  https://api.aurabase.cloud/v1/db/<project_id>/schema/tables/messages/realtime \
  -H "apikey: <service_role_key>" \
  -H 'Content-Type: application/json' \
  -d '{"enabled": true}'

나머지 Supabase 호환 SDK 표면(인증, 저장소, 정책 RLS)에 대해서는 마이그레이션 가이드를 참조하세요. aura-realtime가 단일 작업 공간 상자에서 여러 바이너리로 분할되는 방법은 Cargo 아키텍처 세부 정보를 참조하세요.

배포할 준비가 되셨나요?

5분 만에 백엔드를 완성하세요.

신용카드 불필요 · 500MB 무료 · 50,000 MAU