ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Iceberg 表维护完全指南:快照过期、孤儿文件清理与数据/元数据重写实战

Apache Iceberg 表维护完全指南:快照过期、孤儿文件清理与数据/元数据重写实战 数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/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_ONLYdeleteWith(ConsumerString)自定义删除实现例如将待删文件收集起来而不实际删除executeDeleteWith(ExecutorService)/planWith(ExecutorService)并行执行删除或规划cleanExpiredMetadata(boolean)清理未使用的分区规格partition specs、Schema 等元数据提交时这些变更会应用到最新的表元数据发生提交冲突时Iceberg 会将变更重新应用到新的最新元数据并重试提交。过期操作也会删除不再被有效快照使用的 Manifest 文件以及被过期快照删除掉的数据文件。注意过期操作不允许删除当前快照。对于大表还可以使用 Spark Action 并行执行过期源码入口见 SparkActions.expireSnapshotsTable table ...; SparkActions .get() .expireSnapshots(table) .expireOlderThan(tsToExpire) .execute();重要过期快照只是把它们从元数据中移除使它们不再可用于时间旅行。数据文件只有在不再被任何可能用于时间旅行或回滚的快照引用时才会真正被删除。因此定期过期快照是释放底层数据文件的前提。删除旧元数据文件Remove old metadata filesIceberg 使用 JSON 文件跟踪表元数据。表的每次变更都会产生一个新的元数据文件以提供原子性。默认情况下旧元数据文件会被保留以支撑历史回溯对提交频繁的表如流式写入的表需要定期清理元数据文件。每个元数据文件会在metadata-log字段中跟踪更旧的元数据文件。被跟踪的元数据文件数量由write.metadata.previous-versions-max定义。要在提交后自动删除更旧的元数据文件请在表属性中设置write.metadata.delete-after-commit.enabledtrue。这会保留若干被跟踪的元数据文件最多write.metadata.previous-versions-max个并且每次创建新元数据文件时删除最旧的一个。注意该机制只会删除被metadata-log跟踪的元数据文件不会删除孤儿元数据文件。未被跟踪的元数据文件需要通过孤儿文件删除流程清理。属性定义可对照 core/src/main/java/org/apache/iceberg/TableProperties.java 中的常量与其默认值属性默认值描述write.metadata.delete-after-commit.enabledfalse控制每次表提交后是否删除最旧的被跟踪版本元数据文件write.metadata.previous-versions-max100最多跟踪的上一版本元数据文件数量示例说明设write.metadata.delete-after-commit.enabledfalse、write.metadata.previous-versions-max10提交 100 次后会有 10 个被跟踪的元数据文件与 90 个孤儿元数据文件。这 90 个孤儿文件无法通过设置write.metadata.delete-after-commit.enabledtrue删除因为它们已不再被跟踪只能通过孤儿文件删除流程清理。设write.metadata.delete-after-commit.enabledtrue、write.metadata.previous-versions-max20提交 21 次后会有 20 个被跟踪的元数据文件最旧的元数据文件在写入者提交时被删除。此后每次新提交都会删除最旧的一个元数据文件。更多写入属性请参见 表写入属性配置。删除孤儿文件Delete orphan files在 Spark 及其他分布式处理引擎中任务或作业失败可能留下未被表元数据引用的文件某些情况下普通的快照过期也无法判定某个文件不再需要并将其删除。要清理表目录下的这些“孤儿”文件可使用deleteOrphanFilesActionTable table ...; SparkActions .get() .deleteOrphanFiles(table) .execute();从 api/src/main/java/org/apache/iceberg/actions/DeleteOrphanFiles.java 的接口定义看该 Action 的核心语义是一个文件如果无法被任何有效快照到达即视为孤儿判定方式是对底层存储进行目录列举listing因此该操作代价较高。常用配置项配置默认说明location(String)表根目录指定扫描孤儿文件的目录不设置则扫描整个表目录可能同时删除孤儿数据与元数据文件olderThan(long)3 天前的时间戳只删除早于该时间戳的孤儿文件避免误删正在写入、尚未被元数据引用的文件deleteWith(ConsumerString)表的 FileIO自定义删除实现例如只收集不删除prefixMismatchModeERROR当元数据引用的文件与列举到的文件在 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 天除非必要不要调短。安全警告 2Iceberg 使用路径的字符串表示来判断哪些文件需要删除。在某些文件系统上路径可能随时间变化但仍指向同一文件。例如更换 HDFS 集群的 authorities 后创建时使用的旧路径 URL 与当前列举结果不再匹配运行 RemoveOrphanFiles 会导致数据丢失。请务必确保 MetadataTables 中的条目与 Hadoop FileSystem API 列举的条目一致以免误删。可选维护Optional Maintenance部分表需要额外维护。例如流式查询可能产生大量小文件应压缩为更大的数据文件某些表可以从重写 Manifest 文件中获益让查询定位数据更快。压缩数据文件Compact data filesIceberg 跟踪表中的每个数据文件。数据文件越多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-bytes100 GB单个文件组file group最多处理的数据量用于把超大分区拆成可独立执行的小组max-concurrent-file-group-rewrites5同时重写的文件组数量partial-progress.enabledfalse是否允许在整体重写完成前先提交已完成的文件组会产生多次提交默认单次提交partial-progress.max-commits10启用 partial progress 时最多允许产生的提交数rewrite-job-ordernone作业组执行顺序bytes-asc/bytes-desc/files-asc/files-desc/noneuse-starting-sequence-numbertrue新数据文件使用压缩开始时的快照序号避免与更高序号的新增 equality delete 产生提交冲突remove-dangling-deletesfalse压缩后从当前快照移除不适用于任何存活数据文件的悬挂删除文件equality 与 position 类型都会移除output-spec-id当前表分区规格重写输出使用的分区规格 ID可借此重组数据布局该 Action 还支持多种重写策略binPack()默认装箱策略、sort()按表排序顺序或自定义SortOrder重排、zOrder(String... columns)与hilbert(String... columns)多维空间填充曲线重排。重写 ManifestRewrite manifestsIceberg 使用 Manifest List 与 Manifest 文件中的元数据加速查询规划、裁剪不必要的数据文件。元数据树本质上是对表数据的一层索引。Manifest 在元数据树中会按照被添加的顺序自动压缩因此当写入模式与读取过滤条件一致时查询会更快。例如按小时分区写入的数据配合时间范围查询过滤条件就是对齐的。当表的写入模式与查询模式不一致时可以通过rewriteManifests单机版或rewriteManifestsActionSpark 并行版重写元数据将数据文件重新分组到 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(PredicateManifestFile)只重写满足谓词的 Manifest默认重写全部sortBy(ListString partitionFields)按指定的分区字段需为转换后的列名例如bucket(N, data)应传data_bucket而非data排序重写的 Manifest使 manifest_list 中的 manifest 指向包含相近分区值的数据文件可显著减少规划时跳过的 Manifest 数量stagingLocation(String)指定暂存 Manifest 的写入位置默认写到表的元数据目录。重写 Position Delete 文件Rewrite position delete filesIceberg 可以重写 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 statisticscomputeTableStatsAction 收集表各列的 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 statisticscomputePartitionStatsAction 为分区表计算分区统计信息并将生成的分区统计文件注册到表元数据中。统计信息从最近一个包含分区统计文件的快照开始增量计算到所选快照默认当前快照如果不存在任何历史分区统计文件则执行一次全量计算。Table table ...; SparkActions .get() .computePartitionStats(table) .execute();Spark SQL 中对应compute_partition_stats存储过程。其他 ActionOther Actions以下 Action 不属于日常表维护范畴但在删除表后的存储清理与表迁移复制场景中非常有用。删除可达文件Delete reachable filesdeleteReachableFilesAction 会删除表元数据文件引用的所有文件数据文件、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 pathrewriteTablePathAction 会暂存一份表的元数据文件副本其中所有以源前缀开头的绝对路径都被替换为目标前缀。它可作为将表完整或增量复制到新位置的起点例如用于灾难恢复。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 表在长时间、高频写入下保持稳定、紧凑且查询高效的状态。赞分享数据湖大数据数据存储【免费下载链接】icebergApache Iceberg项目地址https://gitcode.com/gh_mirrors/icebe/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表的性能会随着数据湖大数据数据存储上一篇WePY 小程序组件化开发框架完整指南预编译原理、类 Vue 开发与快速上手下一篇G-Helper终极指南免费轻量级华硕笔记本控制工具完全教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进