PROD欧洲主权BaaS平台打开仪表板 →

工程 · 10 最小读取值

PostgreSQL CDC 和 NATS:实时架构

Affane Daylami · Fondateur · 2026年7月24日

返回博客

通过 WebSocket 连接的 Aurabase 客户端会在 COMMIT 后几毫秒看到插入到数据库中的一行 — 无需轮询,无需配置 Webhook。该机制分三个阶段。逻辑复制插件将 Postgres WAL 解码为 JSON,只有一个进程有权读取它,NATS 将每个事件中继到集群中的所有 WebSocket 服务器。以下是它在 aura-realtimeservice 的实际代码中的连接方式。

该英文文本是根据法文原文自动生成的,尚未经过审查。
该页面已自动翻译。英文版具有权威性。

这里没有什么是理想化的模式。每个摘录均来自出版之日的 aura-realtime箱。这里有一个明显的区别:NATS JetStream 不是在此管道中传输 CDC 事件的东西。下面我们详细介绍 JetStream 在这里实际执行的操作,以及在其位置上承载 CDC 流量的内容。下面进一步记录了经过验证的 Kubernetes 陷阱,它能够使整个机制保持安静,而不会引发任何错误。

要点
  • Aurabase 的 Postgres CDC 通过 wal2json 和 pg_logical_slot_get_changes 读取 WAL,而不是 pgoutput二进制协议,也不是 Debezium/Kafka Connect 桥。
  • 逻辑槽 aura_cdc_slot 只能容纳一个读者: aura-realtime-cdc-worker 通过 Kubernetes Lease (coordination.k8s.io/v1) 选举出来,具有“始终领导者”模式,不包括 K8s。
  • 到 ws-front 副本的扇出通过 NATS 核心 (aura.realtime.>上的简单发布/订阅),而不是通过持久的 JetStream 流。 JetStream 在同一服务中专门提供跨实例存在 KV 服务。
  • 每个事件都由订阅者在传输之前通过相关线路上的真实 SET LOCAL ROLE 请求在 RLS 中重新检查,而不是缓存的近似值。
  • 检查存储库历史记录中的陷阱:Postgres Kubernetes 映像中没有 wal2json 插件,槽创建失败,并且 CDC 保持静默不活动状态。
#
发动机选择

wal2json,不是 pgoutput,不是 Debezium

大多数 CDC Postgres 管道都会通过 pgoutput(二进制逻辑复制协议),然后通过像 Debezium 这样的连接器将其转换为 Kafka。 Aurabase 跳过此步骤: aura-realtime 服务直接将 WAL 解码为逻辑复制 和 wal2json 插件,这会生成可用的 JSON,而无需中间转换阶段。

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 清除超出此限制),则会自动重新创建槽。假设并记录未同时消耗的事件的丢失。

每次轮询一批 1000 个更改,动态表过滤器每 3 秒刷新一次。只有项目明确启用实时的表才输入 add-tables — 这可以避免解码任何订阅者不感兴趣的表的 WAL。

#
单一制片人

cdc-worker:由 Kubernetes Lease 选举产生,不重复

Postgres 逻辑复制槽一次只能容纳一个活动驱动器。并行运行同一个 aura_cdc_slot 的多个使用者会破坏更改的顺序,而不仅仅是重复它们。 Aurabase 将 aura-realtime 拆分为两个单独的二进制文件,以在不牺牲 WebSocket 水平规模的情况下解决此约束。

Cargo.tomltoml
# cdc-worker:当选领导者,在 NATS 上启动 CDC + 发布。没有 WS 服务器。
[[bin]]
name = "aura-realtime-cdc-worker"

# ws-front:N 个副本(HPA),消耗 NATS,RLS 过滤器,服务 WS/SSE。
[[bin]]
name = "aura-realtime-ws-front"

cdc-worker 仅通过持有 Lease Kubernetes 对象 (coordination.k8s.io/v1) 来获得读取插槽的权限。他尝试创建它,或者在它过期的情况下窃取它,然后在其寿命的三分之一后更新它。当租约丢失(续订失败、持有者变更)时,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, ...
{
  // 更新所有lease_duration_secs / 3
  // 不包括 K8s:工作(令牌)仅一次,令牌从未取消
}
#
真正的扇出

