- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
导读
Apache Iceberg 的表维护(Maintenance)是保证数据湖表健康运行的核心工作:每次写入都会产生新的快照(Snapshot),若不定期清理,元数据与数据文件会持续膨胀,拖慢查询规划并推高存储成本。本文基于 docs/docs/maintenance.md 展开,结合本仓库源码,系统讲解 Iceberg 推荐的例行维护操作(过期快照、清理历史元数据文件、删除孤儿文件)与可选优化操作(数据文件压缩、Manifest 重写、Position Delete 重写、悬挂删除文件清理、表统计计算),并覆盖删表后的存储清理与表路径重写等特殊 Action。读完本文,你将掌握如何用 Java API 与 Spark Actions 完成全套 Iceberg 表维护,并理解每个操作背后的元数据机制与安全边界。
所有维护操作都基于
Table实例(或 Spark 中的SparkActions)执行。加载已有表的完整方式可参考 Java API 快速入门。
为什么 Iceberg 表需要维护
Iceberg 采用“写时复制 + 元数据版本化”的架构:每次提交都会生成一个新的metadata.json(元数据文件)、Manifest List 与 Manifest 文件,数据文件本身不可变(immutable)。这种设计带来了快照隔离(Snapshot Isolation)、时间旅行(Time Travel)与原子提交能力,但也意味着:
- 快照无限累积:每个快照都引用着一批数据文件,只有显式过期(expire)才会释放底层文件;
- 元数据文件不断新增:每次提交产生一个新 JSON 文件用于保证原子性,高频写入(如流式作业)会快速堆积;
- 失败残留孤儿文件:分布式任务失败、部分写入中止时,会留下未被任何元数据引用的文件。
因此,维护的本质是:在保留时间旅行能力的前提下,回收不再被任何有效快照引用的数据与元数据文件。
推荐的例行维护(Recommended Maintenance)
过期快照(Expire Snapshots)
每次对 Iceberg 表的写入都会创建一个新的快照,即表的一个版本。快照可用于时间旅行查询,也可将表回滚到任意有效快照。快照会不断累积,直到通过expireSnapshots操作过期它们。定期过期快照是推荐的维护手段,目的是:
- 删除不再需要的数据文件;
- 保持表元数据体积小巧。
以下示例过期 1 天以前的快照:
Table table = ...; long tsToExpire = System.currentTimeMillis() - (1000 * 60 * 60 * 24); // 1 day table.expireSnapshots() .expireOlderThan(tsToExpire) .commit();从源码看,ExpireSnapshots接口(api/src/main/java/org/apache/iceberg/ExpireSnapshots.java)提供了丰富的链式配置:
| 方法 | 作用 |
|---|---|
expireSnapshotId(long) | 按快照 ID 精确过期某个快照 |
expireOlderThan(long) | 过期所有早于给定时间戳(毫秒)的快照 |
retainLast(int) | 保留当前快照最近的前 N 个祖先快照,即使它们早于过期时间戳也不会被过期 |
cleanupLevel(CleanupLevel) | 控制过期时的文件清理级别:NONE(仅移除快照元数据,不做文件清理)、METADATA_ONLY(只清理 Manifest、Manifest List、统计文件等元数据文件,保留数据文件)、ALL(元数据与数据文件都清理,默认值)。当数据文件被多个表共享(例如通过 add-files 过程加入)时,建议使用METADATA_ONLY |
deleteWith(Consumer<String>) | 自定义删除实现(例如将待删文件收集起来而不实际删除) |
executeDeleteWith(ExecutorService)/planWith(ExecutorService) | 并行执行删除或规划 |
cleanExpiredMetadata(boolean) | 清理未使用的分区规格(partition specs)、Schema 等元数据 |
提交时,这些变更会应用到最新的表元数据;发生提交冲突时,Iceberg 会将变更重新应用到新的最新元数据并重试提交。过期操作也会删除不再被有效快照使用的 Manifest 文件,以及被过期快照删除掉的数据文件。注意:过期操作不允许删除当前快照。
对于大表,还可以使用 Spark Action 并行执行过期(源码入口见 SparkActions.expireSnapshots):
Table table = ...; SparkActions .get() .expireSnapshots(table) .expireOlderThan(tsToExpire) .execute();重要:过期快照只是把它们从元数据中移除,使它们不再可用于时间旅行。数据文件只有在不再被任何可能用于时间旅行或回滚的快照引用时才会真正被删除。因此,定期过期快照是释放底层数据文件的前提。
删除旧元数据文件(Remove old metadata files)
Iceberg 使用 JSON 文件跟踪表元数据。表的每次变更都会产生一个新的元数据文件以提供原子性。默认情况下,旧元数据文件会被保留以支撑历史回溯;对提交频繁的表(如流式写入的表),需要定期清理元数据文件。
每个元数据文件会在metadata-log字段中跟踪更旧的元数据文件。被跟踪的元数据文件数量由write.metadata.previous-versions-max定义。
要在提交后自动删除更旧的元数据文件,请在表属性中设置write.metadata.delete-after-commit.enabled=true。这会保留若干被跟踪的元数据文件(最多write.metadata.previous-versions-max个),并且每次创建新元数据文件时删除最旧的一个。注意:该机制只会删除被metadata-log跟踪的元数据文件,不会删除孤儿元数据文件。未被跟踪的元数据文件需要通过孤儿文件删除流程清理。
属性定义可对照 core/src/main/java/org/apache/iceberg/TableProperties.java 中的常量与其默认值:
| 属性 | 默认值 | 描述 |
|---|---|---|
write.metadata.delete-after-commit.enabled | false | 控制每次表提交后是否删除最旧的被跟踪版本元数据文件 |
write.metadata.previous-versions-max | 100 | 最多跟踪的上一版本元数据文件数量 |
示例说明:
- 设
write.metadata.delete-after-commit.enabled=false、write.metadata.previous-versions-max=10:提交 100 次后,会有 10 个被跟踪的元数据文件与 90 个孤儿元数据文件。这 90 个孤儿文件无法通过设置write.metadata.delete-after-commit.enabled=true删除(因为它们已不再被跟踪),只能通过孤儿文件删除流程清理。 - 设
write.metadata.delete-after-commit.enabled=true、write.metadata.previous-versions-max=20:提交 21 次后,会有 20 个被跟踪的元数据文件,最旧的元数据文件在写入者提交时被删除。此后每次新提交都会删除最旧的一个元数据文件。
更多写入属性请参见 表写入属性配置。
删除孤儿文件(Delete orphan files)
在 Spark 及其他分布式处理引擎中,任务或作业失败可能留下未被表元数据引用的文件;某些情况下,普通的快照过期也无法判定某个文件不再需要并将其删除。
要清理表目录下的这些“孤儿”文件,可使用deleteOrphanFilesAction:
Table table = ...; SparkActions .get() .deleteOrphanFiles(table) .execute();从 api/src/main/java/org/apache/iceberg/actions/DeleteOrphanFiles.java 的接口定义看,该 Action 的核心语义是:一个文件如果无法被任何有效快照到达,即视为孤儿;判定方式是对底层存储进行目录列举(listing),因此该操作代价较高。
常用配置项:
| 配置 | 默认 | 说明 |
|---|---|---|
location(String) | 表根目录 | 指定扫描孤儿文件的目录;不设置则扫描整个表目录,可能同时删除孤儿数据与元数据文件 |
olderThan(long) | 3 天前的时间戳 | 只删除早于该时间戳的孤儿文件,避免误删正在写入、尚未被元数据引用的文件 |
deleteWith(Consumer<String>) | 表的 FileIO | 自定义删除实现(例如只收集不删除) |
prefixMismatchMode | ERROR | 当元数据引用的文件与列举到的文件在 scheme/authority 上不一致时的处理策略:ERROR(抛异常,推荐)、IGNORE(跳过不匹配文件)、DELETE(将不匹配文件视为孤儿删除,极危险) |
equalSchemes(Map)/equalAuthorities(Map) | 空 | 声明等价 scheme / authority,例如Map("s3a,s3,s3n", "s3")、Map("s1name,s2name", "servicename"),用于解决前缀不匹配冲突 |
该 Action 在数据与元数据目录中文件很多时可能耗时较长,建议定期执行,但无需过于频繁。
安全警告 1:如果孤儿文件保留间隔(retention interval)短于任何一次写入完成所需的时间,那么进行中的文件可能被视为孤儿而被删除,从而损坏表。默认间隔为 3 天,除非必要不要调短。
安全警告 2:Iceberg 使用路径的字符串表示来判断哪些文件需要删除。在某些文件系统上,路径可能随时间变化但仍指向同一文件。例如更换 HDFS 集群的 authorities 后,创建时使用的旧路径 URL 与当前列举结果不再匹配,运行 RemoveOrphanFiles 会导致数据丢失。请务必确保 MetadataTables 中的条目与 Hadoop FileSystem API 列举的条目一致,以免误删。
可选维护(Optional Maintenance)
部分表需要额外维护。例如流式查询可能产生大量小文件,应压缩为更大的数据文件;某些表可以从重写 Manifest 文件中获益,让查询定位数据更快。
压缩数据文件(Compact data files)
Iceberg 跟踪表中的每个数据文件。数据文件越多,Manifest 文件中存储的元数据就越多;小数据文件还会带来不必要的元数据开销与较差的查询效率(文件打开成本高)。
Iceberg 可以通过 Spark 的rewriteDataFilesAction 并行压缩数据文件,将小文件合并为更大的文件,从而降低元数据开销与运行时文件打开成本:
Table table = ...; SparkActions .get() .rewriteDataFiles(table) .filter(Expressions.equal("date", "2020-08-18")) .option("target-file-size-bytes", Long.toString(500 * 1024 * 1024)) // 500 MB .execute();files元数据表非常适合用于检查数据文件大小、判断何时需要压缩分区。
结合 api/src/main/java/org/apache/iceberg/actions/RewriteDataFiles.java 接口,常用选项如下:
| 选项 | 默认值 | 说明 |
|---|---|---|
target-file-size-bytes | 表属性write.target-file-size-bytes(默认 512 MB,见 TableProperties.java) | 重写目标输出文件大小 |
max-file-group-size-bytes | 100 GB | 单个文件组(file group)最多处理的数据量,用于把超大分区拆成可独立执行的小组 |
max-concurrent-file-group-rewrites | 5 | 同时重写的文件组数量 |
partial-progress.enabled | false | 是否允许在整体重写完成前先提交已完成的文件组(会产生多次提交,默认单次提交) |
partial-progress.max-commits | 10 | 启用 partial progress 时最多允许产生的提交数 |
rewrite-job-order | none | 作业组执行顺序:bytes-asc/bytes-desc/files-asc/files-desc/none |
use-starting-sequence-number | true | 新数据文件使用压缩开始时的快照序号,避免与更高序号的新增 equality delete 产生提交冲突 |
remove-dangling-deletes | false | 压缩后从当前快照移除不适用于任何存活数据文件的悬挂删除文件(equality 与 position 类型都会移除) |
output-spec-id | 当前表分区规格 | 重写输出使用的分区规格 ID,可借此重组数据布局 |
该 Action 还支持多种重写策略:binPack()(默认装箱策略)、sort()(按表排序顺序或自定义SortOrder重排)、zOrder(String... columns)与hilbert(String... columns)(多维空间填充曲线重排)。
重写 Manifest(Rewrite manifests)
Iceberg 使用 Manifest List 与 Manifest 文件中的元数据加速查询规划、裁剪不必要的数据文件。元数据树本质上是对表数据的一层索引。
Manifest 在元数据树中会按照被添加的顺序自动压缩,因此当写入模式与读取过滤条件一致时查询会更快。例如,按小时分区写入的数据,配合时间范围查询过滤条件就是对齐的。
当表的写入模式与查询模式不一致时,可以通过rewriteManifests(单机版)或rewriteManifestsAction(Spark 并行版)重写元数据,将数据文件重新分组到 Manifest 中。
以下示例重写小 Manifest,并按第一个分区字段分组数据文件:
Table table = ...; SparkActions .get() .rewriteManifests(table) .rewriteIf(file -> file.length() < 10 * 1024 * 1024) // 10 MB .execute();依据 api/src/main/java/org/apache/iceberg/actions/RewriteManifests.java,还有以下配置可组合使用:
specId(int):仅重写指定分区规格 ID 的 Manifest,默认使用表默认规格;rewriteIf(Predicate<ManifestFile>):只重写满足谓词的 Manifest,默认重写全部;sortBy(List<String> partitionFields):按指定的分区字段(需为转换后的列名,例如bucket(N, data)应传data_bucket而非data)排序重写的 Manifest,使 manifest_list 中的 manifest 指向包含相近分区值的数据文件,可显著减少规划时跳过的 Manifest 数量;stagingLocation(String):指定暂存 Manifest 的写入位置,默认写到表的元数据目录。
重写 Position Delete 文件(Rewrite position delete files)
Iceberg 可以重写 position delete 文件,这有两个目的:
- 小规模压缩(Minor compaction):把小的 position delete 文件合并为较大的文件,减少 Manifest 中存储的元数据体积与打开小 delete 文件的开销;
- 过滤悬挂记录(Filter dangling records):当某个 position delete 文件被选中重写时,丢弃其中引用已不再存活(live)数据文件的记录。
Table table = ...; SparkActions .get() .rewritePositionDeletes(table) .execute();只有被选中重写的 position delete 文件会受影响。如需仅从元数据中移除被判定为悬挂的整个 delete 文件,请使用removeDanglingDeleteFiles。该 Action 在 Spark SQL 中也有对应的rewrite_position_delete_files存储过程,例如:
CALL catalog_name.system.rewrite_position_delete_files(table => 'db.sample', options => map('min-input-files','2'));移除悬挂的删除文件(Remove dangling delete files)
如果某个 delete 文件的删除操作不再作用于任何存活数据文件,它就是“悬挂的”(dangling)。removeDanglingDeleteFilesAction 会扫描当前快照,并将可判定为悬挂的整个 delete 文件从元数据中移除:
- 分区中没有存活数据文件的 delete 文件;
- 数据序号(data sequence number)小于同分区内任何数据文件的 position delete 文件;其引用的数据文件不再存活时,删除向量(deletion vector)也会被移除;
- 数据序号小于或等于同分区内任何数据文件的 equality delete 文件。
这是一个仅元数据的操作:悬挂 delete 文件从表元数据中被丢弃,不重写任何数据或 delete 文件。该 Action 移除的是整个 delete 文件的引用;它不会从同时包含有效记录的 delete 文件中过滤悬挂记录。对于未分区表,该 Action 是空操作(no-op),因为悬挂删除已在每次提交时全表清理。
Table table = ...; SparkActions .get() .removeDanglingDeleteFiles(table) .execute();Spark SQL 没有对应的存储过程。在压缩时,rewriteDataFiles在remove-dangling-deletes为true时也可以移除悬挂 delete 文件;rewritePositionDeletes同样会在其重写的 position delete 文件中过滤悬挂记录。
计算表统计信息(Compute table statistics)
computeTableStatsAction 收集表各列的 Distinct Values 数量(NDV)统计信息,写入 Puffin 统计文件 并在表元数据中注册。查询引擎可以利用这些统计信息进行基于代价的优化(cost-based optimization)。
默认情况下,统计信息基于表当前快照、对所有顶层原始类型列收集。该 Action 可以配置为使用指定快照和/或列子集:
Table table = ...; SparkActions .get() .computeTableStats(table) .columns("col1", "col2") .execute();Spark SQL 中对应compute_table_stats存储过程,例如:
CALL catalog_name.system.compute_table_stats(table => 'my_table', snapshot_id => 'snap1', columns => array('col1', 'col2'));计算分区统计信息(Compute partition statistics)
computePartitionStatsAction 为分区表计算分区统计信息,并将生成的分区统计文件注册到表元数据中。统计信息从最近一个包含分区统计文件的快照开始,增量计算到所选快照(默认当前快照);如果不存在任何历史分区统计文件,则执行一次全量计算。
Table table = ...; SparkActions .get() .computePartitionStats(table) .execute();Spark SQL 中对应compute_partition_stats存储过程。
其他 Action(Other Actions)
以下 Action 不属于日常表维护范畴,但在删除表后的存储清理与表迁移复制场景中非常有用。
删除可达文件(Delete reachable files)
deleteReachableFilesAction 会删除表元数据文件引用的所有文件:数据文件、delete 文件、Manifest、Manifest List 与元数据文件。它用于在禁用 purge 删除表之后清理底层存储(例如 catalog 自身无法删除文件时)。
String metadataLocation = ...; // 表的 metadata.json 文件路径 FileIO io = ...; // 能读取并删除该表文件的 FileIO SparkActions .get() .deleteReachableFiles(metadataLocation) .io(io) .execute();注意:该 Action 接收的是metadata.json文件的路径而不是Table实例,因为它设计用于表已从 catalog 中删除之后。
危险:此操作会不可逆地删除表的所有可达文件。仅在表已被删除且数据不再需要时使用。如果其他表与该表共享文件(例如由
snapshotAction 或过程创建的表),其数据将被破坏。
重写表路径(Rewrite table path)
rewriteTablePathAction 会暂存一份表的元数据文件副本,其中所有以源前缀开头的绝对路径都被替换为目标前缀。它可作为将表完整或增量复制到新位置的起点,例如用于灾难恢复。
Table table = ...; SparkActions .get() .rewriteTablePath(table) .rewriteLocationPrefix("s3://bucket/old-table-location", "s3://bucket/new-table-location") .execute();该 Action 返回最新的已重写metadata.json名称,以及一个包含所有待复制文件源路径与目标路径的文件清单位置。它会将更新过路径的元数据与 position delete 文件写入暂存目录,但不会复制任何文件到目标位置。需要另用文件复制工具按计划中的清单复制文件。
Spark SQL 中对应rewrite_table_path存储过程。
维护操作的执行入口与调度建议
本仓库中,所有 Spark 版维护 Action 都通过 SparkActions 统一暴露,包括expireSnapshots、deleteOrphanFiles、rewriteDataFiles、rewriteManifests、rewritePositionDeletes、removeDanglingDeleteFiles、computeTableStats、computePartitionStats、deleteReachableFiles、rewriteTablePath等,调用方式均为SparkActions.get().<action>(table).<配置>().execute(),与本文各示例一致。
在实际生产环境中,建议将维护任务纳入定期调度(如每天执行快照过期与孤儿文件清理,按需执行数据文件压缩),并遵循以下原则:
- 先规划、后执行:多数 Action 支持
plan()先行产出计划,确认无误后再execute(); - 孤儿文件清理保留足够窗口:保留间隔必须大于最长写入任务的完成时间,默认 3 天不要随意调短;
- 监控元数据体积与文件分布:利用
files、manifests、snapshots等元数据表观察压缩需求与清理效果; - 高危操作严格隔离:
deleteReachableFiles只用于已 drop 且不再需要数据的表;修改 HDFS authority 等路径变化场景下,务必核对路径一致性再运行孤儿文件删除。
小结
Iceberg 的维护体系围绕“快照版本化 + 元数据树”这一核心设计展开:expireSnapshots释放过期快照引用的文件,deleteOrphanFiles清理元数据之外的残留文件,write.metadata.delete-after-commit.enabled控制提交时的元数据自动回收;而rewriteDataFiles、rewriteManifests、rewritePositionDeletes等优化操作则针对查询性能与元数据开销进行主动调优。将例行维护与可选优化结合,配合 Spark Actions 的并行能力,即可让 Iceberg 表在长时间、高频写入下保持稳定、紧凑且查询高效的状态。
- 数据湖
- 大数据
- 数据存储
【免费下载链接】iceberg
Apache Iceberg
相关推荐
Apache Iceberg Spark 存储过程(Procedures)完全指南:快照管理、数据维护与表迁移实战
Apache Iceberg Spark 存储过程(Procedures)完全指南:快照管理、数据维护与表迁移实战 Apache Iceberg 为 Spark
数据湖大数据数据存储StarRocks Iceberg Catalog Procedures 完整指南:快照管理、数据维护与元数据运维
StarRocks Iceberg Catalog Procedures 完整指南:快照管理、数据维护与元数据运维 StarRocks 的 Iceberg Ca
数据库OLAP数据仓库大数据湖仓一体数据分析Apache Iceberg表维护终极指南:7个快照清理与文件优化技巧
在大数据时代, Apache Iceberg 作为开源的大数据表格式,正在彻底改变数据湖的管理方式。但你知道吗?如果不进行定期维护,Iceberg表的性能会随着
数据湖大数据数据存储
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考