PostgreSQL码农集散地

突发: DuckDB Iceberg 数据湖插件支持实时DML

本期播客

突发: DuckDB Iceberg 数据湖插件支持实时DML

还在忍受数据湖只读不写的痛点?DuckDB-Iceberg v1.4.2全面颠覆!我们实现了Iceberg表的ACID全支持:插入、更新、删除,都在事务中保持一致。是时候放弃残缺的工具了。性能革命已至,即刻了解!

下面来自duckdb的官方博客

  • https://duckdb.org/2025/11/28/iceberg-writes-in-duckdb

一定要看到最后, 有彩蛋, 别看DuckDB Iceberg 数据湖插件支持实时 update/delete 了, 搞不好有哪些坑呢? 文末提出了几个挑战性问题!

DuckDB Iceberg 数据湖支持实时DML了

提要 (TL;DR): 我们为 DuckDB-Iceberg 扩展 (DuckDB-Iceberg extension) 发布了多项功能和改进:现在已全面支持 insert、update 和 delete 语句。

在过去几个月里,DuckDB Labs 团队一直致力于 DuckDB-Iceberg 扩展 (DuckDB-Iceberg extension) 的开发,并在 v1.4.0 中发布了完整的读取支持和初步的写入支持。今天,我们很高兴地宣布,针对 Iceberg v2 表 (Iceberg v2 tables) 的 删除 (delete) 和 更新 (update) 支持已在 v1.4.2 中推出!

Iceberg 开放表格式 (Iceberg open table format) 在过去两年中变得极其流行,许多数据库宣布支持这种最初在 Netflix 开发的开放表格式。在过去一年中,DuckDB 团队将 Iceberg 集成 (Iceberg integration) 列为优先事项,今天我们很高兴地宣布朝着这个方向迈出了又一步。在这篇博客文章中,我们将介绍 DuckDB v1.4.2 中 DuckDB-Iceberg 的当前功能集。

入门 (Getting Started)

要试用新的 DuckDB-Iceberg 功能,您需要连接到您偏好的 Iceberg REST 目录服务 (Iceberg REST Catalog) 。有多种连接到 Iceberg REST 目录服务的方法:如果您想连接到 Apache Polaris 或 Lakekeeper 等目录,请查看 连接到 REST 目录服务 (Connecting to REST Catalogs) ,如果您想连接到 Amazon S3 表 (Amazon S3 Tables) ,请查看 连接到 S3 表 (Connecting to S3 Tables) 页面。

ATTACH 'warehouse_name' AS iceberg_catalog (  
    TYPE iceberg,  
    other options  
);  

插入、删除和更新 (Inserts, Deletes and Updates)

在 DuckDB v1.4.0 中已经添加了创建表和向表中插入数据的支持:您可以使用标准的 DuckDB SQL 语法将数据插入到您的 Iceberg 表中。

CREATETABLE iceberg_catalog.default.simple_table (  
    col1 INTEGER,  
    col2 VARCHAR
);  
INSERTINTO iceberg_catalog.default.simple_table  
VALUES (1, 'hello'), (2, 'world'), (3, 'duckdb is great');  

您也可以使用任何 DuckDB 表扫描函数将数据插入到 Iceberg 表中:

INSERTINTO iceberg_catalog.default.more_data  
SELECT * FROM read_parquet('path/to/parquet');  

从 v1.4.2 开始,标准 SQL 语法也适用于 删除 (deletes) 和 更新 (updates) :

DELETEFROM iceberg_catalog.default.simple_table WHERE col1 = 2;  
UPDATE iceberg_catalog.default.simple_table SET col1 = col1 + 5WHERE col1 = 1;  
SELECT * FROM iceberg_catalog.default.simple_table;  

┌───────┬─────────────────┐  
│ col1  │      col2       │  
│ int32 │     varchar     │  
├───────┼─────────────────┤  
│     3 │ duckdb is great │  
│     6 │ hello           │  
└───────┴─────────────────┘  

