PRODPlataforma BaaS soberana europeaAbrir panel →

Ingeniería · 10 lectura mínima

PostgreSQL CDC con NATS: arquitectura en tiempo real

Affane Daylami · Fondateur · 24 de julio de 2026

volver al blog

Un cliente de Aurabase conectado a través de WebSocket ve aparecer una línea insertada en la base unos milisegundos después del COMMIT, sin sondeo, sin webhook que configurar. El mecanismo consta de tres etapas. Un complemento de replicación lógica decodifica el WAL de Postgres en JSON, solo un proceso tiene derecho a leerlo y NATS transmite cada evento a todos los servidores WebSocket del clúster. Así es como está conectado en el código real del servicio aura-realtime.

Este texto en inglés se generó automáticamente a partir del original en francés y aún no ha sido revisado.
Esta página fue traducida automáticamente. La versión en inglés es autorizada.

Nada aquí es un patrón idealizado. Cada extracto viene tal cual de la caja aura-realtime, en la fecha de publicación. Aquí hay una distinción clara: NATS JetStream no es lo que transporta los eventos de los CDC en este canal. A continuación detallamos qué hace realmente JetStream aquí y qué transporta el tráfico de los CDC en su lugar. A continuación se documenta una trampa de Kubernetes verificada, capaz de silenciar todo el mecanismo sin generar ningún error.

Lo esencial
  • El CDC de Postgres de Aurabase lee el WAL a través de wal2json y pg_logical_slot_get_changes, no el protocolo binario pgoutput, ni el puente Debezium/Kafka Connect.
  • La ranura lógica aura_cdc_slot solo tolera un lector: aura-realtime-cdc-worker se elige mediante un arrendamiento de Kubernetes (coordination.k8s.io/v1), con un modo "siempre líder" que excluye los K8.
  • La distribución de las réplicas ws-front pasa a través del núcleo NATS (pub/sub simple en aura.realtime.>), no a través de una secuencia JetStream persistente. JetStream, en este mismo servicio, brinda exclusivamente presencia entre instancias KV.
  • El suscriptor vuelve a verificar cada evento en RLS, justo antes de la transmisión, a través de una solicitud SET LOCAL ROLE real en la línea en cuestión, no una aproximación almacenada en caché.
  • Error comprobado en el historial del repositorio: sin el complemento wal2json en la imagen de Postgres Kubernetes, la creación de la ranura falla y el CDC permanece silenciosamente inactivo.
#
Elección del motor

wal2json, no pgoutput, no Debezium

La mayoría de las canalizaciones de CDC Postgres pasan por pgoutput, el protocolo de replicación de lógica binaria, y luego por un conector como Debezium que lo traduce a Kafka. Aurabase omite este paso: el servicio aura-realtime decodifica directamente el WAL en una replicación lógica con el complemento wal2json, que produce JSON utilizable sin una etapa de traducción intermedia.

cdc/postgres.rs (consultas reales)sql
-- Creado una vez si está ausente, recreado si WAL expiró
SELECT pg_create_logical_replication_slot(
  'aura_cdc_slot', 'wal2json');

-- Encuesta repetida (retroceso de 100 ms cuando el lote regresa vacío)
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');

Tres requisitos previos de Postgres, verificados en el código: la función de la aplicación debe llevar REPLICATION, wal_level = logical debe estar activa y max_slot_wal_keep_size debe limitar la retención de WAL. Sin este último límite, un consumidor lento hace que el volumen del disco crezca indefinidamente. El servicio también monitorea wal_status de la ranura: si cambia a lost (WAL purgado más allá de este límite), la ranura se recrea automáticamente. La pérdida de eventos no consumidos mientras tanto se asume y se registra.

Un lote de 1000 cambios por encuesta, con un filtro de tabla dinámico que se actualiza cada 3 segundos. Solo las tablas en las que un proyecto ha habilitado explícitamente el tiempo real ingresan a add-tables; esto evita decodificar el WAL de tablas que no son de interés para ningún suscriptor.

#
un solo productor

cdc-worker: elegido mediante un contrato de arrendamiento de Kubernetes, no duplicado

Una ranura de replicación lógica de Postgres tolera sólo una unidad activa a la vez. Ejecutar varios consumidores del mismo aura_cdc_slot en paralelo rompería el orden de los cambios, no solo los duplicaría. Aurabase dividió aura-realtime en dos binarios separados para resolver esta restricción sin sacrificar la escala horizontal de WebSocket.