NATS 核心用于 CDC,JetStream 用于存在

这是此类管道中最常被误解的一点:NATS 和 JetStream 是两个独立的事物,并且 CDC 事件不会通过持久的 JetStream 流传递。 cdc-worker 使用 Client::publish_with_headers — NATS 核心发布/订阅 API发布每个事件,该事件最多传递一次 ,无需持久性或重播,与 JetStream 分开。

NATS 心脏(发布/订阅)将 CDC 事件扇出到 ws-front 副本最多一次,没有持久性或重放
NATS JetStream (KV)跨实例存在(aura_presence)共享状态 60 秒,不是要重放的 CDC 流
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.>,将主题解码为原始通道名称。然后,它将事件重新发布到其本地 tokio::sync::broadcast,以用于连接到它的 WebSocket。 Aura-Origin 标头携带发送实例的标识符:每个副本都会忽略它自己发布的消息,这避免了在没有中央协调的情况下出现重复。

这种选择有一个直接的后果:NATS 核心最多交付一次(最多一次),而无需长时间排队。如果事件通过时 ws-front 副本与 NATS 短暂断开连接,则它不会赶上。每个通道的历史记录(Channel.history,由 HISTORY_SIZE界定的内存缓冲区)位于每个副本本地,而不是集群级别。这是一个可以接受的折衷方案,因为 Postgres 逻辑复制槽仍然是上游持久的事实来源。确保在使用之前不会丢失数据库更改的是 wal2json,而不是 NATS。

JetStream 确实存在于 aura-realtime 中 — 但用途完全不同。它提供跨实例存在 KV 存储(aura_presence,内存存储,60 秒 max_age),同步所有 ws-front副本之间谁在哪个通道上在线。短暂的共享状态,而不是要重播的 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";
-- JWT 声称订阅者注入 GUC
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 读取。仅验证 sub_id 落在事件的 allowed_sub_ids 中;空 vec 意味着在失败关闭中没有人。

阿斯图塞

DELETE 每行都会逃避此检查 - 该行已消失,不可能在 RLS 下重新读取它 - 并返回到该通道的所有订阅者。作为回报,DELETE 的有效负载仅公开标识列(主键),而绝不会公开已删除行的内容。

经常妨碍迁移的细节:如果表具有 REPLICA IDENTITY FULL,则 wal2json 仅捕获 UPDATE 或 DELETE 的完整内容。如果没有它,只有 INSERT 会恢复。这就是为什么在表上启用实时的端点会自动设置此 ALTER TABLE — 对存储库的实际修复,而不是单独的复选框。

#
陷阱已验证

为什么CDC可以在Kubernetes中保持沉默

该文件的历史记录了一个真实的陷阱,而不是理论上的案例。 Kubernetes 集群上默认使用的 Postgres 映像(pgvector/pgvector 库存社区映像)不包含 wal2json插件。如果没有它,pg_create_logical_replication_slot('aura_cdc_slot', 'wal2json') 就会失败,并且客户端上没有任何内容明确表示这一点。 WebSocket 保持开放状态,接受订阅,但不会发生任何数据库事件。

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 标志(只要插槽未确认正常,则为 false)。这是一个值得关注的二进制信号,而不是仅根据客户端事件的缺失来推断 CDC 已失效。

#
电流限制

此管道尚未涵盖的内容:专用集群

此机制读取单个 POSTGRES_REPLICATION_URL变量,因此是单个 Postgres 服务器和单个逻辑槽。这与共享集群上的每个项目模型 (project_<uuid>) 的多租户架构一致。在这个模型上,一切正常:单个 cdc-worker 可以看到同一集群中所有项目的更改,并按图表进​​行路由。

对于专用 CNPG 集群(其自己的独立 Postgres 实例)上的项目,今天使用的映像(从 CloudNativePG 库存映像构建)仅添加 pg_graphql,而不是 wal2json。每个专用集群也没有运行 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 如何在单个工作区 crate 中拆分为多个二进制文件,请参阅 Cargo 架构详细信息。

准备好部署了吗?

五分钟内完成您的后端。

无需信用卡 · 500 MB 免费 · 50,000 MAU