Iceberg 写入支持目前有两个限制:

  1. 更新 (update) 支持仅限于未 分区 (partitioned) 且未 排序 (sorted) 的表。尝试使用 DuckDB-Iceberg 对已 分区 (partitioned) 或已 排序 (sorted) 的表执行 更新 (update) 、插入 (insert) 或 删除 (delete) 操作将导致错误。
  2. 对于 DELETE 和 UPDATE 语句,DuckDB-Iceberg 仅写入 位置删除 (positional deletes) 。写入时复制 (Copy-on-write) 功能尚未支持。

表属性函数 (Functions for Table Properties)

目前,DuckDB-Iceberg 仅支持 读取时合并 (merge-on-read) 语义。在 Iceberg 表元数据 (Table Metadata) 中,表属性 (table properties) 可用于描述允许哪种形式的 删除 (deletes) 或 更新 (updates) 。DuckDB-Iceberg 将遵守用于更新和删除的 write.update.mode 和 write.delete.mode表属性 (table properties) 。如果一个表具有这些属性,并且它们不是 读取时合并 (merge-on-read) ,DuckDB 将抛出错误,并且 UPDATE 或 DELETE 将不会被提交。v1.4.2 版本引入了三个新函数,用于添加、删除和查看 Iceberg 表的 表属性 (table properties) :

  • set_iceberg_table_properties
  • iceberg_table_properties
  • remove_iceberg_table_properties

您可以按如下方式使用它们:

-- to set table properties  
CALL set_iceberg_table_properties(iceberg_catalog.default.simple_table, {  
'write.update.mode': 'merge-on-read',  
'write.file.size': '100000kb'
});  
-- to read table properties  
SELECT * FROM iceberg_table_properties(iceberg_catalog.default.simple_table);  

┌───────────────────┬───────────────┐  
│        key        │     value     │  
│      varchar      │    varchar    │  
├───────────────────┼───────────────┤  
│ write.update.mode │ merge-on-read │  
│ write.file.size   │ 100000kb      │  
└───────────────────┴───────────────┘  

-- to remove table properties  
CALL remove_iceberg_table_properties(  
    iceberg_catalog.default.simple_table,  
    ['some.other.property']  
);  

Iceberg 表元数据 (Iceberg Table Metadata)

DuckDB-Iceberg 还允许您使用 iceberg_metadata() 和 iceberg_snapshots() 函数查看 Iceberg 表的 元数据 (metadata) 。

SELECT * FROM iceberg_metadata(iceberg_catalog.default.table_1);  

┌──────────────────────┬──────────────────────┬──────────────────┬─────────┬──────────────────┬─────────────────────────────────────────────────────────────┬─────────────┬──────────────┐  
│    manifest_path     │ manifest_sequence_…  │ manifest_content │ status  │     content      │                         file_path                           │ file_format │ record_count │  
│       varchar        │        int64         │     varchar      │ varchar │     varchar      │                          varchar                            │   varchar   │    int64     │  
├──────────────────────┼──────────────────────┼──────────────────┼─────────┼──────────────────┼─────────────────────────────────────────────────────────────┼─────────────┼──────────────┤  
│ s3://warehouse/def…  │                    1 │ DATA             │ ADDED   │ EXISTING         │ s3://<storage_location>/simple_table/data/019a6ecc-9e9e-7…  │ parquet     │            3 │  
│ s3://warehouse/def…  │                    2 │ DELETE           │ ADDED   │ POSITION_DELETES │ s3://<storage_location>/simple_table/data/d65b1db8-9fa8-4…  │ parquet     │            1 │  
│ s3://warehouse/def…  │                    3 │ DELETE           │ ADDED   │ POSITION_DELETES │ s3://<storage_location>/simple_table/data/8d1b92dc-5f6e-4…  │ parquet     │            1 │  
│ s3://warehouse/def…  │                    3 │ DATA             │ ADDED   │ EXISTING         │ s3://<storage_location>/simple_table/data/019a6ecf-5261-7…  │ parquet     │            1 │  
└──────────────────────┴──────────────────────┴──────────────────┴─────────┴──────────────────┴─────────────────────────────────────────────────────────────┴─────────────┴──────────────┘  

SELECT * FROM iceberg_snapshots(iceberg_catalog.default.simple_table);  

