云数据库 PostgreSQL 支持 DuckDB 加速功能,本文为您介绍 DuckDB 加速功能的使用指引。
支持版本
v18.4_r1.11及以上版本。
前置参数配置
在 postgresql.conf 中设置以下参数,修改完成后需重启实例。
参数 | 建议值 | 说明 |
shared_preload_libraries | 包含 pg_extension_base | 加载 lake 引擎与 duck server |
wal_level | logical | 必须设置为 logical,replica 级别不足以支撑同步表所需的逻辑复制;若未开启,force_sync 等同步操作将报错 |
max_replication_slots | ≥ 16 | 同步表共享复制槽池(如 sync_pool_slot_5_0),不是每张表独立一个;实际占用数量远小于同步表数量,预留10 ~ 16通常足够 |
max_wal_senders | ≥ max_replication_slots + 8 | 与复制槽配套,给同步引擎和其他逻辑复制/主从留出余量 |
max_worker_processes | ≥ 32 | 1个 database 共用1个 apply worker,并需为全局 sync manager 预留资源 |
配置示例:
shared_preload_libraries = 'pg_extension_base'wal_level = logicalmax_replication_slots = 64max_wal_senders = 64max_worker_processes = 64
检查当前值
输入:
SELECT name, settingFROM pg_settingsWHERE name IN ('shared_preload_libraries','wal_level','max_replication_slots','max_wal_senders','max_worker_processes');
输出:
name | setting | context--------------------------------+----------------------------------------------------------------------------------------------------------------------------------------------+------------max_replication_slots | 10 | postmastermax_wal_senders | 27 | postmastermax_worker_processes | 8 | postmastershared_preload_libraries | pg_stat_statements,pg_stat_log,wal2json,decoderbufs,decoder_raw,rds_server_handler,tencentdb_pwdcheck,auto_explain,pgaudit,pg_extension_base | postmasterwal_level | logical | postmaster(5 rows)
其他约束:
不支持 TEMP 表、UNLOGGED 表以及分区子表;分区表须对父表调用 create_sync。
建议同步表具备主键;若无主键且存在 UPDATE / DELETE 场景,须提前设置 REPLICA IDENTITY FULL。
查看当前值
输入:
SELECT name, setting, contextFROM pg_settingsWHERE name IN ('pg_lake_engine.enable_duckdb_features','pg_lake_engine.duckdb_replication','pg_lake.enabled','pg_lake_table.enable_full_query_pushdown','pg_lake_table.enable_heap_query_pushdown','pg_lake_table.enable_heap_filter_pushdown','lake_sync.routing_debug')ORDER BY name;
输出:
name | setting | context--------------------------------------------+---------+------------lake_sync.routing_debug | off | userpg_lake.enabled | on | sighuppg_lake_engine.duckdb_replication | on | sighuppg_lake_engine.enable_duckdb_features | on | postmasterpg_lake_table.enable_full_query_pushdown | on | userpg_lake_table.enable_heap_filter_pushdown | on | userpg_lake_table.enable_heap_query_pushdown | on | user(7 rows)
创建插件
输入:
CREATE EXTENSION pg_lake CASCADE;
输出:
NOTICE: installing required extension "pg_lake_table"NOTICE: installing required extension "pg_lake_engine"NOTICE: installing required extension "pg_extension_base"NOTICE: installing required extension "pg_map"NOTICE: installing required extension "pg_lake_iceberg"NOTICE: installing required extension "btree_gist"NOTICE: installing required extension "pg_lake_copy"NOTICE: installing required extension "pg_lake_sync"CREATE EXTENSION
CASCADE 会自动安装以下依赖扩展:
pg_lake_engine — duck_server 管理
pg_lake_table — 列存表访问方法(USING duckdb)
pg_lake_sync — 同步表引擎(USING sync / lake_sync.*)
验证安装
输入:
SELECT extname, extversion FROM pg_extension WHERE extname LIKE 'pg_lake%';
输出:
extname | extversion-------------------+------------pg_extension_base | 3.3pg_lake | 3.3pg_lake_copy | 3.3pg_lake_engine | 3.5pg_lake_iceberg | 3.4pg_lake_sync | 3.2pg_lake_table | 3.4(7 rows)
查看 pg_lake 版本(比查询 pg_extension 更直接)
输入:
SELECT lake.version();
输出:
version----------------main (2d50ce2)(1 row)
如何卸载 pg_lake
可通过级联 DROP 一次卸载所有 pg_lake 相关扩展:
输入(在测试库执行):
DROP EXTENSION pg_lake_engine CASCADE;SELECT extname FROM pg_extension WHERE extname LIKE 'pg_lake%' OR extname LIKE 'pg_extension%' ORDER BY extname;
输出:
NOTICE: drop cascades to 5 other objectsDETAIL: drop cascades to extension pg_lake_icebergdrop cascades to extension pg_lake_tabledrop cascades to extension pg_lake_copydrop cascades to extension pg_lake_syncdrop cascades to extension pg_lakeextname-------------------pg_extension_base(1 row)
创建两种加速表
云数据库 PostgreSQL 提供两种基于 DuckDB 的加速表类型,应根据是否需要 OLTP 事务能力进行选择:
类型 | 底层 | 事务 / 写入模型 | 适用场景 |
同步表 USING sync | Heap(行存) + DuckDB 镜像 | 事务写入 heap,分析查询走镜像 | HTAP:OLTP 与分析并存 |
列存表 USING duckdb | 纯 DuckDB foreign table | 支持 INSERT / UPDATE / DELETE / COPY,但不具备行级事务隔离与索引 | 纯分析场景,数据以批为单位持续追加 |
同步表
同步表 = 一张 heap 表 + 一张 DuckDB 镜像表,写入走 heap,分析走镜像,逻辑复制异步同步。
同步表提供两种创建方式:新建即为同步表,或将已有堆表转换为同步表。
方式一:直接新建同步表
输入:
CREATE TABLE orders (o_id bigint PRIMARY KEY,o_customer text,o_amount numeric(12,2),o_ts timestamptz DEFAULT now()) USING sync;
输出:
CREATE TABLE
heap 侧确认。
输入:
SELECT n.nspname AS schema, c.relname, a.amname AS access_methodFROM pg_class cJOIN pg_namespace n ON n.oid = c.relnamespaceLEFT JOIN pg_am a ON a.oid = c.relamWHERE c.relname = 'orders';
输出:
schema | relname | access_method--------+---------+---------------public | orders | heap(1 row)
系统自动生成:
heap 表 orders (业务操作对象)。
镜像表 orders_duckdb_sync (在 PG 中可见,业务查询无需直接引用)。
方式二:将已有堆表转换为同步表
lake_sync.create_sync(regclass, mirror_storage text DEFAULT 'duckdb') 支持两种调用形式,二选一即可:输入 (默认镜像存储类型即 duckdb,可省略第二参数):
-- customers 为普通堆表SELECT lake_sync.create_sync('customers'::regclass);
输出 (函数返回类型为 void,值列显示为空):
create_sync-------------(1 row)
或显式指定镜像存储类型 (效果等价):
SELECT lake_sync.create_sync('customers'::regclass, 'duckdb');
说明:
同一张表不能被重复 create_sync。若目标 heap 表已存在对应的 _duckdb_sync 镜像,再次调用会报错:
ERROR: mirror table public.customers_duckdb_sync already exists
若需要重建镜像,请先 lake_sync.reset_sync(oid) 释放当前同步关系,或用 lake_sync.rebuild_sync_tables('duckdb') 全库重建。
转换后校验:heap 侧仍是 heap 表,镜像表已生成。
输入:
SELECT n.nspname AS schema, c.relname, a.amname AS access_methodFROM pg_class cJOIN pg_namespace n ON n.oid = c.relnamespaceLEFT JOIN pg_am a ON a.oid = c.relamWHERE c.relname = 'customers';
输出:
schema | relname | access_method--------+-----------+---------------public | customers | heap(1 row)
输入:
SELECT heap_table_oid::regclass AS heap_table,iceberg_table_oid::regclass AS mirror_table,sync_state, mirror_storage_typeFROM lake_sync.sync_tablesORDER BY heap_table_oid::regclass::text;
输出:
heap_table | mirror_table | sync_state | mirror_storage_type------------+-----------------------+------------+---------------------customers | customers_duckdb_sync | ACTIVE | duckdborders | orders_duckdb_sync | ACTIVE | duckdb(2 rows)
create_sync 执行以下5项操作:
1. 创建 DuckDB 镜像表 customers_duckdb_sync。
2. 按需设置 REPLICA IDENTITY。
3. 锁定 heap 表并全量导入现有数据 (数据量大时存在阻塞窗口,建议在业务低峰期执行)。
4. 将表加入共享逻辑复制槽与 publication。
5. 状态迁移至 ACTIVE,进入增量同步阶段。
创建完成后的常规操作
强制镜像追赶至最新 LSN
输入:
SELECT lake_sync.force_sync('orders'::regclass);
输出:
force_sync------------t(1 row)
开启自动路由 (默认关闭。开关只影响通过 heap 表名的 SELECT 是否被自动路由到镜像;直接引用 xxx_duckdb_sync 镜像表的查询任何时候都走 DuckDB,与该开关无关)。
输入:
SET lake_sync.auto_routing_enabled = on;
输出:
SET
追平后校验行数一致
输入:
SELECT 'heap' AS src, count(*) FROM ordersUNION ALLSELECT 'mirror' AS src, count(*) FROM orders_duckdb_sync;
输出:
src | count--------+-------heap | 3mirror | 3(2 rows)
列存表
支持的数据类型
列存表底层由 DuckDB 承载,大部分 PostgreSQL 常见标量与数组类型均可作为列类型。建表前建议先确认字段类型是否在下表范围内:
类别 | 支持类型 |
整数 | smallint / integer / bigint |
浮点 / 定点 | real / double precision / numeric(p, s) |
文本 | text / varchar(n) / char(n) |
布尔 | boolean |
时间 | date / time / timestamp / timestamptz |
二进制 | bytea |
JSON | json / jsonb |
唯一标识 | uuid |
一维数组 | int[] / bigint[] / text[] / numeric[] 等 |
实测建表与写入(在 de-postgres-7mvblukj 环境):
输入:
CREATE TABLE t_types (c_int integer,c_bigint bigint,c_num numeric(12,2),c_text text,c_ts timestamp,c_tstz timestamptz,c_bool boolean,c_bytea bytea,c_json json,c_jsonb jsonb,c_uuid uuid,c_arr_int int[],c_arr_text text[]) USING duckdb;INSERT INTO t_types VALUES (1, 9999999999, 3.14, 'hello','2026-07-22 10:00:00', '2026-07-22 10:00:00+08',true, '\\xdeadbeef', '{"a":1}', '{"a":1}','11111111-1111-1111-1111-111111111111',ARRAY[1,2,3], ARRAY['x','y']);SELECT * FROM t_types;
输出:
CREATE TABLEINSERT 0 1c_int | c_bigint | c_num | c_text | c_ts | c_tstz | c_bool | c_bytea | c_json | c_jsonb | c_uuid | c_arr_int | c_arr_text-------+------------+-------+--------+---------------------+------------------------+--------+------------+---------+----------+--------------------------------------+-----------+------------1 | 9999999999 | 3.14 | hello | 2026-07-22 10:00:00 | 2026-07-22 10:00:00+08 | t | \\xdeadbeef | {"a":1} | {"a": 1} | 11111111-1111-1111-1111-111111111111 | {1,2,3} | {x,y}(1 row)
创建
输入:
CREATE TABLE sales_analytics (ts timestamp,region text,amount numeric(12,2)) USING duckdb;
输出:
CREATE TABLE
也支持基于查询结果直接创建(CTAS)和基于已有表结构创建(LIKE)。
CTAS:从其他表灌入并落成列存
输入:
CREATE TABLE sales_summary USING duckdb ASSELECT product, sum(price*100) AS totalFROM salesGROUP BY product;SELECT * FROM sales_summary ORDER BY product;
输出:
CREATE TABLE ASproduct | total---------+---------apple | 100.000banana | 200.000(2 rows)
复用已有表结构
输入:
CREATE TABLE sales_copy (LIKE sales) USING duckdb;
输出:
CREATE TABLE
后续数据写入
建表完成后,支持标准 DML 语法继续追加、修改、删除数据。
操作 | 是否支持 | 推荐写法 |
批量导入 | ✅ 推荐 | COPY sales_analytics FROM '...' WITH (FORMAT csv); |
从其他表灌入 | ✅ 推荐 | INSERT INTO sales_analytics SELECT ... FROM staging; |
单行/小批量插入 | ✅ 支持,性能一般 | INSERT INTO sales_analytics VALUES (...); |
更新 | ✅ 支持 | UPDATE sales_analytics SET amount = ... WHERE ...; |
删除 | ✅ 支持 | DELETE FROM sales_analytics WHERE ...; |
清空重灌 | ✅ 推荐用于全量刷新 | TRUNCATE + INSERT ... SELECT |
DDL 变更 | ✅ 支持 | ALTER TABLE ... ADD/DROP/RENAME COLUMN |
以下逐项实测。
批量插入
输入:
INSERT INTO sales VALUES (1,'apple',1.5), (2,'banana',2.3);SELECT count(*) AS row_cnt FROM sales;
输出:
INSERT 0 2row_cnt---------2(1 row)
查询
输入:
SELECT * FROM sales ORDER BY id;
输出:
id | product | price----+---------+-------1 | apple | 1.52 | banana | 2.3(2 rows)
DDL 变更 (RENAME / ADD / ALTER TYPE / DEFAULT)
输入:
ALTER TABLE sales RENAME COLUMN product TO product_name;ALTER TABLE sales ADD COLUMN amount int;ALTER TABLE sales ADD COLUMN category text;ALTER TABLE sales ADD COLUMN status text DEFAULT 'active';ALTER TABLE sales ALTER COLUMN price TYPE numeric(10,2);
输出:
ALTER TABLEALTER TABLEALTER TABLEALTER TABLEALTER TABLE
UPDATE
输入:
UPDATE sales SET price = 2 WHERE product_name = 'apple';SELECT id, product_name, price FROM sales ORDER BY id;
输出:
UPDATE 1id | product_name | price----+--------------+-------1 | apple | 22 | banana | 2.3(2 rows)
DELETE
输入:
DELETE FROM sales WHERE product_name = 'banana';SELECT id, product_name, price FROM sales ORDER BY id;
输出:
DELETE 1id | product_name | price----+--------------+-------1 | apple | 2(1 row)
UPDATE / DELETE 的限制:必须支持 full query pushdown
列存表的 UPDATE / DELETE 必须能够整条 SQL 下推到 DuckDB 执行;若 SET 子句、WHERE 子句里出现不可下推的表达式(比如引用普通 heap 表的相关子查询、WITH ... UPDATE、部分不可下推的函数等),会直接报错。
查询
查询由 DuckDB 引擎执行并完成算子下推,常见的过滤、聚合、GROUP BY、JOIN 均可下推。
输入:
EXPLAINSELECT o_customer, sum(o_amount) AS totalFROM orders_duckdb_syncGROUP BY o_customerORDER BY total DESCLIMIT 5;
输出:
QUERY PLAN-----------------------------------------------------------------------------Custom Scan (Query Pushdown)Engine: DuckDB-> PROJECTION-> HASH_GROUP_BYGroups: #0Aggregates: sum_no_overflow(#1)-> PROJECTION-> SEQ_SCANType: Sequential ScanTable: postgres__public__orders_duckdb_sync
系统目录表现
列存表在 pg_class 中的表现:relkind = 'f'(foreign table),amname 为空。
输入:
SELECT n.nspname AS schema, c.relname AS name, c.relkind,COALESCE(a.amname, '-') AS amname,CASE c.relkind WHEN 'r' THEN 'ordinary table'WHEN 'f' THEN 'foreign table' END AS relkind_descFROM pg_class cJOIN pg_namespace n ON n.oid = c.relnamespaceLEFT JOIN pg_am a ON a.oid = c.relamWHERE n.nspname = 'public'AND c.relname IN ('demo_col','demo_heap','demo_sync','demo_sync_duckdb_sync')ORDER BY name;
输出:
schema | name | relkind | amname | relkind_desc--------+-----------------------+---------+--------+----------------public | demo_col | f | - | foreign tablepublic | demo_heap | r | heap | ordinary tablepublic | demo_sync | r | heap | ordinary tablepublic | demo_sync_duckdb_sync | f | - | foreign table(4 rows)
同步表的刷新方式
镜像表默认采用异步逻辑复制。系统提供6种刷新与维护 API:
使用场景 | API | 行为说明 |
强制镜像追赶至最新 LSN | lake_sync.force_sync(oid) | 通知 worker 立即刷新缓冲区并阻塞等待;返回 true 表示成功 |
暂停指定表的同步(用于大批量写入前) | lake_sync.pause_sync(oid) | 状态迁移至 PAUSED,后续写入不再同步至镜像 |
恢复同步(或从 ERROR 状态恢复) | lake_sync.resume_sync(oid) | 状态迁移至 ACTIVE;从 ERROR 恢复时自动重启 apply |
镜像数据出现不一致,需要重建 | lake_sync.reset_sync(oid) | 清空镜像表并从 heap 全量重导 |
全库镜像批量重建 | lake_sync.rebuild_sync_tables('duckdb' | 'iceberg' | 'all') | 全库镜像 drop、重建并全量重导 |
查询当前同步延迟 | lake_sync.get_sync_lag(oid) | 返回 interval 类型 |
下面按 API 逐条演示。
force_sync — 强制镜像追赶至最新 LSN
输入:
INSERT INTO orders SELECT g, 'cust'||g, g*10, now()FROM generate_series(100, 1099) g;SELECT lake_sync.force_sync('orders'::regclass);SELECT 'heap' AS src, count(*) FROM ordersUNION ALLSELECT 'mirror' AS src, count(*) FROM orders_duckdb_sync;
输出:
INSERT 0 1000force_sync------------t(1 row)src | count--------+-------heap | 1003mirror | 1003(2 rows)
pause_sync — 暂停指定表的同步
暂停后 heap 继续写入,mirror 停留在暂停前的位置。
输入:
SELECT lake_sync.pause_sync('orders'::regclass);SELECT table_name, sync_state FROM pg_catalog.sync_status WHERE table_name='orders';-- 暂停期间写入 1 行INSERT INTO orders VALUES (9999, 'paused_row', 99.99, now());SELECT 'heap' AS src, count(*) FROM ordersUNION ALLSELECT 'mirror' AS src, count(*) FROM orders_duckdb_sync;
输出:
NOTICE: Sync paused for table with OID 17645pause_sync------------(1 row)table_name | sync_state------------+------------orders | PAUSED(1 row)src | count--------+-------heap | 1004mirror | 1003(2 rows)
resume_sync — 恢复同步(或从 ERROR 状态恢复)
输入:
SELECT lake_sync.resume_sync('orders'::regclass);SELECT table_name, sync_state FROM pg_catalog.sync_status WHERE table_name='orders';
输出:
NOTICE: Sync resumed for table with OID 17645resume_sync-------------(1 row)table_name | sync_state------------+------------orders | ACTIVE(1 row)
reset_sync — 清空镜像并从 heap 全量重导
输入:
SELECT lake_sync.reset_sync('orders'::regclass::oid);SELECT pg_sleep(3);SELECT lake_sync.force_sync('orders'::regclass::oid);SELECT 'heap' AS src, count(*) FROM ordersUNION ALLSELECT 'mirror' AS src, count(*) FROM orders_duckdb_sync;
输出:
NOTICE: Shared replication slot sync_pool_slot_5_0 is not advanced by single-table reset_syncNOTICE: Sync reset completed for table with OID 17645reset_sync------------(1 row)pg_sleep----------(1 row)force_sync------------t(1 row)src | count--------+-------heap | 1004mirror | 1004(2 rows)
rebuild_sync_tables('duckdb') — 全库镜像重建
输入:
SELECT * FROM lake_sync.rebuild_sync_tables('duckdb');
输出:
NOTICE: pg_lake_sync: invalidating DuckDB mirror for public.customersNOTICE: pg_lake_sync: DuckDB mirror invalidated for public.customers, bgworker signaledNOTICE: pg_lake_sync: invalidating DuckDB mirror for public.ordersNOTICE: pg_lake_sync: DuckDB mirror invalidated for public.orders, bgworker signaledNOTICE: pg_lake_sync: invalidate complete: 2 invalidated, 0 skipped, 0 failedschema_name | heap_table | mirror_table | mirror_type | rows_imported | status-------------+------------+-----------------------+-------------+---------------+------------------------------------public | customers | customers_duckdb_sync | duckdb | 0 | ok: invalidated, bgworker signaledpublic | orders | orders_duckdb_sync | duckdb | 0 | ok: invalidated, bgworker signaled(2 rows)
get_sync_lag — 查询当前同步延迟
输入:
SELECT lake_sync.get_sync_lag('orders'::regclass::oid) AS orders_lag,lake_sync.get_sync_lag('customers'::regclass::oid) AS customers_lag;
输出:
orders_lag | customers_lag-----------------+-----------------00:00:03.320283 | 00:00:03.307904(1 row)
查询路由
开关(默认关闭)
会话级:
输入:
SET lake_sync.auto_routing_enabled = on;
输出:
SET
持久化(供新会话默认生效):
ALTER DATABASE <你的库名> SET lake_sync.auto_routing_enabled = on;
输出:
ALTER DATABASE
执行后请重新登录再
SHOW lake_sync.auto_routing_enabled 确认值为 on。业务查询使用的表名
业务查询统一使用原表名(如 orders),系统会根据查询谓词自动路由到 heap 或 mirror。不建议在业务代码中直接引用镜像表 orders_duckdb_sync。
路由规则
SQL 特征 | 路由目标 |
等值条件 + 索引列(点查) | heap(索引扫描) |
范围条件 / 无 WHERE / 全表扫描 | mirror(DuckDB) |
聚合操作( count / sum / group by) | mirror |
多表 JOIN 且带谓词下推 | mirror |
所有 INSERT / UPDATE / DELETE | 始终走 heap(路由仅影响 SELECT) |
说明:
对子查询和 CTE 递归应用路由规则:当同一条 SQL 含有子查询或 WITH CTE 时,系统会对每个子查询独立应用上述规则,而不是“整条 SQL 走同一目标”。例如外层是范围扫描走 mirror,内层却是索引点查,内层仍可能被路由至 heap。
通过 EXPLAIN 验证路由结果
等值查询 → heap 索引扫描
输入:
SET lake_sync.auto_routing_enabled = on;EXPLAIN (COSTS OFF) SELECT * FROM demo_sync WHERE id = 1;
输出:
QUERY PLAN----------------------------------------------Index Scan using demo_sync_pkey on demo_syncIndex Cond: (id = 1)(2 rows)
→ 路由至 heap,使用主键索引扫描。
范围查询 → mirror(DuckDB)
输入:
EXPLAIN (COSTS OFF) SELECT * FROM demo_sync WHERE id > 1;
输出:
QUERY PLAN-------------------------------------------------------------Custom Scan (Query Pushdown)Engine: DuckDB-> SEQ_SCANType: Sequential ScanTable: postgres__public__demo_sync_duckdb_syncFilters: id>1Estimated Cardinality: 4(7 rows)
→ 路由至 mirror,查询由 DuckDB 执行,Filters: id>1 表示谓词已下推。
聚合查询 → mirror(算子下推)
输入:
EXPLAIN (COSTS OFF) SELECT count(*) FROM demo_sync;
输出:
QUERY PLAN-------------------------------------------------------------------------Custom Scan (Query Pushdown)Engine: DuckDB-> UNGROUPED_AGGREGATEAggregates: count_star()-> PROJECTIONProjections: 42Estimated Cardinality: 20-> SEQ_SCANType: Sequential ScanTable: postgres__public__demo_sync_duckdb_syncProjections:Estimated Cardinality: 20(12 rows)
→ 路由至 mirror,聚合算子完成下推。
关闭路由后的对比
输入:
RESET lake_sync.auto_routing_enabled; -- 或 SET ... = offEXPLAIN (COSTS OFF) SELECT count(*) FROM demo_sync;
输出:
QUERY PLAN-----------------------------Aggregate-> Seq Scan on demo_sync(2 rows)
→ 回退至 heap 全表扫描。
手动指定路由目标
强制查询 mirror(直接引用镜像表)
输入:
SELECT o_customer, sum(o_amount) AS totalFROM orders_duckdb_syncGROUP BY o_customerORDER BY total DESCLIMIT 5;
输出:
o_customer | total------------+----------cust1099 | 10990.00cust1098 | 10980.00cust1097 | 10970.00cust1096 | 10960.00cust1095 | 10950.00(5 rows)
强制查询 heap(关闭自动路由)
输入:
SET lake_sync.auto_routing_enabled = off;SELECT o_id, o_customer, o_amount FROM orders WHERE o_id = 100;
输出:
SETo_id | o_customer | o_amount------+------------+----------100 | cust100 | 1000.00(1 row)
注意事项
路由仅对 SELECT 生效;所有 DML 均写入 heap。
镜像存在同步延迟,对实时性敏感的查询应先调用 force_sync 强制追赶,或显式关闭路由以直接访问 heap。
该能力还依赖 DuckDB 查询下推开关
pg_lake_engine.enable_duckdb_features = on(默认已启用);关闭该 GUC 后,即使路由把查询指向了 mirror,也会退回到普通 Foreign Scan 而非 Custom Scan (Query Pushdown),列存加速失效。自动路由 vs 查询下推:两个开关的差异
客户最容易混淆的问题:“我明明打开了 auto_routing_enabled,为什么这条 SQL 还是慢?”—— 通常原因是查询下推开关没打开。这两个开关控制的层级不同,组合起来才决定 SQL 的最终执行路径:
开关 | GUC | 层级 | 决定的事 |
自动路由 | lake_sync.auto_routing_enabled | 表名重写 | 通过 heap 表名的 SELECT 是否被改写到镜像表 xxx_duckdb_sync |
查询下推 | pg_lake_table.enable_full_query_pushdown | 执行计划 | 整条 SQL(仅涉及 pg_lake 表且所有算子可下推时)是否打包给 DuckDB 一次执行 |
三种组合的实测对比(以同步表 t_init_case 为例):
组合 A:
auto_routing = on + full_query_pushdown = on → 理想路径,SQL 走镜像并整条下推给 DuckDB。输入:
SET lake_sync.auto_routing_enabled = on;SET pg_lake_table.enable_full_query_pushdown = on;EXPLAIN (COSTS OFF) SELECT count(*) FROM t_init_case;
输出:
QUERY PLAN----------------------------------------------------------------------Custom Scan (Query Pushdown)Engine: DuckDB-> UNGROUPED_AGGREGATEAggregates: count_star()-> PROJECTION-> SEQ_SCANType: Sequential ScanTable: postgres__public__t_init_case_duckdb_sync
组合 B:
auto_routing = on + full_query_pushdown = off → SQL 已路由到镜像,但每一步都要在 PG 与 DuckDB 之间往返,退化为 Foreign Scan,列存加速失效。输入:
SET lake_sync.auto_routing_enabled = on;SET pg_lake_table.enable_full_query_pushdown = off;EXPLAIN (COSTS OFF) SELECT count(*) FROM t_init_case;
输出:
QUERY PLAN----------------------------------------------------------------------Foreign ScanRelations: Aggregate on (t_init_case_duckdb_sync)Engine: DuckDB-> UNGROUPED_AGGREGATEAggregates: count_star()-> SEQ_SCANTable: postgres__public__t_init_case_duckdb_sync
组合 C:
auto_routing = off → SQL 不路由,直接走 heap。输入:
SET lake_sync.auto_routing_enabled = off;EXPLAIN (COSTS OFF) SELECT count(*) FROM t_init_case;
输出:
QUERY PLAN-------------------------------------------------------------Aggregate-> Index Only Scan using t_init_case_pkey on t_init_case
排障建议:如果发现“路由已经开了但性能还是不行”,按此顺序检查:
1.
SHOW pg_lake_engine.enable_duckdb_features; 是否 on(总电源)。2.
SHOW pg_lake_table.enable_full_query_pushdown; 是否 on。3.
EXPLAIN 输出的第一行是不是 Custom Scan (Query Pushdown);若是 Foreign Scan,说明下推被拦下了。使用 routing_debug 排查“为什么这条 SQL 没走 mirror”
会话内开启 lake_sync.routing_debug,并把 client_min_messages 提到 DEBUG1,可以看到路由决策的完整链路:
输入:
SET lake_sync.auto_routing_enabled = on;SET lake_sync.routing_debug = on;SET client_min_messages = 'DEBUG1';SELECT count(*) FROM t_init_case WHERE id > 100;
输出:
DEBUG: pg_lake_sync: found sync table t_init_case (heap=49289, iceberg=49297)DEBUG: pg_lake_sync: skipping index for non-equality operator ">"DEBUG: pg_lake_sync: routing t_init_case to iceberg table t_init_case_duckdb_sync (sequential scan preferred)DEBUG: pg_lake_sync: query rewritten to use underlying table
日志能明确告诉你:识别到同步表 → 因为 > 不是等值操作跳过索引 → 决定路由到镜像 → 已改写。
排障完成后建议 RESET ALL 关闭 DEBUG 输出,以免污染业务日志。
数据库中同步表 / 列存表的盘点与识别
重要说明:同步表的 heap 侧在 pg_class 中 relkind = 'r'、amname = 'heap',与普通堆表在系统目录层面无法直接区分。识别同步表的权威来源为 lake_sync.sync_tables 或视图 pg_catalog.sync_status。
查看全部同步表状态
输入:
SELECT * FROM pg_catalog.sync_status;
输出:
schema_name | table_name | sync_state | last_sync_lsn | current_wal_lsn | last_sync_time | sync_lag_seconds | created_at | replication_slot_name | publication_name-------------+------------+------------+---------------+-----------------+-------------------------------+------------------+------------------------------+-----------------------+-------------------public | demo_sync | ACTIVE | 0/713DF50 | 0/7143D78 | 2026-07-22 16:08:33.260193+08 | 7 | 2026-07-22 16:08:33.14298+08 | sync_pool_slot_5_0 | sync_pool_pub_5_0(1 row)
字段说明:
sync_state — 同步状态,可能值:INITIALIZING — 初始化中。在 create_sync / reset_sync / rebuild_sync_tables 执行全量导入阶段短暂出现;当数据量大或后台负载重时窗口更明显。此状态下 apply worker 尚未开始处理增量,业务查询若走镜像可能读到不完整数据。ACTIVE — 正常增量同步中(稳态)。PAUSED — 已通过 pause_sync 显式暂停,heap 写入不会同步到镜像。ERROR — 同步失败(常见:镜像端存在冲突、逻辑复制槽异常),需要 resume_sync 重试,若仍失败则 reset_sync 重建镜像。sync_lag_seconds — 当前同步延迟(秒),监控告警的核心指标。last_sync_lsn / current_wal_lsn — 镜像最后同步位点与库当前 WAL 位点。replication_slot_name / publication_name — 底层复制槽与发布名称。同步表 heap 名与镜像名对照
输入:
SELECTn.nspname AS heap_schema,c.relname AS heap_table,c.relname || '_duckdb_sync' AS mirror_table_name,s.mirror_storage_type,s.sync_state,s.last_sync_time,s.last_sync_lsnFROM lake_sync.sync_tables sJOIN pg_class c ON c.oid = s.heap_table_oidJOIN pg_namespace n ON n.oid = c.relnamespaceORDER BY 1,2;
输出:
heap_schema | heap_table | mirror_table_name | mirror_storage_type | sync_state-------------+------------+-----------------------+---------------------+------------public | demo_sync | demo_sync_duckdb_sync | duckdb | ACTIVE(1 row)
全库列存表(USING duckdb)
输入:
SELECT n.nspname AS schema, c.relname AS nameFROM pg_class cJOIN pg_namespace n ON n.oid = c.relnamespaceWHERE c.relkind = 'f' -- 列存表与同步镜像均为 foreign tableAND c.relname NOT LIKE '%\\_duckdb\\_sync' ESCAPE '\\' -- 排除同步镜像AND n.nspname NOT IN ('pg_catalog','information_schema')ORDER BY 1,2;
输出(演示环境含 demo_col 列存表和 demo_sync_duckdb_sync 镜像,仅保留列存表):
schema | name--------+----------public | demo_col(1 row)
全库表类型汇总视图
输入:
SELECTn.nspname AS schema,c.relname AS name,CASEWHEN s.heap_table_oid IS NOT NULL THEN 'sync (heap side)'WHEN c.relname LIKE '%\\_duckdb\\_sync' ESCAPE '\\' THEN 'sync (duckdb mirror)'WHEN c.relkind = 'f' THEN 'duckdb columnstore'WHEN c.relkind = 'r' THEN 'plain heap'ELSE c.relkind::textEND AS kindFROM pg_class cJOIN pg_namespace n ON n.oid = c.relnamespaceLEFT JOIN lake_sync.sync_tables s ON s.heap_table_oid = c.oidWHERE n.nspname NOT IN ('pg_catalog','information_schema','pg_toast','lake_sync','lake_engine')AND c.relkind IN ('r','f')ORDER BY 1,2;
输出:
schema | name | kind--------+-----------------------+----------------------public | demo_col | duckdb columnstorepublic | demo_heap | plain heappublic | demo_sync | sync (heap side)public | demo_sync_duckdb_sync | sync (duckdb mirror)(4 rows)
常用监控 SQL
延迟超过60秒的同步表
输入:
SELECT schema_name, table_name, sync_lag_secondsFROM pg_catalog.sync_statusWHERE sync_lag_seconds > 60ORDER BY sync_lag_seconds DESC;
输出(演示环境健康,应为空集)
schema_name | table_name | sync_lag_seconds-------------+------------+------------------(0 rows)
处于异常状态的同步表
输入:
SELECT * FROM pg_catalog.sync_status WHERE sync_state <> 'ACTIVE';
输出(演示环境健康,应为空集):
schema_name | table_name | sync_state | ...-------------+------------+------------+-----(0 rows)
复制槽占用情况(用于判断是否需要调大 max_replication_slots)
输入:
SELECT count(*) AS used_slots,current_setting('max_replication_slots')::int AS max_slotsFROM pg_replication_slots;
输出:
used_slots | max_slots------------+-----------1 | 10(1 row)
常见问题
Q1. force_sync 一直不返回怎么办?
force_sync 内部会通知 apply worker 刷新缓冲区并阻塞等待追赶到当前 LSN。如果长时间不返回,通常是以下两种原因:
wal_level 未设置为 logical:同步表依赖逻辑复制槽,若 wal_level 为 replica,worker 无法建立复制槽,force_sync 会挂住。
排查:
SHOW wal_level;-- 期望输出:logical;若为 replica,需在 postgresql.conf 修改并重启实例
apply worker 未运行或已崩溃:检查表状态是否为 ERROR,或 max_worker_processes 是否被占满:
SELECT * FROM pg_catalog.sync_status WHERE table_name = 'your_table';SELECT current_setting('max_worker_processes')::int AS max_worker_processes,(SELECT count(*) FROM pg_stat_activity WHERE backend_type = 'background worker') AS running_bgworkers;
若 sync_state = ERROR,调 lake_sync.resume_sync(...) 恢复;仍失败则 lake_sync.reset_sync(...) 重建镜像。
Q2. 同步状态显示 ERROR 怎么处理?
同步表进入 ERROR 通常是逻辑复制流程中出现无法自动恢复的冲突(如 heap 表 DDL 改变了列类型、镜像端被外部改动等)。
恢复路径(按优先级):
1. resume_sync:轻量恢复,内部会重启 apply worker、重试上一批变更。适合“暂时性错误”。
输入:
SELECT lake_sync.resume_sync('your_table'::regclass);SELECT sync_state FROM pg_catalog.sync_status WHERE table_name='your_table';
输出:
NOTICE: Sync resumed for table with OID 17645sync_state------------ACTIVE(1 row)
2. reset_sync:重量级修复,会清空镜像并从 heap 全量重导。适合镜像数据已经和 heap 不一致的情况。
3. rebuild_sync_tables('duckdb'):全库同类型镜像一次性重建(用于批量恢复场景)。
Q3. UPDATE 报错“require full query pushdown”怎么改?
报错来自列存表的 UPDATE / DELETE:整条 SQL 必须能下推到 DuckDB 执行。常见触发原因:
SET 或 WHERE 子句里出现了引用普通 heap 表的相关子查询。使用了不可下推的 PG 函数(自定义 PL/pgSQL 函数、部分特殊内置函数)。
使用了修改型 CTE(
WITH ... UPDATE ... SELECT 之类的嵌套 DML)。修复思路:
简化表达式,把不可下推的部分拆到独立 SQL 里完成。
或者让所有相关表都变成 pg_lake 表(列存表 / 同步表),整条 SQL 就落在 DuckDB 内部。
或者退回到 heap:直接对同步表的 heap 侧执行 UPDATE(所有 DML 本来就走 heap,只有列存表原生模式才有这个限制)。
Q4. 如何确认一条查询是否真的走了 DuckDB 下推?
用 EXPLAIN,重点看第一行:
Custom Scan (Query Pushdown) → 已经下推,整条 SQL 打包给 DuckDB 执行,列存加速生效。
Foreign Scan → 只有部分算子在 DuckDB 侧执行,PG 与 DuckDB 之间需要多次往返,通常是 enable_full_query_pushdown 被关掉,或 SQL 中含不可下推算子。
Seq Scan / Index Scan on <heap_table> → 完全没走镜像,自动路由未生效(检查 lake_sync.auto_routing_enabled)。
进一步排查参见本文中自动路由 vs 查询下推:两个开关的差异与使用 routing_debug 排查“为什么这条 SQL 没走 mirror”:先看两个开关,再开 routing_debug 抓路由决策日志。