Cargo.tomltoml
# cdc-worker: líder electo, lanza CDC + publica en NATS. Sin servidor WS.
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front: N réplicas (HPA), consume NATS, filtro RLS, sirve WS/SSE.
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker solo obtiene el derecho de leer la ranura manteniendo un objeto de arrendamiento de Kubernetes (coordination.k8s.io/v1). Intenta crearlo o robarlo si ha caducado y luego lo renueva después de un tercio de su vida útil. Cuando se pierde el contrato de arrendamiento (falla la renovación, se cambia de titular), un CancellationToken corta el trabajo en curso y el proceso vuelve a un ciclo de adquisición. Fuera de Kubernetes (durante las pruebas locales o el desarrollo fuera del clúster, por ejemplo), el cliente K8s no se conecta y el servicio cambia al modo "siempre líder". Útil en desarrollo, incorrecto si lo olvidas en producción con varias réplicas.

cdc/leader.rs (firmas reales)rust
pub async fn run_with_lease<F, Fut>(cfg: LeaseConfig, work: F) -> anyhow::Result<()>
where F: Fn(CancellationToken) -> Fut, ...
{
  // renueva todos los arrendamientos_duración_segundos / 3
  // excluyendo K8: trabajo (token) solo una vez, el token nunca se cancela
}
#
Fan-out real

Núcleo NATS para CDC, JetStream para presencia

Este es el punto que más a menudo se malinterpreta en este tipo de canalización: NATS y JetStream son dos cosas separadas, y los eventos CDC no pasan a través de una transmisión JetStream persistente. cdc-worker publica cada evento con Client::publish_with_headers — la API principal de publicación/subscripción de NATS, la que se entrega como máximo una vez sin persistencia ni repetición, separada de JetStream.

Corazón NATS (pub/sub)Distribución de eventos de CDC en réplicas de ws-frontComo máximo una vez, sin persistencia ni repetición
NATS JetStream (KV)Presencia entre instancias (aura_presence)Estado compartido 60, no una transmisión de CDC para reproducir
channels/manager.rs → nats/mod.rs (extraits réels)rust
// cdc-worker publica (ChannelManager::publish_and_relay)
nc.publish_with_headers(subject, headers, payload).await

// ws-front escucha todo el subárbol (NatsRelay::start_relay)
client.subscribe("aura.realtime.>").await

Cada réplica ws-front se suscribe al subárbol aura.realtime.>y decodifica el tema al nombre del canal original. Luego vuelve a publicar el evento en su tokio::sync::broadcastlocal, para los WebSockets conectados a él. Un encabezado Aura-Origin lleva el identificador de la instancia emisora: cada réplica ignora los mensajes que ella misma ha publicado, lo que evita duplicados sin coordinación central.

Esta elección tiene una consecuencia directa: el núcleo NATS entrega, en el mejor de los casos, una vez (como máximo una vez), sin una cola de larga duración. Si una réplica ws-front se desconecta brevemente de NATS cuando pasa un evento, no se pone al día. El historial por canal (Channel.history, un búfer en memoria delimitado por HISTORY_SIZE) vive localmente en cada réplica, no en el nivel del clúster. Este es un compromiso aceptable precisamente porque la ranura de replicación lógica de Postgres sigue siendo la fuente duradera de verdad. Es wal2json el que garantiza que no se pierdan cambios de base de datos antes del consumo, no NATS.

JetStream existe en aura-realtime, pero para un uso completamente diferente. Sirve el almacén KV de presencia entre instancias (aura_presence, almacén de memoria, max_agede 60 segundos ), que sincroniza quién está en línea en qué canal entre todas las réplicas ws-front. Un estado compartido de corta duración, no un flujo de eventos de los CDC para reproducir.

#
Seguridad

RLS revisado nuevamente por el suscriptor, en cada evento

Un evento de CDC no se transmite directamente a todos los suscriptores de un canal. El sondeador enruta cada cambio a uno de los trabajadores CDC_NUM_SHARDS (4 de forma predeterminada, codificados por project_id), que primero sondea pg_policies. Una tabla sin una política RLS permite a todos los suscriptores (el comportamiento de Supabase), mientras que una tabla con políticas activa una verificación por suscriptor.

subscriptions.rs::check_rls_visible_batch (real, simplificado)sql
BEGIN;
SET LOCAL ROLE "aura_authenticated";
-- Reclamaciones de JWT del suscriptor inyectadas en GUC
SELECT set_config('request.jwt.claims', $1, true);
SELECT EXISTS(SELECT 1 FROM "schema"."table" WHERE id::text = $2);
ROLLBACK;