┌─────────────────┬─────────────────────┬─────────────────────────┬──────────────────────────────────────────────────────────────────────────────────────────────────────────────┐  
│ sequence_number │     snapshot_id     │      timestamp_ms       │                                                manifest_list                                                 │  
│     uint64      │       uint64        │        timestamp        │                                                   varchar                                                    │  
├─────────────────┼─────────────────────┼─────────────────────────┼──────────────────────────────────────────────────────────────────────────────────────────────────────────────┤  
│               1 │ 1790528822676766947 │ 2025-11-10 17:24:55.075 │ s3://<storage_location>/simple_table/data/snap-1790528822676766947-f09658c4-ca52-4305-943f-6a8073529fef.avro │  
│               2 │ 6333537230056014119 │ 2025-11-10 17:27:35.602 │ s3://<storage_location>/simple_table/data/snap-6333537230056014119-316d09bc-549d-46bc-ae13-a9fab5cbf09b.avro │  
│               3 │ 7452040077415501383 │ 2025-11-10 17:27:52.169 │ s3://<storage_location>/simple_table/data/snap-7452040077415501383-93dee94e-9ec1-45fa-aec2-13ef434e50eb.avro │  
└─────────────────┴─────────────────────┴─────────────────────────┴──────────────────────────────────────────────────────────────────────────────────────────────────────────────┘  

时间旅行 (Time Travel)

使用 AT (VERSION => ...) 或 AT (TIMESTAMP => ...) 语法,还可以通过 快照 ID (snapshot ids) 或 时间戳 (timestamps) 进行 时间旅行 (Time travel) 。

-- via snapshot id  
SELECT *  
FROM iceberg_catalog.default.simple_table AT (  
VERSION => snapshot_id  
);  

┌───────┬─────────────────┐  
│ col1  │      col2       │  
│ int32 │     varchar     │  
├───────┼─────────────────┤  
│     1 │ hello           │  
│     3 │ duckdb is great │  
└───────┴─────────────────┘  

-- via timestamp  
SELECT *  
FROM iceberg_catalog.default.simple_table AT (  
TIMESTAMP => '2025-11-10 17:27:45.602'
);  

┌───────┬─────────────────┐  
│ col1  │      col2       │  
│ int32 │     varchar     │  
├───────┼─────────────────┤  
│     1 │ hello           │  
│     3 │ duckdb is great │  
└───────┴─────────────────┘  

查看 Iceberg REST 目录服务的请求 (Viewing Requests to the Iceberg REST Catalog)

您可能还会好奇 DuckDB 向 Iceberg REST 目录服务 (Iceberg REST Catalog) 发出了哪些请求。为此,请启用 HTTP 日志记录 (HTTP logging) ,运行您的工作负载,然后从 HTTP 日志中查询。

CALL enable_logging('HTTP');  
SELECT * FROM iceberg_catalog.default.simple_table;  
SELECT request.type, request.url, response.status  
FROM duckdb_logs_parsed('HTTP');  

┌─────────┬──────────────────────────────────────────────────────────────────────────────────────────────────────────┬────────────────────┐  
│  type   │                                                                             url                          │       status       │  
│ varchar │                                                                           varchar                        │      varchar       │  
├─────────┼──────────────────────────────────────────────────────────────────────────────────────────────────────────┼────────────────────┤  
│ GET     │ https://<catalog_endpoint>/iceberg/v1/<warehouse>/iceberg-testing/namespaces/default                     │ NULL               │  
│ HEAD    │ https://<catalog_endpoint>/iceberg/v1/<warehouse>/iceberg-testing/namespaces/default/tables/simple_table │ NULL               │  
│ GET     │ https://<catalog_endpoint>/iceberg/v1/<warehouse>/iceberg-testing/namespaces/default/tables/simple_table │ NULL               │  
│ GET     │ https://<storage_endpoint>/data/snap-5943683398986255948-c2217dde-6036-4e07-88f2-…                       │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/f8c95b93-7b6b-4a24-8557-b98b553723d4-m0.avro                             │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/214a7988-da39-4dac-aa3a-4a73d3ead405-m0.avro                             │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/019a7244-c6e8-7bc9-9dd4-7249fcb04959.parquet                             │ PartialContent_206 │  
│ GET     │ https://<storage_endpoint>/data/019a7244-fcb5-7308-96ec-1c9e32509eab.parquet                             │ PartialContent_206 │  
│ GET     │ https://<storage_endpoint>/data/7f14bb06-f57a-42b4-ba7f-053a65152759-m0.avro                             │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/71f8b43d-51e7-40e7-be88-e8d869836ecd-deletes.parq…                       │ PartialContent_206 │  
│ GET     │ https://<storage_endpoint>/data/64f6c6e2-2f54-470e-b990-b201bc615042-m0.avro                             │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/4e54afed-6dd8-4ba0-88fb-16f972ac1d91-deletes.parq…                       │ PartialContent_206 │  
├─────────┴──────────────────────────────────────────────────────────────────────────────────────────────────────────┴────────────────────┤  
│ 12 rows                                                                                                                       3 columns │  
└─────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┘  

