PRODSovereign European BaaS platformOpen Dashboard →

Engineering · 10 min read

PostgreSQL CDC with NATS: real-time architecture

Affane Daylami · Fondateur · July 24, 2026

Back to blog

An Aurabase client connected via WebSocket sees a line inserted in the base appear a few milliseconds after the COMMIT — without polling, without webhook to configure. The mechanism is in three stages. A logical replication plugin decodes the Postgres WAL into JSON, only one process has the right to read it, and NATS relays each event to all WebSocket servers in the cluster. Here's how it's wired in the actual code of the aura-realtimeservice.

This English text was generated automatically from the French original and has not been reviewed yet.

Nothing here is an idealized pattern. Each extract comes as is from the aura-realtimecrate, on the date of publication. There is a clear distinction here: NATS JetStream is not what transports CDC events in this pipeline. We detail below what JetStream actually does here, and what carries CDC traffic in its place. A verified Kubernetes trap, capable of making the entire mechanism silent without raising any errors, is documented further below.

The essentials
  • Aurabase's Postgres CDC reads the WAL via wal2json and pg_logical_slot_get_changes — not the pgoutputbinary protocol, not Debezium/Kafka Connect bridge.
  • The logical slot aura_cdc_slot only tolerates one reader: aura-realtime-cdc-worker is elected via a Kubernetes Lease (coordination.k8s.io/v1), with an "always leader" mode excluding K8s.
  • The fan-out to the ws-front replicas goes through NATS core (simple pub/sub on aura.realtime.>), not through a persistent JetStream stream. JetStream, in this same service, exclusively serves cross-instance presence KV.
  • Each event is rechecked in RLS by subscriber, just before transmission, via a real SET LOCAL ROLE request on the line concerned — not a cached approximation.
  • Checked pitfall in repository history: without the wal2json plugin in the Postgres Kubernetes image, the slot creation fails and the CDC remains silently inactive.
#
Engine choice

wal2json, not pgoutput, not Debezium

Most CDC Postgres pipelines go through pgoutput, the binary logic replication protocol, and then through a connector like Debezium which translates it to Kafka. Aurabase skips this step: the aura-realtime service directly decodes the WAL into logical replication with the wal2json plugin, which produces usable JSON without an intermediate translation stage.

cdc/postgres.rs (real queries)sql
-- Created once if absent, recreated if WAL expired
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Repeated poll (backoff 100ms when the batch returns empty)
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');

Three Postgres prerequisites, verified in the code: the application role must carry REPLICATION, wal_level = logical must be active, and max_slot_wal_keep_size must limit WAL retention. Without this last limit, a slow consumer causes the disk volume to grow indefinitely. The service also monitors wal_status of the slot: if it changes to lost (WAL purged beyond this limit), the slot is automatically recreated. The loss of events not consumed in the meantime is assumed and logged.

A batch of 1000 changes per poll, with a dynamic table filter refreshed every 3 seconds. Only tables where a project has explicitly enabled real time enter add-tables — this avoids decoding the WAL of tables that are of no interest to any subscribers.

#
A single producer

cdc-worker: elected by a Kubernetes Lease, not duplicated

A Postgres logical replication slot tolerates only one active drive at a time. Running multiple consumers of the same aura_cdc_slot in parallel would break the order of changes, not just duplicate them. Aurabase split aura-realtime into two separate binaries to resolve this constraint without sacrificing the horizontal scale of the WebSocket.

Cargo.tomltoml
# cdc-worker: elected leader, launches the CDC + publishes on NATS. No WS server.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N replicas (HPA), consumes NATS, RLS filter, serves WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker only gains the right to read the slot by holding a Lease Kubernetes object (coordination.k8s.io/v1). He attempts to create it, or steal it if it has expired, then renews it after a third of its lifespan. When the lease is lost (renew failed, holder changed), a CancellationToken cuts the work in progress and the process returns to an acquisition loop. Outside of Kubernetes — during local testing or non-cluster development, for example — the K8s client fails to connect, and the service switches to “always leading” mode. Useful in dev, incorrect if you forget it in prod with several replicas.

cdc/leader.rs (real signatures)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // renews all lease_duration_secs / 3
  // excluding K8s: work(token) only once, token never canceled
}
#
Real fan-out

NATS core for CDC, JetStream for presence

This is the most often misunderstood point in this kind of pipeline: NATS and JetStream are two separate things, and CDC events do not pass through a persistent JetStream stream. cdc-worker publishes each event with Client::publish_with_headers — the NATS core pub/sub API, the one delivered at-most-once without persistence or replay, separate from JetStream.

