DuckDB 可以增量加载吗?维度历史呢?
昨天讨论了DuckDB 的 Luajit 扩展里,我塞了个 ETL 引擎,今天继续深入这个话题。
每天凌晨 2 点同步订单表,全量重跑 10 分钟——其实一天就多了几百行新数据。客户维度改个地址,直接 UPDATE 覆盖,历史状态就没了,想查"他去年住哪"翻不到。增量加载 + 维度历史,数据工程师的日常,但每次都要手写水位游标、手写版本表。
第一个坑是效率:数据量上来之后,全量重跑从 10 分钟涨到 1 小时,但增量可能就几百行——99% 的时间在搬没变过的老数据。第二个坑是留痕:财务月底要对账,问"这个客户 6 月的时候注册地址是哪",维度表被后来的 UPDATE 覆盖了,谁都答不上来。
我的 DuckDB 里跑着不少同步任务,之前每次写增量都复制粘贴同一段 WHERE:WHERE ts > (SELECT max(ts) FROM 目标表)。慢不说,还会踩坑——同一时间戳落库的两行,只进一行。维度表更麻烦:客户搬家了,旧的记录直接覆盖,审计要回溯历史版本的时候什么都没有。
所以我把这两件事封装成了 DuckDB 原生函数:给 duckdb-luajit 扩展写了两个 Lua 库,一共 150 行。一个管增量加载,一个管维度历史。
● ● ●
incremental:水位游标,三种模式
incremental 库维护一张水位表 etl_watermark(name, last_ts, last_id, updated_at)——每个目标表一行,记录上次加载到哪。每次运行只取水位之后的新行:
SELECT luajit_s('incremental', {'op':'run',
'target':'tgt_orders', 'source':'src_orders',
'ts_col':'ts', 'mode':'ts'});
-- → {"loaded":3,"last_ts":"2026-08-02 09:00:00","last_id":null}
返回 JSON,SQL 里 json_extract 直接取 loaded 数和最新水位。加载完自动更新水位表,下次自动从断点继续。
三种游标模式,覆盖常见增量场景:
- ●
ts:WHERE ts > last_ts——按时间增量,适合日志、订单 - ●
id:WHERE id > last_id——按自增主键,适合流水 - ●
ts_id:WHERE ts > last_ts OR (ts = last_ts AND id > last_id)——复合游标
第三种是我踩过坑才加的:只按时间游标,同一秒内落库的多行会被漏掉一半;加个自增主键做同刻消歧,不漏不重。
增量加载水位游标流程图
实测:源表 3 行首载全量(水位表里还没有记录,自动当全量处理),再补 2 行只加载 2 行,源没动时返回 {"loaded":0}——重跑多少次都不会重复加载,幂等。想强制重来一遍?清掉水位表对应行就行。
● ● ●
scd2:维度历史版本化
scd2 库是缓慢变化维度类型 2(SCD Type 2):属性一变,旧版本自动闭合、新版本顶上,全程留痕。SCD 类型 1 是直接覆盖(省空间、丢历史),类型 2 是每变一次留一个版本(多占空间、可回溯)——要能回答"6 月时他住哪"这种问题,只能选类型 2。自动建表时带上四个管理列:_valid_from / _valid_to / _is_current / _version:
SELECT luajit_s('scd2', {'op':'run',
'target':'dim_customer', 'source':'src_customer',
'key_col':'cust_id', 'attr_cols':['city','tier'], 'ts_col':'ts'});
-- → {"inserted":1,"closed":1}
客户从杭州搬到北京,库里长这样:
杭州 v1 valid_from=08-01 valid_to=08-02 is_current=false
北京 v2 valid_from=08-02 valid_to=NULL is_current=true
再搬家就是 v3、v4,一条时间线完整保留。想查"去年他住哪"就是 WHERE _is_current = false 翻历史。
属性比对用 md5 指纹——把 city、tier 这些列拼起来算哈希,哈希变了才动。源表没变化时 0 插入 0 关闭,幂等。这正是 dbt snapshots 里 check 策略的思路:全列比对,变了才开新版本。我们只是把比对换成了指纹,几百个属性列也算得动。
● ● ●
白嫖 DuckLake:增量直接写数据湖
写到这我想起来:增量加载最合适的落点其实是数据湖。DuckDB 官方的 DuckLake 扩展(SQL + Parquet 的开源湖格式,2025 年开源)直接 ATTACH 当普通库用——incremental.run 的目标表指向 DuckLake 表,一行代码不用改:
ATTACH 'ducklake:meta.ducklake' AS dl (DATA_PATH 'data/');
CREATE TABLE dl.orders(id BIGINT, ts TIMESTAMP, amount DOUBLE);SELECT luajit_s('incremental', {'op':'run',
'target':'dl.orders', 'source':'src_orders',
'ts_col':'ts', 'mode':'ts'});
-- → {"loaded":3,...} 增量进湖
白赚三件套:
- 1.时间旅行:
FROM dl.orders AT (VERSION => 2)看任意历史快照——月初跑的数能原样复盘,不怕下游改数 - 2.变更数据流:
FROM dl.table_changes('orders', 1, 2)拿到每个快照的 insert/update/delete,下游团队可以直接消费变更流,不用再蹭我们的水位 - 3.数据落 parquet:Spark、Polars 都能直接读;小批量变更自动内联进 catalog 数据库(官方叫 Data Inlining),不写一堆几 KB 的小文件——等数据攒大了,自然落成 parquet 文件
实测组合:首载 3 行 → 增量 2 行 → DuckLake 里 5 行;时间旅行快照 2 是 3 行、快照 3 是 5 行,分毫不差。整个过程 incremental 的代码没动过一个字符——DuckLake 表就是普通目标表,水位逻辑照样跑。
水位归 Lua,存储归湖——各干各的。
● ● ●
和 dbt、dlt 什么关系?
有人会问:这不就是 dbt 的 incremental、dlt 的增量吗?机制同源,定位不同:
- ●dbt:增量模型靠
is_incremental()宏,过滤条件用户自己写——官方文档原话"the rows that you tell dbt to filter for",最常见就是查目标表max(),没有独立水位表。SCD2 是dbt snapshots:check 策略全列比对,生成dbt_scd_id/dbt_valid_from/dbt_valid_to——和我们的指纹 + 版本列一模一样。但 dbt 是整套 transform 层:项目、文档、测试、语义层,适合数据团队。 - ●dlt:Python 库,
incremental()游标跟踪 cursor 字段(timestamp/ID),加载状态存在 pipeline state 里,面向 source→loader,适合写 Python 的工程师。 - ●这两个 Lua 库:DuckDB 原生 SQL 层的单文件函数,独立水位表,零外部依赖。社区 301 个 DuckDB 扩展里没有原生的增量加载和 SCD2——这是真空位。dbt-duckdb 用户照样可以在模型里调它。
不是替代谁,是让 DuckDB 自己先把这两件事干了。
● ● ●
诚实边界
- ●需要 luajit 扩展的普通模式(非 trusted)——函数内部要执行 SQL
- ●表对表:
source和target列序一致 - ●scd2 的源表约定每业务键一行(快照式),不是 CDC 变更流
- ●没有 dbt 的依赖图、文档化、测试框架——就是个函数,别指望它管全家桶
150 行,两件事,配 DuckDB 官方数据湖。代码在 duckdb-luajit-libs 仓库的 libs/etl/ 下。
参考来源:
- ●github.com/alitrack/duckdb-luajit-libs(incremental / scd2 源码与测试)
- ●ducklake.select 官方文档(DuckLake 扩展、Data Inlining、Change Data Feed)
- ●dbt Labs 官方文档:Incremental models、Snapshots
- ●dlt 官方文档:Incremental loading