在这里,我们可以看到对 Iceberg REST 目录服务 (Iceberg REST Catalog) 的调用,随后是对 存储端点 (storage endpoint) 的调用。前三个对 Iceberg REST 目录服务的调用是为了验证模式 (schema) 仍然存在并获取 DuckDB-Iceberg 表的最新 metadata.json。接下来,它会查询 清单列表 (manifest list) 、清单文件 (manifest files) ,最后是包含数据和删除信息的文件。数据和删除文件会存储在本地 缓存 (cache) 中,以加快后续读取速度。

事务 (Transactions)

DuckDB 是一个支持 事务 (transactions) 的 ACID 兼容数据库 (ACID-compliant database) 。DuckDB-Iceberg 的开发也考虑到了这一点。在 事务 (transaction) 中,Iceberg 表将遵循以下条件。

  1. 表在 事务 (transaction) 中首次读取时,其 快照信息 (snapshot information) 会存储在 事务 (transaction) 中,并在该 事务 (transaction) 内部保持 一致 (consistent) 。
  2. 更新 (Updates) 、插入 (inserts) 和 删除 (deletes) 仅在 事务提交 (transaction is committed) 时(即 COMMIT)才会提交到 Iceberg 表中;

第 1 点对于读取 性能 (performance) 非常重要。如果您希望对 Iceberg 表进行 分析 (analytics) ,并且不需要每次都获取表的最新版本,那么在 事务 (transaction) 中运行您的 分析 (analytics) 将避免为每个查询获取最新版本。

-- truncate the logs  
CALL truncate_duckdb_logs();  
CALL enable_logging('HTTP')  
BEGIN;  
-- first read gets latest snapshot information  
SELECT * FROM iceberg_catalog.default.simple_table;  
-- subsequent read reads from local cached data  
SELECT * FROM iceberg_catalog.default.simple_table;  
-- get logs  
SELECT request.type, request.url, response.status  
FROM duckdb_logs_parsed('HTTP');  

┌─────────┬─────────────────────────────────────────────────────────────────────────────────────────────────────────────┬────────────────────┐  
│  type   │                                                  url                                                        │       status       │  
│ varchar │                                                varchar                                                      │      varchar       │  
├─────────┼─────────────────────────────────────────────────────────────────────────────────────────────────────────────┼────────────────────┤  
│ GET     │ https://<catalog_endpoint>/iceberg/v1/<warehouse>/iceberg-testing/namespaces/default                        │ NULL               │  
│ HEAD    │ https://<catalog_endpoint>/iceberg/v1/<warehouse>/iceberg-testing/namespaces/default/tables/simple_table    │ NULL               │  
│ GET     │ https://<catalog_endpoint>/iceberg/v1/<warehouse>/iceberg-testing/namespaces/default/tables/simple_table    │ NULL               │  
│ GET     │ https://<storage_endpoint>/data/snap-5943683398986255948-c2217dde-6036-4e07-88f2-1…                         │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/f8c95b93-7b6b-4a24-8557-b98b553723d4-m0.avro                                │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/214a7988-da39-4dac-aa3a-4a73d3ead405-m0.avro                                │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/019a7244-c6e8-7bc9-9dd4-7249fcb04959.parquet                                │ PartialContent_206 │  
│ GET     │ https://<storage_endpoint>/data/019a7244-fcb5-7308-96ec-1c9e32509eab.parquet                                │ PartialContent_206 │  
│ GET     │ https://<storage_endpoint>/data/7f14bb06-f57a-42b4-ba7f-053a65152759-m0.avro                                │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/71f8b43d-51e7-40e7-be88-e8d869836ecd-deletes.parquet                        │ PartialContent_206 │  
│ GET     │ https://<storage_endpoint>/data/64f6c6e2-2f54-470e-b990-b201bc615042-m0.avro                                │ OK_200             │  
│ GET     │ https://<storage_endpoint>/data/4e54afed-6dd8-4ba0-88fb-16f972ac1d91-deletes.parquet                        │ PartialContent_206 │  
├─────────┴─────────────────────────────────────────────────────────────────────────────────────────────────────────────┴────────────────────┤  
│ 12 rows                                                                                                                          3 columns │  
└────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────────┘  