NATS heart (pub/sub)Fan-out of CDC events to ws-front replicasAt-most-once, without persistence or replay
NATS JetStream (KV)Cross-instance presence (aura_presence)Shared state 60s, not a CDC stream to replay
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker publishes (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front listens to the entire subtree (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Each replica ws-front subscribes to subtree aura.realtime.>, decodes the topic to the original channel name. It then republishes the event to its local tokio::sync::broadcast, for the WebSockets connected to it. A Aura-Origin header carries the identifier of the sending instance: each replica ignores the messages that it has itself published, which avoids duplicates without central coordination.

This choice has a direct consequence: the NATS core delivers at best once (at-most-once), without a long-lasting queue. If an ws-front replica is briefly disconnected from NATS when an event passes, it does not catch up. The per-channel history (Channel.history, an in-memory buffer bounded by HISTORY_SIZE) lives locally at each replica, not at the cluster level. This is an acceptable compromise precisely because the Postgres logical replication slot remains the upstream durable source of truth. It's wal2json that ensures no DB changes are lost before consumption, not NATS.

JetStream does exist in aura-realtime — but for an entirely different use. It serves the cross-instance presence KV store (aura_presence, memory store, 60-second max_age), which synchronizes who is online on which channel between all ws-frontreplicas. A short-lived shared state, not a stream of CDC events to replay.

#
Security

RLS rechecked by subscriber, at each event

A CDC event does not go raw to all subscribers of a channel. The poller routes each change to one of CDC_NUM_SHARDS workers (4 by default, hashed by project_id), which first polls pg_policies. A table without an RLS policy allows all subscribers — the Supabase behavior — while a table with policies triggers a per-subscriber check.

subscriptions.rs::check_rls_visible_batch (real, simplified)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- JWT claims of the subscriber injected into GUC
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

The transaction takes on the physical Postgres role of the subscriber and injects its JWT claims into request.jwt.claims — the one that auth.uid() and auth.role() read on the policy side. It then checks that the line remains visible under this role, then cancels: no writing, a real RLS read. Only validated sub_id land in allowed_sub_ids of the event; an empty vec means no one, in fail-closed.

Astuce

A DELETE escapes this check per line - the line has disappeared, it is impossible to reread it under RLS - and goes back to all subscribers of the channel. In return, the payload of a DELETE only ever exposes the identity columns (primary key), never the content of the deleted row.

A detail that often gets in the way of migration: wal2json only captures the complete content of a UPDATE or a DELETE if the table has REPLICA IDENTITY FULL. Without it, only the INSERT come back up. This is why the endpoint that enables real-time on a table sets this ALTER TABLE automatically — an actual fix to the repository, not a separate checkbox.

#
Trap verified

Why CDC can remain silent in Kubernetes

The filing's history documents a real trap, not a theoretical case. The Postgres image used by default on a Kubernetes cluster (the pgvector/pgvector stock community image) does not include the wal2jsonplugin. Without it, pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') fails, and nothing on the client side explicitly signals this. WebSockets stay open, subscriptions are accepted, but no DB events ever happen.

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

The actual fix was to build and publish a dedicated Postgres image with this package installed, then point the Kubernetes StatefulSet at it instead of the stock image. It also explicitly set wal_level=logical, max_replication_slots and max_slot_wal_keep_size as startup arguments — absent from generic images.

It is precisely to make this type of failure observable that cdc-worker exposes a cdc_active flag (false as long as the slot is not confirmed healthy) on its dedicated /health endpoint. This is a binary signal to watch for — rather than inferring dead CDC from the absence of client-side events alone.

#
Current limit

What this pipeline does not yet cover: dedicated clusters

This mechanism reads a single POSTGRES_REPLICATION_URLvariable, therefore a single Postgres server and a single logical slot. This is consistent with the multi-tenant schema per project model (project_<uuid>) on a shared cluster. On this model, everything works: a single cdc-worker sees the changes of all the projects in the same cluster, and route by diagram.

For a project on a dedicated CNPG cluster — its own isolated Postgres instance — the image used today (built from the CloudNativePG stock image) only adds pg_graphql, not wal2json. Nothing, either, runs a cdc-worker per dedicated cluster. In the current state of the code, the real-time CDC described here therefore remains focused on the shared tier — a real architectural boundary, not a roadmap functionality.

#
Client side

What doesn't change: the SDK's postgres_changes API

This entire mechanism remains invisible from the SDK. Theaurabase-js chainable API does not change depending on whether the event comes from cdc-worker via NATS or, in local dev, from the unsplit aura-realtime binary.

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

But before a single event can happen, the table must be registered with the CDC — this is not automatic at creation. A service_role call on the dedicated endpoint (or the equivalent toggle in Studio) writes the table in realtime_tables of the project's platform schema, and sets REPLICA IDENTITY FULL for you.

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

For the rest of the Supabase compatible SDK surface — auth, storage, policies RLS —, see the migration guide. For how aura-realtime splits into multiple binaries in a single workspace crate, see the Cargo architecture detail.

READY TO DEPLOY?

Your backend in five minutes.

No credit card required · 500 MB free · 50,000 MAU