跳到主要内容

Postgres CDC

监听 PostgreSQL 表变更。presence / postgres_changes 监听必须在 subscribe() 之前绑定。

前置条件

本页代码能跑通的前提是业务库已配置好 数据库侧准备(publication、connector 角色授权、RLS 策略)。这些没做的话,订阅会成功但收不到任何事件,容易误判成 SDK 问题。

const channel = realtime.channel("db-changes");

channel.on("postgres_changes", { event: "*", schema: "public" }, (payload) => {
console.log("public schema 变更", payload);
});

channel.on(
"postgres_changes",
{ event: "INSERT", schema: "public", table: "messages" },
(payload) => {
console.log("新增", payload.new);
}
);

channel.on(
"postgres_changes",
{
event: "UPDATE",
schema: "public",
table: "users",
filter: "username=eq.Realtime",
},
(payload) => {
console.log("更新", payload.new, payload.old);
}
);

channel.subscribe((status, err) => {
if (status === "SUBSCRIBED") {
console.log("开始接收数据库变更");
}
if (status === "CHANNEL_ERROR") {
console.error(err);
}
});

filter 可以是原始字符串,也可以用 postgresChangesFilter() 构建,二者线格式相同:

import { postgresChangesFilter } from "@cloudbase/js-sdk/realtime-js";

// 字符串
{ event: "UPDATE", schema: "public", table: "users", filter: "id=eq.1" }

// Builder
{
event: "UPDATE",
schema: "public",
table: "users",
filter: postgresChangesFilter().eq("id", 1),
}
运算符字符串Builder含义
eqid=eq.1.eq("id", 1)等于
neqid=neq.1.neq("id", 1)不等于
lt / lte / gt / gteage=gte.18.gte("age", 18)比较
instatus=in.(active,pending).in("status", ["active", "pending"])属于列表
like / iliketitle=like.%foo%.like("title", "%foo%")模式匹配
isdeleted_at=is.null.is("deleted_at", null)IS null/true/false
match / imatchtitle=match.^foo.match("title", "^foo")POSIX 正则
isdistinctvalue=isdistinct.1.isDistinct("value", 1)NULL 安全不等于

取反:字符串加 not. 前缀,或使用 .not(column, operator, value)。多个条件用逗号(AND)连接,Builder 则链式调用。

channel.on(
"postgres_changes",
{
event: "UPDATE",
schema: "public",
table: "orders",
filter: postgresChangesFilter()
.gt("amount", 100)
.not("status", "in", ["draft", "archived"]),
},
(payload) => console.log(payload)
);

用 select 只接收部分列,减小载荷:

channel.on(
"postgres_changes",
{
event: "*",
schema: "public",
table: "users",
select: ["id", "first_name"],
},
(payload) => {
// payload.new 仅含 { id, first_name }
console.log(payload);
}
);

需要等服务端确认 CDC 订阅就绪后再触发 SUBSCRIBED 时,设置 postgres_changes_options.wait:

const channel = realtime.channel("db-changes", {
config: {
postgres_changes_options: { wait: true, timeout: 15000 },
},
});
提示

Realtime 在服务端按单表 WAL 过滤,没有资源嵌套(!inner)和 or() 分组。like / ilike 通配符使用 % 而不是 *。

数据库侧准备​

postgres_changes 走的是业务库的逻辑复制。下面以 public.todos 表为例说明如何开启、查看和删除表监听。

名词约定​

名词说明
publicationPG 逻辑复制的发布对象,realtime 通过它决定监听哪些表,本体系默认名 cloudbase_realtime(名称由租户配置决定,若你的环境指定了其他名称以配置为准)
connector 角色每个库一个,cloudbase_realtime_admin,由服务端创建,负责读 WAL / 跑 poller
订阅角色客户端 JWT 里的角色:anon / authenticated / service_role,RLS 检查以它为准

1. 开启一张表的监听​

最小操作两条:

-- 1. 把表加入 publication(publication 已存在时)
ALTER PUBLICATION cloudbase_realtime ADD TABLE public.todos;

-- 2. 给 connector 角色授权(通常不会自动有,缺了 poller 读不到变更)
GRANT SELECT ON public.todos TO "cloudbase_realtime_admin";

首次开启(publication 尚不存在)​

CREATE PUBLICATION 没有 IF NOT EXISTS,publication 不存在时用:

CREATE PUBLICATION cloudbase_realtime FOR TABLE public.todos;

幂等脚本​

不区分首次 / 增量,可重复执行:

DO $$
BEGIN
IF NOT EXISTS (SELECT 1 FROM pg_publication WHERE pubname = 'cloudbase_realtime') THEN
CREATE PUBLICATION cloudbase_realtime;
END IF;

IF NOT EXISTS (
SELECT 1 FROM pg_publication_tables
WHERE pubname = 'cloudbase_realtime'
AND schemaname = 'public'
AND tablename = 'todos'
) THEN
EXECUTE 'ALTER PUBLICATION cloudbase_realtime ADD TABLE public.todos';
END IF;
END $$;

-- GRANT / REPLICA IDENTITY 天然幂等,直接重复执行
GRANT SELECT ON public.todos TO "cloudbase_realtime_admin";
ALTER TABLE public.todos REPLICA IDENTITY FULL;

2. 业务表授权 + 开启 RLS​

变更事件是按业务表自己的 RLS 策略逐订阅者校验的:

-- 让角色有读权限
GRANT SELECT ON public.todos TO anon, authenticated;

-- 开启行级安全
ALTER TABLE public.todos ENABLE ROW LEVEL SECURITY;

-- 只让用户收到自己的数据
CREATE POLICY todos_select ON public.todos
FOR SELECT TO authenticated
USING (owner_id = auth.uid());

3. 查看当前监听​

查 publication 里开启了哪些表:

SELECT schemaname, tablename
FROM pg_publication_tables
WHERE pubname = 'cloudbase_realtime'
ORDER BY 1, 2;

查当前活跃的订阅(realtime 视角):

SELECT entity, filters, created_at
FROM realtime.subscription
ORDER BY created_at DESC;

4. 删除监听​

移除一张表(常用):

ALTER PUBLICATION cloudbase_realtime DROP TABLE public.todos;

移除后 poller 会在下一个周期感知到,自动停止该表的变更推送;已有客户端订阅不会立刻断开,但不再收到该表事件。

彻底关闭所有监听:

DROP PUBLICATION cloudbase_realtime;
谨慎操作

DROP PUBLICATION 会移除 publication 内全部表,需要重新开启时按首次开启流程重建。不要手动删除复制槽(cloudbase_realtime_replication_slot 等),复制槽由服务端管理。

常见问题​

现象排查
订阅成功但收不到事件① 表是否在 publication 里(见查看当前监听);② connector 角色是否有 SELECT;③ RLS 策略是否放行订阅角色
UPDATE / DELETE 的 old_record 为空表没设 REPLICA IDENTITY FULL
客户端订阅被拒订阅角色(anon / authenticated)对表无 SELECT,补 GRANT
加表后约 10 秒才生效正常,poller 周期性感知 publication 变化,无需重启