DuckDB 1.5 + Iceberg 写入深度拆解:当单机分析引擎终于长出湖仓的手——从 positional deletes 到 DuckLake v0.4 的元数据取舍
本文聊三件事:为什么"能写 Iceberg"比"能读 Iceberg"难一个数量级;DuckDB-Iceberg v1.4.2 的 insert/update/delete 到底是怎么落地的、坑在哪;以及 DuckLake 用一个 SQL 数据库换掉整套文件元数据的取舍是否值得。配可复现的 SQL/Python/Docker 实战,以及一份 18 条生产踩坑清单。
一、背景:湖仓的"只读时代"是怎么被卡住的
1.1 云上数据工程的基本范式
过去五年,云上数据分析的架构其实收敛到了一个非常朴素的形态:计算存储分离。
数据以 Parquet 格式躺在 S3 / GCS / OSS 上,便宜、持久、没有厂商锁定。需要查询的时候拉起计算资源,查完释放。这个范式在只读场景下几乎是完美的:
- 存储成本是块存储的十分之一量级
- Parquet 列式 + 字典编码 + 行组统计,扫描效率极高
- 任何引擎都能读:Spark、Trino、Flink、DuckDB、Polars、Pandas
问题出在写入上。
1.2 裸 Parquet 做不到的五件事
假设你有两条 pipeline 同时往同一张"表"(S3 上的一个目录前缀)写数据。如果两边都在尝试覆盖昨天的分区,谁先写完谁被覆盖?中间挂掉了留下半个文件怎么办?
裸 Parquet 目录有五个致命缺失:
- 没有并发写入协调:两个 writer 同时改同一个前缀,最后状态取决于运气
- 没有原子提交:写 1000 个文件写到第 700 个挂了,读者会看到一个不完整的表
- 没有 schema evolution:加一列、改类型、重命名,全靠约定和祈祷
- 没有时间旅行:数据写错了,昨天的状态找不回来
- 没有增量查询:下游想知道"上次消费之后新增了什么",只能比文件名和 mtime
Apache Iceberg(最初在 Netflix 开发)、Delta Lake、Apache Hudi 都是在回答同一个问题:如何在不放弃"对象存储 + 开放格式"这个核心理念的前提下,把数据库的表语义装回数据湖。这套东西被叫做 Lakehouse。
Iceberg 的答案是:用一堆 JSON 和 Avro 文件,在对象存储里手搓一个多版本并发控制系统。
1.3 DuckDB 的位置:从"读客户端"到"写参与者"
DuckDB 长期以来的定位是嵌入式 OLAP 引擎——SQLite 之于 OLTP,DuckDB 之于 OLAP。它的杀手场景是:
-- 一行 SQL 直接查 S3 上 200 个 Parquet 文件,不需要任何集群
SELECT date_trunc('day', ts) AS d, count(*)
FROM read_parquet('s3://bucket/logs/2026/*/*.parquet')
GROUP BY 1 ORDER BY 1;
对 Iceberg,DuckDB 早期只做了只读支持:能 scan、能读 manifest、能时间旅行。这在很多场景里够用(做 BI 查询、做数据探索),但一旦你想在 DuckDB 里做 ETL 的最后一公里——"把这批清洗好的数据写回 Iceberg 表"——就必须切回 Spark 或 PyIceberg。
DuckDB v1.4.0 加入了 CREATE TABLE 和 INSERT。DuckDB-Iceberg v1.4.2 补上了最后一块:DELETE 和 UPDATE。
这意味着一个非常有意思的可能性:一个几十 MB 的单进程二进制,可以作为 Iceberg 表的一等写入方参与生产链路。
后面我们会看到,这个"可以"是有明确边界的,而理解这些边界比知道语法重要得多。
二、核心概念:Iceberg 的四层结构与两种删除语义
要理解 DuckDB 的写入实现,必须先把 Iceberg 的物理结构搞清楚。很多人对 Iceberg 的理解停留在"就是带元数据的 Parquet",这个理解不足以让你判断性能问题出在哪。
2.1 四层寻址链
Iceberg 一张表的完整寻址链是四跳:
Catalog(目录服务)
↓ 记录:表名 → 当前 metadata.json 的路径(这一跳是唯一需要原子性的地方)
metadata.json(表元数据)
↓ 记录:schema 列表、分区 spec 列表、快照列表、当前快照 ID
snap-<id>.avro(manifest list,快照清单)
↓ 记录:这个快照包含哪些 manifest 文件,以及每个 manifest 的分区范围统计
<uuid>-m0.avro(manifest,清单文件)
↓ 记录:具体哪些 data file / delete file 属于本表,附带每列的 min/max/null_count
019a6ecc-....parquet(data file / delete file)
看懂这条链,你就能理解为什么 Iceberg 的小事务代价那么高:改一行数据,理论上需要写一个新的 data file 或 delete file、一个新的 manifest、一个新的 manifest list、一个新的 metadata.json,然后在 catalog 里做一次 CAS(compare-and-swap)把指针指过去。
一次 UPDATE = 至少 4 个新文件 + 1 次 catalog CAS。
这就是为什么 Iceberg 天生适合批量写、天生不适合高频小事务。这不是实现问题,是架构选择的必然结果。
2.2 snapshot 与 sequence_number
每次提交生成一个新的 snapshot(快照),带一个单调递增的 sequence_number。这个 sequence number 不只是版本号,它承担了一个关键的语义职责:决定 delete file 对哪些 data file 生效。
规则是:一个 delete file 只对 sequence_number 小于或等于它自己的 data file 生效。
这条规则解决了一个很微妙的问题。假设:
- seq=1:写入 data file A,包含 3 行
- seq=2:写入 delete file D,删除 A 的第 2 行
- seq=3:写入 data file B,也包含 3 行
如果没有 sequence number 约束,D 会不会误删 B 里的行?位置删除是按 (file_path, position) 定位的,理论上不会撞。但等值删除(equality delete)是按列值定位的,如果 B 里恰好有一行的主键跟 D 删掉的那行一样(比如"删除后重新插入"),没有 seq 约束就会被错误地重复删除。
Iceberg 用 sequence_number 把"删除"的作用域钉死在它诞生之前的数据上。 这是 Iceberg 元数据设计里最容易被忽略但最精妙的一处。
2.3 merge-on-read vs copy-on-write
这是湖仓写入最核心的一组权衡,值得单独讲。
Copy-on-Write(CoW,写时复制):
删一行 → 把包含这行的整个 data file 读出来 → 去掉这行 → 写一个新文件 → 在新快照里用新文件替换旧文件。
- 写放大巨大:删 1 行可能要重写 128MB
- 读性能完美:读者只需要扫 data file,零额外开销
Merge-on-Read(MoR,读时合并):
删一行 → 只写一个很小的 delete file,记录"某文件的第 N 行已删除" → 读者在扫描时把 delete file 加载进来,做反连接过滤。
- 写放大极小:写几 KB 就完事
- 读放大逐渐累积:delete file 越多,每次查询要额外读的文件越多、要维护的过滤结构越大
DuckDB-Iceberg 目前只实现 merge-on-read。 这个决定非常关键,它直接决定了你的运维策略——后面性能优化一节会详细展开。
2.4 positional delete vs equality delete
MoR 又分两种删除文件:
Positional delete(位置删除):记录 (file_path, row_position)。
file_path | pos
s3://wh/tbl/data/019a6ecc-9e9e-7....parquet | 1
s3://wh/tbl/data/019a6ecc-9e9e-7....parquet | 5
优点:读取时过滤极快(就是个按文件分组的有序位置集合,可以直接做游标推进)。缺点:写入方必须知道被删除行的物理位置,也就是必须先扫一遍数据。
Equality delete(等值删除):记录 id = 42。
优点:写入方不需要读数据,纯 append,Flink CDC 场景的首选。缺点:读取时要对每个 data file 做一次反连接,代价高得多,而且必须依赖 sequence number 严格约束作用域。
DuckDB-Iceberg 只写 positional deletes。
把这两条结论合起来看,DuckDB 的 Iceberg 写入定位就非常清楚了:
它是一个"知道自己在删什么"的批量写入方,不是一个流式 CDC 汇聚点。
DELETE 语句会先扫描定位,再写位置删除文件。这也解释了后面会讲到的一个限制——分区表和排序表不支持写入。
三、架构分析:DuckDB-Iceberg v1.4.2 的写入实现拆解
3.1 连接:Iceberg REST Catalog 是唯一入口
DuckDB 通过 ATTACH 把一个 Iceberg warehouse 挂成一个数据库:
INSTALL iceberg;
LOAD iceberg;
-- 通用 REST Catalog(Lakekeeper / Apache Polaris / Nessie / Gravitino)
ATTACH 'warehouse_name' AS iceberg_catalog (
TYPE iceberg,
ENDPOINT 'http://localhost:8181/catalog',
SECRET my_secret
);
注意这里的设计取舍:DuckDB 走 REST Catalog,不直接读 Hive Metastore、不支持 hadoop catalog 那种"扫目录找 metadata.json"的模式(只读场景下有 iceberg_scan 可以直接指路径,但写入必须有 catalog)。
原因很直接:写入需要原子的指针交换。REST Catalog 协议里有明确的 commit-table 语义,带 requirements 断言(比如"当前快照 ID 必须还是 X"),服务端做 CAS。没有这个,多写入方并发就无从保证。
v1.5.0 还补了两个细节:CREATE TABLE 时可以直接指定表属性,以及可以通过附加 HTTP 头连接 Google BigLake 这类需要特殊认证头的目录。
3.2 INSERT:标准语法,走 DuckDB 原生执行器
CREATE TABLE iceberg_catalog.default.simple_table (
col1 INTEGER,
col2 VARCHAR
);
INSERT INTO iceberg_catalog.default.simple_table
VALUES (1, 'hello'), (2, 'world'), (3, 'duckdb is great');
-- 更实用的形式:从任意 DuckDB 表函数灌数据
INSERT INTO iceberg_catalog.default.more_data
SELECT * FROM read_parquet('path/to/*.parquet');
第二种写法是真正的生产用法,它把 DuckDB 的全部读取能力接到了 Iceberg 写入端:read_csv、read_json、read_parquet、postgres_scan、mysql_scan、read_duckdb、S3 上的任意文件……都可以直接作为 Iceberg 表的数据源。
这实际上让 DuckDB 变成了一个零依赖的 Iceberg 入湖工具。以前这活儿要么写 PyIceberg 脚本,要么起 Spark。
3.3 UPDATE / DELETE:v1.4.2 的主角
DELETE FROM iceberg_catalog.default.simple_table WHERE col1 = 2;
UPDATE iceberg_catalog.default.simple_table SET col1 = col1 + 5 WHERE col1 = 1;
SELECT * FROM iceberg_catalog.default.simple_table;
┌───────┬─────────────────┐
│ col1 │ col2 │
│ int32 │ varchar │
├───────┼─────────────────┤
│ 3 │ duckdb is great │
│ 6 │ hello │
└───────┴─────────────────┘
语法上完全是标准 SQL,没有任何 Iceberg 特有的方言。但底层发生的事情值得展开:
DELETE 的执行路径:
- 扫描表,定位满足 WHERE 条件的行,同时记录它们的
(file_path, position) - 把这些位置按 file_path 分组、排序
- 写出一个 positional delete file(Parquet 格式,两列:file_path、pos)
- 生成新 manifest(content=DELETE)
- 生成新 manifest list,把旧的 DATA manifest 和新的 DELETE manifest 都列进去
- 生成新 metadata.json
- 向 catalog 提交 commit-table
UPDATE 的执行路径:UPDATE 在 MoR 语义下被分解为 delete + insert:
- 定位待更新行,写 positional delete file 标记旧行删除
- 把更新后的行作为新数据写一个新的 data file
- 两者在同一个快照里一起提交
这就是为什么 UPDATE 之后 iceberg_metadata() 里你会同时看到新的 DELETE manifest 和新的 DATA manifest。
3.4 表属性守门:DuckDB 会拒绝违约提交
Iceberg 表元数据里有两个属性描述了"这张表允许什么形式的更新":
write.update.modewrite.delete.mode
取值是 copy-on-write 或 merge-on-read。
DuckDB-Iceberg 会遵守这两个属性。如果表上写着 copy-on-write,而你在 DuckDB 里执行 UPDATE,DuckDB 会直接报错,并且不提交。
这个设计我个人非常认可。它避免了一个极其危险的场景:Spark 侧按 CoW 语义配置了表(下游读者可能是某个不支持读 delete file 的老引擎),DuckDB 偷偷写进去一堆 positional delete file,然后下游读到的是已删除的脏数据——因为它不知道要去 merge delete file。
在异构引擎共写一张表的环境里,表属性是唯一的契约载体。 DuckDB 选择尊重契约而不是尽力而为,是对的。
v1.4.2 为此引入了三个函数:
-- 设置表属性
CALL set_iceberg_table_properties(iceberg_catalog.default.simple_table, {
'write.update.mode': 'merge-on-read',
'write.file.size': '100000kb'
});
-- 读取表属性
SELECT * FROM iceberg_table_properties(iceberg_catalog.default.simple_table);
┌───────────────────┬───────────────┐
│ key │ value │
├───────────────────┼───────────────┤
│ write.update.mode │ merge-on-read │
│ write.file.size │ 100000kb │
└───────────────────┴───────────────┘
-- 移除表属性
CALL remove_iceberg_table_properties(
iceberg_catalog.default.simple_table,
['some.other.property']
);
3.5 事务:快照钉住(snapshot pinning)是个大特性
这是全篇我认为最值得单独强调的一点,因为它既是一致性保证,又是最大的性能杠杆。
DuckDB-Iceberg 的事务语义是:
- 事务内第一次读某张 Iceberg 表时,快照信息被记录进事务,并在整个事务内保持一致。
- UPDATE / INSERT / DELETE 只在 COMMIT 时才真正提交到 Iceberg 表。
第 1 点意味着:
BEGIN;
-- 第一次读:走 catalog 拿最新 metadata.json,解析 manifest 链
SELECT count(*) FROM iceberg_catalog.default.events;
-- 后续读:直接用事务内钉住的快照 + 本地缓存的数据文件
SELECT ... FROM iceberg_catalog.default.events GROUP BY ...;
SELECT ... FROM iceberg_catalog.default.events WHERE ...;
COMMIT;
没有 BEGIN 包裹的话,每个查询都要重新走一遍 catalog → metadata.json → manifest list → manifest 这条四跳链。
在对象存储上,这四跳每跳都是一次 HTTP 往返,每次几十到几百毫秒。一个 BI 面板刷十个 widget,如果每个都独立走一遍元数据解析,光元数据往返就能吃掉好几秒——而实际数据扫描可能只要 200ms。
结论:做分析型查询、不需要每次拿最新数据时,一定用显式事务包起来。 这不是可选优化,这是数量级差异。
3.6 观测:HTTP 日志是排查请求放大的唯一手段
DuckDB 提供了一个非常实用的观测能力——把对 catalog 和存储端点的 HTTP 请求全部记下来,然后用 SQL 查询这些日志:
CALL enable_logging('HTTP');
SELECT * FROM iceberg_catalog.default.simple_table;
SELECT request.type, request.url, response.status
FROM duckdb_logs_parsed('HTTP');
一次简单的全表扫描,输出大致是这样:
GET .../iceberg/v1/<warehouse>/.../namespaces/default -- 验证 namespace
HEAD .../namespaces/default/tables/simple_table -- 探测表存在
GET .../namespaces/default/tables/simple_table -- 拿 metadata.json
GET .../data/snap-5943683398986255948-....avro -- manifest list
GET .../data/f8c95b93-....-m0.avro -- manifest 1
GET .../data/214a7988-....-m0.avro -- manifest 2
GET .../data/019a7244-c6e8-....parquet (206 PartialContent) -- data file
GET .../data/019a7244-fcb5-....parquet (206 PartialContent) -- data file
GET .../data/7f14bb06-....-m0.avro -- delete manifest
GET .../data/71f8b43d-....-deletes.parquet (206 PartialContent) -- delete file
GET .../data/64f6c6e2-....-m0.avro
GET .../data/4e54afed-....-deletes.parquet (206 PartialContent)
12 个请求,读一张只有几行的表。
几个值得注意的信号:
206 PartialContent是好事:说明 DuckDB 在用 HTTP Range 请求做 Parquet 的部分读取(先读 footer 拿行组元数据,再只读需要的列块),不是傻乎乎拉整个文件- delete file 的数量直接映射到请求数:上面这张表有两个 delete manifest + 两个 delete file,就是四个额外请求。如果你的表被 UPDATE 打了几百次,这里就是几百个请求
- data file 和 delete file 会进本地缓存,后续读取会快很多,但缓存是进程级的,冷启动无法避免
这个日志能力配合下面的 SQL,就是一份现成的请求放大诊断脚本:
-- 按请求类型和资源类别统计请求数
SELECT
request.type,
CASE
WHEN request.url LIKE '%.avro' THEN 'manifest'
WHEN request.url LIKE '%deletes.parquet' THEN 'delete_file'
WHEN request.url LIKE '%.parquet' THEN 'data_file'
WHEN request.url LIKE '%metadata%' THEN 'metadata_json'
ELSE 'catalog_api'
END AS category,
count(*) AS reqs
FROM duckdb_logs_parsed('HTTP')
GROUP BY 1, 2
ORDER BY reqs DESC;
如果 delete_file 的请求数超过 data_file 的三倍,你的 compaction 已经欠债了。 这是我用得最多的一条经验阈值。
3.7 元数据自省函数
SELECT * FROM iceberg_metadata(iceberg_catalog.default.table_1);
┌────────────────────┬──────────────────────┬──────────────────┬─────────┬──────────────────┬─────────────┬──────────────┐
│ manifest_path │ manifest_sequence_no │ manifest_content │ status │ content │ file_format │ record_count │
├────────────────────┼──────────────────────┼──────────────────┼─────────┼──────────────────┼─────────────┼──────────────┤
│ s3://warehouse/... │ 1 │ DATA │ ADDED │ EXISTING │ parquet │ 3 │
│ s3://warehouse/... │ 2 │ DELETE │ ADDED │ POSITION_DELETES │ parquet │ 1 │
│ s3://warehouse/... │ 3 │ DELETE │ ADDED │ POSITION_DELETES │ parquet │ 1 │
│ s3://warehouse/... │ 3 │ DATA │ ADDED │ EXISTING │ parquet │ 1 │
└────────────────────┴──────────────────────┴──────────────────┴─────────┴──────────────────┴─────────────┴──────────────┘
这张表把前面讲的架构全部映射出来了,一行一行读:
- seq=1 一个 DATA manifest,3 条记录 → 最初的 INSERT
- seq=2 一个 DELETE manifest,POSITION_DELETES,1 条 →
DELETE WHERE col1 = 2 - seq=3 同时有 DELETE(1 条)和 DATA(1 条)→ 这就是 UPDATE 被分解成 delete + insert 的物证
SELECT * FROM iceberg_snapshots(iceberg_catalog.default.simple_table);
┌─────────────────┬─────────────────────┬─────────────────────────┐
│ sequence_number │ snapshot_id │ timestamp_ms │
├─────────────────┼─────────────────────┼─────────────────────────┤
│ 1 │ 1790528822676766947 │ 2025-11-10 17:24:55.075 │
│ 2 │ 6333537230056014119 │ 2025-11-10 17:27:35.602 │
│ 3 │ 7452040077415501383 │ 2025-11-10 17:27:52.169 │
└─────────────────┴─────────────────────┴─────────────────────────┘
3.8 时间旅行
-- 按快照 ID
SELECT * FROM iceberg_catalog.default.simple_table AT (VERSION => 6333537230056014119);
-- 按时间戳
SELECT * FROM iceberg_catalog.default.simple_table AT (TIMESTAMP => '2025-11-10 17:27:45.602');
时间旅行不只是"看历史"这个玩具功能,它在生产里有三个硬用途:
- 误操作恢复:
INSERT INTO tbl SELECT * FROM tbl AT (VERSION => 上一个好快照)之前先把当前数据备份 - 可复现报表:月度报表钉在月末那个快照上,避免"同一份 SQL 今天跑出来的数不一样"
- 增量对账:对比两个快照的差异来验证 pipeline 是否漏数
3.9 当前的两条硬限制
必须明确说清楚,这两条限制会直接决定你能不能用:
限制一:UPDATE / INSERT / DELETE 不支持已分区(partitioned)或已排序(sorted)的表。 对这类表执行写操作会直接报错。
限制二:DELETE 和 UPDATE 只写 positional deletes,copy-on-write 尚未支持。
限制一是真正的拦路虎。生产环境里的 Iceberg 事实表几乎都是分区的(按天、按小时、按业务维度)。这意味着 v1.4.2 的写入能力目前主要覆盖:
- 未分区的维度表 / 小事实表
- 中间态的 staging 表
- 数据量在单机可控范围内的分析型表
为什么分区表这么难?因为分区写入需要正确处理隐藏分区变换(Iceberg 的 hidden partitioning:days(ts)、bucket(16, id)、truncate(10, name)),写入方必须对每一行计算分区值、按分区分组、为每个分区生成独立的 data file,并在 manifest 里正确填写 partition 字段。排序表还要额外维护全局或局部有序性。这是一套不小的工程量。
判断:这是时间问题,不是路线问题。 但在它落地之前,不要按"DuckDB 可以写生产 Iceberg 事实表"来做架构设计。
四、代码实战
4.1 实战一:五分钟起一套本地湖仓(Lakekeeper + MinIO)
要练手 Iceberg 写入,最省事的是本地起一套 REST Catalog + S3 兼容存储。
# docker-compose.yml
version: '3.8'
services:
minio:
image: minio/minio:latest
command: server /data --console-address ":9001"
environment:
MINIO_ROOT_USER: minioadmin
MINIO_ROOT_PASSWORD: minioadmin
ports:
- "9000:9000"
- "9001:9001"
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:9000/minio/health/live"]
interval: 5s
retries: 10
mc-init:
image: minio/mc:latest
depends_on:
minio:
condition: service_healthy
entrypoint: >
/bin/sh -c "
mc alias set local http://minio:9000 minioadmin minioadmin;
mc mb --ignore-existing local/warehouse;
exit 0;
"
catalog-db:
image: postgres:16
environment:
POSTGRES_USER: catalog
POSTGRES_PASSWORD: catalog
POSTGRES_DB: catalog
healthcheck:
test: ["CMD-SHELL", "pg_isready -U catalog"]
interval: 5s
retries: 10
lakekeeper:
image: quay.io/lakekeeper/catalog:latest
depends_on:
catalog-db:
condition: service_healthy
environment:
LAKEKEEPER__PG_DATABASE_URL_READ: postgresql://catalog:catalog@catalog-db:5432/catalog
LAKEKEEPER__PG_DATABASE_URL_WRITE: postgresql://catalog:catalog@catalog-db:5432/catalog
RUST_LOG: info
command: ["serve"]
ports:
- "8181:8181"
docker compose up -d
# 等 lakekeeper 起来,然后初始化 warehouse(具体 API 见 Lakekeeper 文档)
curl -s http://localhost:8181/health
4.2 实战二:DuckDB 侧连接与建表
INSTALL iceberg;
LOAD iceberg;
INSTALL httpfs;
LOAD httpfs;
-- S3(MinIO)凭据
CREATE OR REPLACE SECRET minio_secret (
TYPE S3,
KEY_ID 'minioadmin',
SECRET 'minioadmin',
ENDPOINT 'localhost:9000',
URL_STYLE 'path',
USE_SSL false
);
-- 挂载 Iceberg catalog
ATTACH 'demo_warehouse' AS lake (
TYPE iceberg,
ENDPOINT 'http://localhost:8181/catalog'
);
CREATE SCHEMA IF NOT EXISTS lake.ops;
-- 未分区表(记住:分区表当前不支持写)
CREATE TABLE lake.ops.order_events (
event_id BIGINT,
order_id BIGINT,
status VARCHAR,
amount DECIMAL(18, 2),
updated_at TIMESTAMP
);
-- 建表时直接钉死 MoR 语义,避免后续被别的引擎改成 CoW
CALL set_iceberg_table_properties(lake.ops.order_events, {
'write.update.mode': 'merge-on-read',
'write.delete.mode': 'merge-on-read',
'write.file.size': '134217728' -- 128MB 目标文件大小
});
4.3 实战三:批量入湖 + 观察 delete file 累积
-- 从本地 Parquet 灌 100 万行
INSERT INTO lake.ops.order_events
SELECT
i AS event_id,
(i % 200000) AS order_id,
['created','paid','shipped','done'][1 + i % 4] AS status,
round((random() * 1000)::DECIMAL(18,2), 2) AS amount,
now() - (i % 86400) * INTERVAL '1 second' AS updated_at
FROM range(1, 1000001) t(i);
-- 基线:此刻的文件构成
SELECT manifest_content, content, count(*) AS manifests, sum(record_count) AS records
FROM iceberg_metadata(lake.ops.order_events)
GROUP BY 1, 2;
现在做一件"看起来很自然、实际很危险"的事——循环小批量更新:
-- ⚠️ 反面教材:50 次小事务,每次更新 2 万行
-- 真实生产里不要这么写,这里是为了观察 delete file 增长
UPDATE lake.ops.order_events SET status = 'refunded' WHERE order_id % 200 = 0;
UPDATE lake.ops.order_events SET status = 'refunded' WHERE order_id % 200 = 1;
-- ... 重复 50 次
-- 再看一次文件构成
SELECT
manifest_content,
content,
count(*) AS manifests,
sum(record_count) AS records
FROM iceberg_metadata(lake.ops.order_events)
GROUP BY 1, 2
ORDER BY 1, 2;
你会看到 POSITION_DELETES 的 manifest 数量线性增长。然后测一下查询代价:
CALL truncate_duckdb_logs();
CALL enable_logging('HTTP');
SELECT status, count(*), sum(amount)
FROM lake.ops.order_events
GROUP BY status;
-- 请求放大诊断
SELECT
CASE
WHEN request.url LIKE '%deletes.parquet' THEN 'delete_file'
WHEN request.url LIKE '%.parquet' THEN 'data_file'
WHEN request.url LIKE '%.avro' THEN 'manifest'
ELSE 'catalog_api'
END AS category,
count(*) AS reqs
FROM duckdb_logs_parsed('HTTP')
GROUP BY 1 ORDER BY reqs DESC;
这个实验是整篇文章里最值得你亲手跑一遍的部分。 你会亲眼看到 MoR 的读放大是怎么从"完全无感"变成"查询主要成本"的。
4.4 实战四:正确的批量 upsert 姿势(Python)
小事务是 Iceberg 的敌人。正确做法是攒批 + 单事务 + 一次性合并。
"""
duckdb_iceberg_upsert.py
把一批变更以「单事务、单次 DELETE + 单次 INSERT」的方式写进 Iceberg。
关键点:绝不在循环里逐条 UPDATE。
"""
import duckdb
import pandas as pd
from contextlib import contextmanager
def connect() -> duckdb.DuckDBPyConnection:
con = duckdb.connect()
con.execute("INSTALL iceberg; LOAD iceberg;")
con.execute("INSTALL httpfs; LOAD httpfs;")
con.execute("""
CREATE OR REPLACE SECRET minio_secret (
TYPE S3,
KEY_ID 'minioadmin', SECRET 'minioadmin',
ENDPOINT 'localhost:9000', URL_STYLE 'path', USE_SSL false
);
""")
con.execute("""
ATTACH 'demo_warehouse' AS lake (
TYPE iceberg,
ENDPOINT 'http://localhost:8181/catalog'
);
""")
return con
@contextmanager
def iceberg_tx(con):
"""显式事务:钉住快照 + 原子提交。异常自动回滚。"""
con.execute("BEGIN")
try:
yield con
con.execute("COMMIT")
except Exception:
con.execute("ROLLBACK")
raise
def upsert(con, table: str, changes: pd.DataFrame, key: str = "event_id"):
"""
单事务内完成 upsert:
1) 用 anti-join 语义删掉将被覆盖的旧行(一次 DELETE → 一个 delete file)
2) 一次 INSERT 灌入全部新行(一个 data file)
对比逐行 UPDATE:文件数从 2N 降到 2,catalog CAS 从 N 次降到 1 次。
"""
con.register("staging_changes", changes)
with iceberg_tx(con):
con.execute(f"""
DELETE FROM {table}
WHERE {key} IN (SELECT {key} FROM staging_changes)
""")
con.execute(f"""
INSERT INTO {table}
SELECT * FROM staging_changes
""")
con.unregister("staging_changes")
def health_check(con, table: str) -> dict:
"""delete/data manifest 比例,作为 compaction 触发信号。"""
rows = con.execute(f"""
SELECT manifest_content, count(*) AS n
FROM iceberg_metadata({table})
GROUP BY 1
""").fetchall()
stats = {c: n for c, n in rows}
data = stats.get("DATA", 0) or 1
delete = stats.get("DELETE", 0)
return {
"data_manifests": stats.get("DATA", 0),
"delete_manifests": delete,
"delete_ratio": round(delete / data, 2),
"needs_compaction": delete / data > 3.0,
}
if __name__ == "__main__":
con = connect()
batch = pd.DataFrame({
"event_id": [1, 2, 3, 4, 5],
"order_id": [1001, 1002, 1003, 1004, 1005],
"status": ["refunded"] * 5,
"amount": [10.5, 20.0, 30.25, 40.0, 50.75],
"updated_at": pd.Timestamp.utcnow().tz_localize(None),
})
upsert(con, "lake.ops.order_events", batch)
print(health_check(con, "lake.ops.order_events"))
核心思想:把 N 次小事务压成 1 次大事务。 文件数从 2N 降到 2,catalog CAS 从 N 次降到 1 次,写放大和后续读放大同时下降一个数量级。
4.5 实战五:用时间旅行做误操作恢复
-- 场景:某个 pipeline 把 status 全刷错了
-- 第一步:找到出事前的快照
SELECT sequence_number, snapshot_id, timestamp_ms
FROM iceberg_snapshots(lake.ops.order_events)
ORDER BY sequence_number DESC
LIMIT 10;
-- 第二步:先验证目标快照的数据是对的(别急着覆盖)
SELECT status, count(*)
FROM lake.ops.order_events AT (VERSION => 6333537230056014119)
GROUP BY status;
-- 第三步:差异对账 —— 到底影响了多少行
WITH good AS (
SELECT event_id, status FROM lake.ops.order_events AT (VERSION => 6333537230056014119)
),
now_ AS (
SELECT event_id, status FROM lake.ops.order_events
)
SELECT count(*) AS changed_rows
FROM good g JOIN now_ n USING (event_id)
WHERE g.status IS DISTINCT FROM n.status;
-- 第四步:单事务恢复
BEGIN;
CREATE OR REPLACE TEMP TABLE snapshot_backup AS
SELECT * FROM lake.ops.order_events AT (VERSION => 6333537230056014119);
DELETE FROM lake.ops.order_events
WHERE event_id IN (SELECT event_id FROM snapshot_backup);
INSERT INTO lake.ops.order_events SELECT * FROM snapshot_backup;
COMMIT;
注意第三步的差异对账。直接 rollback 是新手做法,因为你不知道出事之后有没有正常的新数据也一起被回滚掉。 先算清楚影响面,再决定是整表回滚还是按 key 精准修复。
4.6 实战六:VARIANT 类型接半结构化日志(DuckDB 1.5)
DuckDB 1.5.0 引入了原生 VARIANT 类型(灵感来自 Snowflake)。与 JSON 的关键区别是:VARIANT 存的是带类型的二进制数据,压缩率和查询性能都更好。
CREATE TABLE events (id INTEGER, data VARIANT);
INSERT INTO events VALUES
(1, 42::VARIANT),
(2, 'hello world'::VARIANT),
(3, [1, 2, 3]::VARIANT),
(4, {'name': 'Alice', 'age': 30}::VARIANT);
SELECT id, data, variant_typeof(data) AS vtype FROM events;
┌───────┬────────────────────────────┬───────────────────┐
│ id │ data │ vtype │
├───────┼────────────────────────────┼───────────────────┤
│ 1 │ 42 │ INT32 │
│ 2 │ hello world │ VARCHAR │
│ 3 │ [1, 2, 3] │ ARRAY(3) │
│ 4 │ {'name': Alice, 'age': 30} │ OBJECT(name, age) │
└───────┴────────────────────────────┴───────────────────┘
-- 点号提取嵌套字段(也可用 variant_extract)
SELECT data.name FROM events WHERE id = 4;
真实场景里最有价值的用法是处理 schema 漂移的埋点日志:
-- 上游埋点字段天天变,用 VARIANT 兜住,再按需投影出强类型列
CREATE TABLE raw_tracking (
received_at TIMESTAMP,
payload VARIANT
);
INSERT INTO raw_tracking
SELECT now(), json_payload::VARIANT
FROM read_json('s3://bucket/tracking/2026-08-*.json');
-- 下游按需投影,字段缺失返回 NULL 而不是炸掉
CREATE OR REPLACE VIEW tracking_typed AS
SELECT
received_at,
payload.user_id::BIGINT AS user_id,
payload.event::VARCHAR AS event_name,
payload.props.page::VARCHAR AS page,
payload.props.duration::DOUBLE AS duration_ms
FROM raw_tracking;
这套「VARIANT 落原始层 + View 做类型投影」的模式,比传统「上游改字段就要改建表语句」的做法健壮得多。 而且 Parquet 里的 VARIANT 可以直接读,支持"切碎"存储(shredding),也就是把高频出现的子字段单独抽成列存,兼顾灵活性和扫描性能。
需要留意:Iceberg v3 的 VARIANT 支持要到 DuckDB 1.5.1 才推出。所以现在 VARIANT 主要用在 DuckDB 原生表和 Parquet 上,还不能直接写进 Iceberg 表。
4.7 实战七:read_duckdb 跨库通配符查询
1.5.0 新增了一个不起眼但极其顺手的表函数:
-- 以前:每个库都要 ATTACH 一遍,然后手写 UNION ALL
ATTACH 'numbers1.db' AS db1;
ATTACH 'numbers2.db' AS db2;
ATTACH 'numbers3.db' AS db3;
-- SELECT ... UNION ALL SELECT ...
-- 现在:一行搞定
SELECT min(i), max(i) FROM read_duckdb('numbers*.db');
┌────────┬────────┐
│ min(i) │ max(i) │
├────────┼────────┤
│ 1 │ 5 │
└────────┴────────┘
它支持单文件、通配符、甚至配合 httpfs 读远程文件。最实用的场景是分库分表的历史数据分析:
-- 每天一个日志库,一口气查全年
SELECT
regexp_extract(filename, 'logs_(\d{8})', 1) AS day,
count(*) AS events,
count(DISTINCT user_id) AS uv
FROM read_duckdb('archive/logs_2026*.db')
GROUP BY 1
ORDER BY 1;
对于"手上一堆 .db 文件不知道里面有什么"的探索场景,省掉 ATTACH 这一步的心智负担比想象中大。
五、架构对比:DuckLake 用 SQL 数据库换掉文件元数据,值得吗
讲完 Iceberg,必须讲 DuckLake,因为它是对同一个问题的根本不同的回答,而且 DuckDB 团队自己在推。
5.1 核心分歧:元数据放哪里
Iceberg 的选择:元数据全部放对象存储。
snapshot JSON
↓
manifest list (Avro)
↓
manifest files (Avro)
↓
Parquet 数据文件
好处:不依赖任何外部系统,只要能读文件就能读数据。 任何引擎、任何语言,实现一个 Avro 解析器就能接入。这是 Iceberg 能成为事实标准的根本原因。
坏处:元数据操作代价高。 每次提交要写一串文件;查询要走四跳;列出快照要读 metadata.json;找某个分区要顺序扫 manifest。所有"数据库里一条索引查询就搞定"的事,在这里都是若干次对象存储往返。
DuckLake 的选择:元数据全部放标准 SQL 数据库,数据仍然是开放格式的 Parquet。
DuckLake 的核心洞察非常辛辣:
既然 Iceberg REST Catalog 最终也要用一个数据库来存表指针,那为什么不干脆从一开始就用数据库管理全部元数据?
这句话很难反驳。Iceberg 花了大量复杂度在"用文件模拟数据库",但生产部署里你照样要跑一个 Postgres 支撑的 catalog 服务。既然依赖已经存在,为什么只用它存一个指针?
DuckLake 的做法是:用 SQL 表来记录 schema、快照、文件列表、统计信息、分区信息。
对比效果:
| 操作 | Iceberg | DuckLake |
|---|---|---|
| 提交一次变更 | 写 delete/data file + manifest + manifest list + metadata.json + catalog CAS | 数据文件 + 一次 SQL 事务(几条 INSERT) |
| 列出所有快照 | 读 metadata.json 并解析 | SELECT * FROM ducklake_snapshot |
| 找某分区的文件 | 顺序扫 manifest(可用分区统计裁剪) | 一条带 WHERE 的索引查询 |
| 小事务代价 | 高(多次对象存储写 + 元数据放大) | 低(一次数据库事务) |
| 跨引擎互操作 | 极强(事实标准) | 需要引擎支持 DuckLake 规范 |
| 外部依赖 | 只需对象存储(理论上) | 必须有一个 SQL 数据库 |
5.2 DuckLake v0.4 的三个关键增强
DuckDB 1.5.0 把 DuckLake 规范更新到 v0.4,为即将发布的 DuckLake 1.0 铺路。三个新东西值得说:
1. 删除内联(delete inlining)
这是对 MoR 读放大问题的直接正面回答。既然元数据已经在数据库里了,小规模的删除为什么还要单独写一个 Parquet 文件?直接把删除标记内联到元数据表的行里。
这一下解决了本文第 4.3 节那个实验暴露的核心问题:高频小删除不再产生大量小 delete file。Iceberg 只能靠定期 compaction 来收拾,DuckLake 从源头上避免了债务积累。
2. 排序表(sorted tables)
在元数据里记录排序键和每个文件的排序范围,查询时可以做更激进的文件裁剪和归并优化。回想一下:DuckDB-Iceberg 目前不支持对排序表写入,而 DuckLake 原生支持——这个对比很能说明"元数据放数据库"带来的实现成本差异。
3. 宏(macros)
在元数据层存可复用的 SQL 宏定义。把"这个指标怎么算"这件事从各个下游 SQL 里收拢到表定义旁边,本质上是把轻量语义层下沉到湖仓元数据里。
另外,DuckDB 1.5.0 还把 DuckLake 扩展的体积减小了 30%(Excel 扩展减了 60%),这对 serverless / Lambda 这类冷启动敏感的部署有实际意义。
5.3 我的选型决策树
你需要多引擎共写同一张表(Spark 写、Trino 读、Flink 流入)?
├── 是 → Iceberg。这是唯一答案,生态就是护城河,别折腾。
└── 否 ↓
你的写入模式是「高频小事务」(分钟级甚至秒级提交)?
├── 是 → DuckLake(或 Iceberg + 强制攒批 + 激进 compaction)
│ Iceberg 硬扛小事务,运维成本会持续咬你
└── 否 ↓
你已经有一个可靠的 Postgres 且愿意让它成为元数据关键路径?
├── 是 → DuckLake。元数据查询快一个数量级,运维模型更简单。
└── 否 → Iceberg。至少你只需要对象存储可用。
补充判断:
- 需要「换引擎不迁移数据」的长期可选性 → Iceberg
- 团队只有 1~3 人、不想维护额外服务 → DuckLake + 现成托管 Postgres
- 已经在 S3 Tables / BigLake / Unity Catalog 上 → Iceberg / Delta,跟着平台走
一句话总结这个取舍:Iceberg 用实现复杂度换互操作性,DuckLake 用一个外部依赖换元数据效率。 没有对错,只有你的约束是什么。
5.4 顺带一提:Delta Lake 侧的进展
DuckDB 1.5.0 对 Delta Lake 也做了增强:支持写入 Unity Catalog、幂等写入(idempotent writes)、表检查点(checkpoints)。
幂等写入这个能力容易被忽略但很重要——它让 pipeline 重试变得安全。如果一个任务写到一半失败被调度器重试,没有幂等保证的话你会得到重复数据;有了它,同一个 txnAppId + txnVersion 的写入只会生效一次。
任何有重试机制的写入链路,都应该优先选支持幂等写入的路径。 这条经验在所有存储系统上都通用。
六、性能优化:八个真正有效的杠杆
6.1 用显式事务钉住快照(收益:数量级)
前面讲过,这是最大的单点优化。再强调一次量化直觉:
-- ❌ 十个独立查询 = 十次完整元数据解析 = 十 × 四跳 HTTP 往返
SELECT ... FROM lake.ops.t WHERE a = 1;
SELECT ... FROM lake.ops.t WHERE a = 2;
-- ...
-- ✅ 一次元数据解析,后续全部走缓存
BEGIN;
SELECT ... FROM lake.ops.t WHERE a = 1;
SELECT ... FROM lake.ops.t WHERE a = 2;
-- ...
COMMIT;
代价是你在事务内看不到别人的新提交。对分析型负载,这几乎总是可接受的——而且"报表期间数据不变"本来就是你想要的语义。
BI 后端接 Iceberg,务必在连接层做事务包裹。 我见过因为这一点让面板从 8 秒降到 900ms 的案例。
6.2 攒批写入,杀死小事务(收益:数量级)
一次写 100 万行和一百次各写 1 万行,最终数据量一样,但:
- 文件数:几个 vs 几百个
- manifest 数:几个 vs 几百个
- catalog CAS:1 次 vs 100 次(每次都有冲突重试风险)
- 后续查询的元数据解析代价:低 vs 高
规则:Iceberg 的提交频率应该以分钟计,不是以秒计。 上游是流的话,用一个攒批层(Kafka consumer 攒够 N 条或 T 秒再落一次)。
6.3 监控 delete/data 比例并定期 compaction(收益:防止劣化)
MoR 的读放大是累积型债务。加一条监控:
-- 建议接入定时任务,比例超阈值就告警或触发 compaction
SELECT
sum(CASE WHEN manifest_content = 'DELETE' THEN 1 ELSE 0 END) AS delete_manifests,
sum(CASE WHEN manifest_content = 'DATA' THEN 1 ELSE 0 END) AS data_manifests,
sum(CASE WHEN manifest_content = 'DELETE' THEN 1 ELSE 0 END)::DOUBLE
/ greatest(sum(CASE WHEN manifest_content = 'DATA' THEN 1 ELSE 0 END), 1) AS ratio
FROM iceberg_metadata(lake.ops.order_events);
经验阈值:
ratio < 1:健康1 ≤ ratio < 3:观察ratio ≥ 3:该 compaction 了,查询已经在为删除文件付费
注意:DuckDB-Iceberg 目前不提供 compaction 能力(rewrite_data_files 这类操作还得靠 Spark 的 Iceberg procedures 或 PyIceberg)。这是异构组合方案里必须提前规划的运维缺口——别等到查询变慢才想起来没人负责压实。
6.4 调 write.file.size(收益:显著)
CALL set_iceberg_table_properties(lake.ops.order_events, {
'write.file.size': '134217728' -- 128MB
});
目标文件大小的权衡:
- 太小(< 16MB):文件数爆炸,元数据膨胀,对象存储请求数上升,每个文件的 footer 读取都是固定开销
- 太大(> 512MB):并行度下降(一个文件难以被多线程切分得很好),单文件重写代价高,谓词裁剪粒度粗
- 甜点区:128MB ~ 256MB,这跟 HDFS 时代的块大小经验值收敛到同一个区间不是巧合——都是"并行度 × 元数据开销 × 单次 IO 效率"三者的平衡点
6.5 让 Parquet 统计信息真正起作用(收益:显著)
Iceberg 的 manifest 里存了每个文件每列的 min/max/null_count,Parquet 行组里也有一份。查询裁剪能不能生效,取决于数据的物理布局是否与查询谓词对齐。
-- ❌ 数据按插入顺序乱序存放,min/max 区间几乎覆盖全域,裁剪失效
-- 每个文件的 updated_at 范围都是 [最早, 最晚],谁都跳不掉
-- ✅ 写入前先排序,让每个文件的 min/max 区间尽量不重叠
INSERT INTO lake.ops.order_events
SELECT * FROM staging
ORDER BY updated_at; -- 高频过滤列排在前面
"写入时多花一次排序,换取后续每次查询都能跳过 90% 的文件",几乎总是划算的交易。 这是列存系统里最被低估的优化,比调任何参数都有效。
如果有多个常用过滤列,考虑 Z-order 或 Hilbert 曲线排序(需要 Spark 侧支持),能在多维上同时保持局部性。
6.6 善用本地缓存与 Range 请求(收益:中等,基本自动)
DuckDB 会把读到的 data file 和 delete file 缓存在本地,后续读取显著加速。前面 HTTP 日志里的 206 PartialContent 就是 Range 请求在工作——只读 Parquet 需要的列块,不拉整个文件。
你能做的:
- 长驻进程 > 一次性进程。每次新起进程都是冷缓存,把 DuckDB 跑成常驻服务(或复用连接)能吃到缓存红利
- SELECT 明确列名,不要
SELECT *。列存下,多读一列就是多一次 Range 请求 + 多一份解压
6.7 非阻塞检查点(DuckDB 1.5,收益:17% 量级)
这条针对的是 DuckDB 原生存储(.duckdb 文件),不是 Iceberg,但同样重要。
1.5.0 实现了非阻塞检查点:检查点期间可以并发读取、写入、插入(带索引)和删除。TPC-H SF100 吞吐量提升 17%。聚合函数也有优化,last 函数提速 40%。
对"用 DuckDB 做本地暂存表、定期落湖"这种混合架构,检查点不再是写入毛刺的来源了。以前长事务遇上检查点会出现明显的延迟尖刺,现在这个问题基本消失。
6.8 网络栈换成 libcurl(收益:稳定性)
httpfs 扩展的 HTTP 后端从 httplib 换成了 curl。原有配置项(超时、重试)保持兼容,但稳定性和安全性上了一个台阶——尤其在弱网、代理、TLS 中间件复杂的企业环境里,curl 久经考验的连接处理能减少很多莫名其妙的偶发失败。
如果你之前遇到过"S3 读取偶发超时/连接重置且无法复现",升级到 1.5 值得一试。
七、生产踩坑清单(18 条)
分区表和排序表当前不能写。 生产 Iceberg 事实表几乎都分区,动手前先确认表结构,别写到一半才发现报错。
write.update.mode/write.delete.mode不是 merge-on-read 时,DuckDB 会拒绝提交。 这是保护而非 bug。异构引擎共写场景下,不要为了让 DuckDB 能写就把表属性强改成 MoR——先确认所有下游读者都能正确 merge delete file。只有 positional deletes,没有 copy-on-write。 意味着删除量大的表读放大会持续累积,compaction 必须有人负责。
别在循环里逐条 UPDATE。 每条一个事务 = 一串文件 + 一次 catalog CAS。攒批到单事务,文件数和请求数同时降一个数量级。
分析查询一定用显式事务包起来。 不包的话每个查询都重走四跳元数据链,纯浪费。
DuckDB-Iceberg 不做 compaction。 需要 Spark procedures 或 PyIceberg 补位。架构评审时就要把这个缺口写进方案,别留到线上变慢再补。
iceberg_metadata()里的 delete/data manifest 比例是最好的健康指标。 超过 3:1 就该动手了,接进监控。快照会无限累积,记得配置过期策略。 时间旅行很爽,但每个快照都钉住了一批文件不能被清理,存储成本和元数据解析代价都会涨。
SELECT *在列存 + 对象存储上是双重浪费。 多一列 = 多一次 Range 请求 + 多一份解压 CPU。冷启动无法避免元数据解析。 Serverless / CLI 一次性任务的元数据开销占比可能高得离谱。能长驻就长驻。
写入前按高频过滤列排序。 这是投入产出比最高的优化,让 min/max 统计真正能裁剪文件。
write.file.size别用默认值就完事。 128MB~256MB 是甜点区,小于 16MB 会让元数据成为瓶颈。HTTP 日志是唯一的请求放大真相来源。
enable_logging('HTTP')+duckdb_logs_parsed('HTTP'),用 SQL 查自己的网络行为,这个能力比大多数商业工具的可观测性都直接。Iceberg v3 的 VARIANT 支持要等 1.5.1。 现在 VARIANT 只能用在 DuckDB 原生表和 Parquet 上,别在设计里假设它能直接落 Iceberg。
箭头 lambda 语法
x -> x + 1从 1.5 起会发弃用警告,DuckDB 2.0 默认禁用。 改成 Python 风格lambda x: x + 1,或用lambda_syntax配置过渡。存量脚本要提前批量改。空间扩展的轴序要变。 GEOMETRY 类型已内置到核心。
geometry_always_xy设置:v1.5 默认保持旧行为(纬度/经度)但发警告,v2.0 旧行为报错,v2.1 起新行为(经度/纬度)成为默认。有 GIS 逻辑的项目现在就要开始改,这是会静默算错结果的那种变更——不报错,只是坐标反了。read_duckdb('*.db')会读取匹配库里的全部表。 通配符范围写宽了会拉进意外的数据源,先用小范围确认再放开。v1.4 是 LTS,维护到 2026 年 9 月;下一个版本是 DuckDB 2.0(预计 9 月)。 生产环境要么钉 1.4 LTS 稳一段,要么做好 2.0 破坏性变更(箭头 lambda、轴序)的迁移预案。别在 2.0 发布前夕才开始评估。
八、总结与展望
8.1 这次变更真正的意义
如果只记一件事,我希望是这个:
DuckDB 加上 Iceberg 写入,本质上是在挑战"分析型写入必须依赖分布式计算框架"这个默认假设。
过去要往湖仓写数据,起点就是 Spark 集群或 Flink 作业——不是因为数据量非要那么大,而是因为只有它们实现了表格式的写入协议。协议实现的复杂度,而非计算量,成了事实门槛。
DuckDB 用一个几十 MB 的单进程二进制把这个门槛按下去了一截。对大量"数据量其实只有几十 GB、但被迫用 Spark 只为了能写 Iceberg"的团队,这是实实在在的架构简化:少一个集群,少一套调度,少一批 YARN/K8s 运维,少一堆 JVM 调参。
但边界也很清楚:分区表不能写,就意味着它现在还进不了大多数生产事实表的主链路。它当下的最佳位置是:
- 维度表和小事实表的完整生命周期管理
- ETL 的 staging 层与最后一公里
- 数据探索、临时修数、误操作恢复
- 单机能吃下的分析型表的全量管理
8.2 三个值得盯的方向
方向一:分区表写入。 这是 DuckDB-Iceberg 从"能用"到"能上生产主链路"的唯一门槛。隐藏分区变换的写入实现是硬骨头,但没有理论障碍。这个能力落地那天,值得重新评估你的入湖架构。
方向二:DuckLake 1.0 与"元数据即数据库"范式。 v0.4 的删除内联已经展示了这条路的威力——它从源头消除了 MoR 读放大这个 Iceberg 必须靠运维压实来对抗的顽疾。如果 DuckLake 1.0 能吸引到足够的引擎支持,"元数据放文件"这个 Iceberg 的立身之本可能会被重新审视。
值得注意的行业信号:腾讯云 PostgreSQL 已经把 DuckDB 作为向量化执行引擎内置进 PG 实例,通过一条 SET 语句就能让分析查询走 DuckDB 路径、事务查询继续走 PG 原生路径。当 DuckDB 开始出现在托管数据库的执行层里,"元数据放 SQL 库"这个设计的现实基础就比两年前扎实多了。
方向三:Iceberg v3 与 VARIANT。 半结构化数据落湖一直是个别扭活儿——JSON 字符串列压缩差、查询慢,硬拆成强类型列又扛不住 schema 漂移。Iceberg v3 的 VARIANT 加上 shredding,是这个问题第一个看起来靠谱的答案。DuckDB 1.5.1 会跟上。
8.3 一个更大的判断
把 DuckDB 1.5、DuckLake v0.4、Iceberg 写入这些点连起来看,我认为在发生的是数据栈的"去分布式化"。
过去十年的默认叙事是"数据量在涨,所以需要更大的集群"。但实际情况是:单机能力涨得比大多数公司的数据量更快。 NVMe 的顺序读能到几 GB/s,一台机器插几百 GB 内存不算奢侈,向量化执行引擎把单核效率榨到了接近内存带宽极限。
一个残酷的事实:很多公司的 Spark 集群,处理的数据量单机 DuckDB 一分钟就能扫完。 他们付的是分布式的全部代价(调度延迟、序列化开销、shuffle、运维复杂度、JVM 调参、集群成本),换来的是本来不需要的扩展性。
湖仓格式的开放性在这里是关键前提:因为数据是 Parquet + 开放元数据,你可以在"单机够用时用 DuckDB、真的不够时上 Spark"之间自由切换,而不需要迁移数据。 这个可选性本身就是巨大的架构价值——它把"选错了就要重来"变成了"随时可以换"。
DuckDB 补上 Iceberg 写入,是把这个可选性从"只读"扩展到"读写"。这一步比它看起来更重要。
最后一句实践建议:如果你手上有一套 Iceberg 湖仓,花半小时把第 4.3 节那个 delete file 累积实验跑一遍。看到 delete_file 请求数超过 data_file 请求数的那一刻,你对 merge-on-read 的理解会比读十篇文章都深。
参考:DuckDB 官方博客 Iceberg writes 相关文章、DuckDB 1.5.0 发布说明(代号 Variegata,自 v1.4 以来近 100 位贡献者提交超过 6500 个 commit)、DuckLake v0.4 规范、Apache Iceberg 表规范。