alitrack

DuckDB 的 Luajit 扩展里,我塞了个 ETL 引擎

上周我在 DuckDB 的 Lua 扩展里塞了个 C 编译器(libtcc),能在 SQL 里现场编译 C 代码。这周我干了件更"正经"的事:用这个 Lua 层造了一个完整的 ETL 引擎——从数据接入、清洗转换、质量检查到落库,一条管道跑完,每个环节自动记审计。

起因是一个叫 ETLX 的开源项目。它用 Markdown 文件当 ETL 配置——一份 pipeline.md 同时是配置、是文档、是审计记录,跑完自动生成流程图。这个"pipeline 即文档"的想法很妙,但它的实现方式我不认同:用 Markdown 编排 SQL,表达能力太弱了。

我的 duckdb-luajit 扩展里有完整的 LuaJIT——变量、循环、条件、函数、错误处理,哪一样不比 Markdown 强?为什么不用 Lua 编排 DuckDB?

● ● ●

四大能力:从 ETLX 偷师,用 Lua 重做

1. 审计:每次运行自动落表

ETLX 最值钱的模式叫 AUTO_LOGS——每次运行自动注入日志块(开始/结束/耗时/行数/成功/错误)。我的版本更彻底:审计直接写进 DuckDB 表,SQL 可查。

-- etl_run_log 表:ts / fn / args / duration_ms / rows / success / err
SELECT * FROM etl_run_log ORDER BY ts DESC LIMIT 10;

每跑一步,一条记录:(步骤名, 耗时毫秒, 影响行数, 成功与否, 错误原因)。ETL 跑挂了,先查这张表,不用翻日志。

2. 幂等:防重复加载

数据管道最怕的事:昨天跑了一半挂了,今天重跑一遍,订单表里多了一倍数据。ETLX 的 load_validation 就是防这个——加载前检查目标表状态:

etl.validate('orders_clean', 'throw_if_not_empty')
-- 目标表已有数据 → 直接报错:already has 100 rows (load twice?)

空表才让过,非空就拦——管道永远不重跑。

3. 自愈:表不存在就自动建

etl.insert_auto('orders_clean',
'CREATE TABLE orders_clean (id BIGINT, name VARCHAR, amount DOUBLE)',
"INSERT INTO orders_clean SELECT ...")

目标表不存在?错误消息匹配 does not exist → 自动执行建表 SQL → 重试一次。一行代码,少一个人工介入。

4. SQL 组件化:查询像乐高

etl.q({
select = 'c.name, SUM(o.amount) AS total',
from = 'src_orders o',
join = 'JOIN src_customers c ON o.cust_id = c.id',
where = 'o.amount > 0',
group_by = 'c.name',
order_by = 'total DESC',
})

动态拼 SQL,比字符串拼接干净,比 ORM 透明。

● ● ●

Pipeline 引擎:画布上拖出来的,这里一行 JSON 跑完

Duckle 这类工具的核心是"拖拽画布 → Pipeline JSON → 执行"。我把最后一步做进了 Lua:Pipeline JSON 直接编译执行,还带拓扑排序和环检测。

local pipe = {
name = 'orders_etl',
nodes = {
{ id = 'src', type = 'source', query = "SELECT * FROM orders_raw" },
{ id = 'clean', type = 'transform', query = "SELECT id, upper(name) AS name, amount FROM src WHERE amount > 0",
deps = { 'src' } },
{ id = 'check', type = 'quality', checks = {
{ type = 'not_null', col = 'id' },
{ type = 'unique', col = 'id' },
{ type = 'rowcount_min', min = 1 },
}, deps = { 'clean' } },
{ id = 'enrich', type = 'transform_lua',
fn = function(row) row.name = row.name .. '!'; return row end,
deps = { 'clean' } },
{ id = 'out', type = 'sink', table = 'orders_clean', deps = { 'enrich' } },
},
}
etl.pipeline(pipe)

实测一次 5 节点管道:src=4 → clean=3(过滤负数)→ check=3 → enrich=3 → out=3,每个节点自动记审计。

五种节点类型,对应主流 ETL 工具的组件矩阵:

  • ●source:一条查询建临时表
  • ●transform:SQL 转换(dbt 的 model 干的事)
  • ●transform_lua:Lua 闭包逐行处理——row.name = row.name .. '!' 这种逻辑不用拼 SQL 字符串
  • ●quality:not_null / unique / rowcount / schema 检查,失败即中止 + 记审计
  • ●sink:落正式表

质量检查失败不是静默的——success=false, err=原因 进审计表。DuckDB 的管道,每一步都有据可查。

● ● ●

和 dbt、dlt 是什么关系?

做完我才意识到,这个 Lua ETL 层正好卡在 dbt 和 dlt 之间:

dbt  = SQL 项目管理(lineage / 测试 / 物化)
dlt = Python EL 层(100+ 源 / schema 演化)
etl.lua = 流程层:编排 + 审计 + 质量 + 自愈
  • ●dbt 替代不了:它管的是 100 个 SQL 模型的依赖图,Lua 不碰这个
  • ●dlt 替代不了:它的长处是 100+ 已验证数据源和 schema 演化机制,Lua 重写成本爆炸
  • ●但轻量 EL 场景:一个 HTTP API 拉数据 → Lua 直接写读取逻辑,进程内、零 Python 依赖、零 venv——比 dlt 的 pipeline 对象轻一个数量级

三者的关系不是竞争,是分层。重数据源用 dlt,大规模模型用 dbt,中间这层"编排 + 质量 + 审计"——我用 Lua 写在了 DuckDB 里面。

● ● ●

为什么是 Lua?

有人会问:SQL 不也能做 ETL 吗?能,但 SQL 没有控制流。DuckDB 有 macro,但 macro 不支持循环、条件、错误处理、闭包。

Lua 补的正是这些:pcall 错误处理、闭包传参、表结构自由组合、ffi 拉 C 库。而且因为整个引擎跑在扩展进程内(_duckdb_call 直调同一个连接),没有进程间通信,没有序列化开销。

顺手测了个数独提速当 benchmark:朴素回溯 1236.7ms/题 → 位运算 + MRV 启发式 0.14ms/题,9000 倍——同样的算法,同样的 DuckDB,Lua 层写得快。

● ● ●

结尾

ETLX 教会我"pipeline 即文档"的理念,但我选择用 Lua 而不是 Markdown 来实现它——因为 Lua 本身就是比 Markdown 强得多的"文档":可执行、可测试、有错误处理。

现在 duckdb-luajit 的 libs 仓库里躺着这个 ETL 引擎,一条 install etl 就能用。DuckDB 的 Lua 层里,我造了个 ETL 引擎——下次谁再说"ETL 需要一整套独立工具",我会说:DuckDB 里就有。


参考来源

  • ●https://github.com/alitrack/duckdb-luajit
  • ●https://github.com/alitrack/duckdb-luajit-libs
  • ●https://github.com/realdatadriven/etlx
  • ●https://github.com/slothflowlabs/duckle
  • ●https://github.com/dbt-labs/dbt-core
  • ●https://github.com/dlt-hub/dlt