La transacción asume el rol físico de Postgres del suscriptor e inyecta sus reclamos JWT en request.jwt.claims, el que auth.uid() y auth.role() leen en el lado de la política. Luego verifica que la línea permanezca visible bajo este rol y luego cancela: sin escritura, una lectura RLS real. Sólo los sub_id validados aterrizan en el allowed_sub_ids del evento; un vec vacío significa nadie, en fall-closed.

Astucia

Un DELETE escapa a esta verificación por línea (la línea ha desaparecido, es imposible volver a leerla en RLS) y vuelve a todos los suscriptores del canal. A cambio, la carga útil de DELETE solo expone las columnas de identidad (clave principal), nunca el contenido de la fila eliminada.

Un detalle que a menudo obstaculiza la migración: wal2json solo captura el contenido completo de un UPDATE o un DELETE si la tabla tiene REPLICA IDENTITY FULL. Sin él, solo vuelve a aparecer el INSERT. Es por eso que el punto final que habilita el tiempo real en una tabla establece este ALTER TABLE automáticamente: una solución real para el repositorio, no una casilla de verificación separada.

#
Trampa verificada

Por qué los CDC pueden permanecer en silencio en Kubernetes

El historial de la presentación documenta una trampa real, no un caso teórico. La imagen de Postgres utilizada de forma predeterminada en un clúster de Kubernetes (la imagen de la comunidad de stock pgvector/pgvector) no incluye el complemento wal2json. Sin él, pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') falla y nada en el lado del cliente lo indica explícitamente. Los WebSockets permanecen abiertos, se aceptan suscripciones, pero nunca ocurre ningún evento de base de datos.

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

La solución real fue crear y publicar una imagen de Postgres dedicada con este paquete instalado y luego apuntar Kubernetes StatefulSet a ella en lugar de a la imagen de archivo. También establece explícitamente wal_level=logical, max_replication_slots y max_slot_wal_keep_size como argumentos de inicio, ausentes en las imágenes genéricas.

Es precisamente para hacer observable este tipo de falla que cdc-worker expone un indicador cdc_active (falso siempre que no se confirme que la ranura está en buen estado) en su punto final dedicado /health. Esta es una señal binaria a la que hay que prestar atención, en lugar de inferir CDC muertos únicamente por la ausencia de eventos del lado del cliente.

#
Límite actual

Lo que este pipeline aún no cubre: clusters dedicados

Este mecanismo lee una única variable POSTGRES_REPLICATION_URL, por lo tanto, un único servidor Postgres y una única ranura lógica. Esto es coherente con el esquema multiinquilino por modelo de proyecto (project_<uuid>) en un clúster compartido. En este modelo, todo funciona: un solo cdc-worker ve los cambios de todos los proyectos en el mismo clúster y los ruta por diagrama.

Para un proyecto en un clúster CNPG dedicado (su propia instancia aislada de Postgres), la imagen utilizada hoy (creada a partir de la imagen de archivo de CloudNativePG) solo agrega pg_graphql, no wal2json. Nada ejecuta tampoco un cdc-worker por clúster dedicado. Por lo tanto, en el estado actual del código, el CDC en tiempo real que se describe aquí permanece centrado en el nivel compartido: un límite arquitectónico real, no una funcionalidad de hoja de ruta.

#
Lado del cliente

Lo que no cambia: la API postgres_changes del SDK

Todo este mecanismo permanece invisible desde el SDK. La API encadenableaurabase-js no cambia dependiendo de si el evento proviene de cdc-worker a través de NATS o, en desarrollo local, del binario aura-realtime no dividido.

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

Pero antes de que pueda ocurrir un solo evento, la tabla debe registrarse en el CDC; esto no es automático en el momento de la creación. Una llamada service_role en el punto final dedicado (o el conmutador equivalente en Studio) escribe la tabla en realtime_tables del esquema de plataforma del proyecto y configura REPLICA IDENTITY FULL por usted.

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

Para el resto de la superficie SDK compatible con Supabase (autenticación, almacenamiento, políticas RLS), consulte la guía de migración . Para saber cómo aura-realtime se divide en varios archivos binarios en una única caja de espacio de trabajo, consulte el detalle de la arquitectura de carga .

¿LISTO PARA IMPLEMENTAR?

Tu backend en cinco minutos.

No se requiere tarjeta de crédito · 500 MB gratis · 50,000 MAU