在这里,我们看到了与上一节中看到的所有相同的请求。然而,现在我们处于一个 事务 (transaction) 中,这意味着当我们第二次从 iceberg_catalog.default.simple_table 读取时,我们不需要向 REST 目录服务 (REST Catalog) 查询表更新。这意味着 DuckDB-Iceberg 在第二次读取表时没有执行额外的请求,显著提高了 性能 (performance) 。

结论和未来工作 (Conclusion and Future Work)

凭借这些功能,DuckDB-Iceberg 现在为 Iceberg 表提供了强大的基础支持,使用户能够在他们的 Iceberg 表上发挥 DuckDB 的 分析能力 (analytical powers) 。未来还有更多的工作要做,并且 Iceberg 表规范 (Iceberg table specification) 还有许多 DuckDB 团队希望在 DuckDB-Iceberg 中支持的功能。如果您认为任何功能对您的 分析工作负载 (analytical workloads) 至关重要,请随时通过 DuckDB-Iceberg GitHub 仓库 (DuckDB-Iceberg GitHub repository) 与我们联系,或与我们的工程师取得联系。

以下是近期计划的改进列表(排名不分先后):

  • 性能改进 (Performance improvements)
  • 对 分区表 (partitioned tables) 的 更新 / 删除 / 插入 (Updates / deletes / inserts)
  • 对 排序表 (sorted tables) 的 更新 / 删除 / 插入 (Updates / deletes / inserts)
  • 模式演进 (Schema evolution)
  • 支持 Iceberg v3 表 (Iceberg v3 tables) ,重点关注 二进制删除向量 (binary deletion vectors) 和 行血缘跟踪 (row lineage tracking)

彩蛋

https://deepwiki.com/duckdb/duckdb-iceberg/4.2-insert-and-delete-operations

我的理解是否正确: “更新和删除iceberg数据时,并不是直接在原始文件中进行删除和修改, 而是存储一个新的文件, 用来标记原始文件中已删除的记录. ”?

同时我还有几个问题: 1、DML会不会导致iceberg文件碎片, 碎片会不会影响查询性能? 如何整理碎片? 2、已更新或已删除的旧记录如何进行整理和空间收缩?

看看ai怎么基于代码回复

您的理解是正确的。在Iceberg中,更新和删除操作确实不是直接修改原始文件,而是通过创建新的删除文件来标记已删除的记录  。

关于您的问题:

1. DML会导致文件碎片并影响查询性能

是的,DML操作会导致文件碎片:

  • 每次删除操作都会创建新的删除文件(positional delete files)
  • 测试显示,多次删除会产生多个删除文件

碎片会影响查询性能:

  • 查询时需要处理所有相关的删除文件
  • 删除文件通过ProcessDeletes方法处理,增加了查询开销

2. 空间整理和收缩机制

从当前代码库来看,没有实现自动的空间整理机制:

  • 删除文件会持续累积,没有自动清理过程
  • 旧的数据文件和删除文件会保留在存储中
  • 需要通过外部工具或手动操作进行清理

删除文件的处理流程:

  1. 删除操作收集要删除的行号
  2. 创建新的positional delete文件
  3. 生成manifest条目记录删除文件
  4. 更新快照包含删除文件

Notes

  • 当前实现主要关注删除操作的正确性,空间优化需要依赖外部工具
  • Iceberg v3计划通过deletion vectors改进这一机制,但当前代码主要处理v2格式
  • 查询性能影响主要体现在需要读取和处理额外的删除文件上