适用版本与隔离前提
本文以 Iceberg 1.8.1、Spark 3.5 和 format v2 表为基线,讨论写入提交、小文件重写、快照保留与 orphan 检查。使用 CALL 过程需要匹配的 Iceberg Spark 运行时、已配置 catalog 和 Iceberg SQL 扩展;Spark 3.5 的包不能用于任意 Spark 或 Scala 版本。配置依据 Iceberg Spark Configuration。
所有可修改示例只用于独立 catalog 或 namespace 内的专用测试表,其存储前缀不与生产共享,总数据不超过一千行、二十个数据文件。本文没有运行示例或清理文件,验证状态为 PENDING。先记录 catalog、表 UUID 或受控身份、当前 snapshot、写入应用及数据保留要求。
文件落盘与提交成功是两件事
Spark Task 可能先生成数据文件,随后才由写入流程完成表提交;失败或冲突会留下需要后续判断的文件。客户端收到超时,不等于提交必然失败,也不意味着可以直接再次 append。先通过快照、应用提交审计和业务批次键确认是否已经提交,再决定重试。
Iceberg 的并发提交依赖乐观并发与冲突验证,不是让多个写入者绕过同一 catalog 自行修改文件。避免两个不同 catalog 注册同一存储目录后独立写入。更多原理见 Iceberg 1.8.1 Reliability 和 表元数据与快照。
先明确 append、overwrite 与业务幂等
追加写会增加记录,不自动按业务主键去重。覆盖写要明确被替换的数据范围,动态分区覆盖与静态覆盖的含义不同,分区规范变化也可能改变操作范围。Spark DataFrameWriterV2 的 writeTo API 更便于明确按名称写入及操作类型,但它不能替业务定义重复批次的处理规则。
以输入批次清单、唯一事件键及目标快照共同判定重复或缺失。涉及 MERGE、DELETE 或多个输出系统时,核对对应版本、SQL 扩展与事务边界,不能把一张 Iceberg 表的原子提交推导为多个表或外部系统的整体事务。写入语义见 Iceberg 1.8.1 Spark Writes。
用元数据定位小文件成本
在前述小型测试表中查看文件清单和最近提交;元数据查询仍有规划与读取成本,不建议全库遍历:
SELECT content, file_path, record_count, file_size_in_bytes
FROM REPLACE_WITH_CATALOG.REPLACE_WITH_DB.REPLACE_WITH_TEST_TABLE.files
ORDER BY file_size_in_bytes LIMIT 20;
SELECT committed_at, snapshot_id, operation
FROM REPLACE_WITH_CATALOG.REPLACE_WITH_DB.REPLACE_WITH_TEST_TABLE.snapshots
ORDER BY committed_at DESC LIMIT 10;先确认视图范围与 content 含义,再解释文件数、大小和记录数。文件元数据记录数不一定等于应用可见净行数,尤其要考虑删除文件和历史视图。表文件少但计划慢时,还需检查 manifest 数量与分区信息。字段依据 Spark Queries。
在受控分区演练数据文件重写
示例只改写专用测试表中可能匹配某日的文件,分区列 event_day 为整数日期键,并限制文件组并发。目标 catalog 已启用过程支持:
CALL REPLACE_WITH_CATALOG.system.rewrite_data_files(
table => 'REPLACE_WITH_DB.REPLACE_WITH_TEST_TABLE',
where => 'event_day = 20260101',
options => map('min-input-files', '2',
'max-concurrent-file-group-rewrites', '1',
'partial-progress.enabled', 'false')
);where 选择可能包含匹配数据的文件,不是只改写那些匹配行;总工作量边界来自事先确认的小表规模。输出的 rewritten、added 文件计数描述重写结果,不能单独证明业务数据相同;没有符合条件的文件时可能没有变化。若启用 partial progress,部分文件组可能先提交,更需确认失败后的实际快照。依据 Iceberg 1.8.1 Procedures。
文件大小目标与 Spark 分区配合
目标文件大小不是强制填满承诺。Spark Task 数据量、压缩比和 Iceberg 分区边界都影响输出,一个文件不会跨越 Iceberg 分区。调整 write.target-file-size-bytes 却不改变过小输入任务,通常不能凭空产生大文件;也不能为追求大文件把所有输入压到单个 Task。
先记录写入 distribution、AQE 与 Task 数据量,再评估批次大小、分区基数和压缩。若采用 fanout,要考虑保持多个文件句柄带来的资源成本。只改变一个有证据支持的条件,对照写入延迟、文件分布和下游扫描,避免同时调整十多个参数。Spark 执行分析见 Spark SQL 性能诊断。
保留与 orphan 检查分开进行
快照过期处理历史快照及不再被保留快照需要的文件;orphan 检查寻找不受表元数据引用的对象。二者不同,不能把一次重写产生的旧文件立即当作 orphan 删除。时间旅行、长期读取、流消费者、分支标签和独立备份均要进入保留评审。
下面只对专用测试表预览候选,不删除。日期是显式演练截止点,必须在最长写入窗口之前,并与演练样本修改时间对应:
CALL REPLACE_WITH_CATALOG.system.remove_orphan_files(
table => 'REPLACE_WITH_DB.REPLACE_WITH_TEST_TABLE',
older_than => TIMESTAMP '2025-12-01 00:00:00',
dry_run => true
);dry_run 避免删除,不消除目录枚举成本;返回路径只是候选,空结果也不证明全部历史文件都已健康。路径 scheme、authority 或共享前缀不一致会影响判断,在途写入的文件可能暂时未提交。安全窗口与路径问题见 Iceberg Maintenance。
常见误区与验收方法
不要手工从存储中删除看起来较旧的 Parquet、manifest 或 metadata 文件;不要把删除旧快照当成可随时反悔的容量操作;也不要因 CALL 报错就宣告没有提交。失败结果需要结合操作日志、当前 snapshot 与 partial progress 状态解释。
以同一输入窗口比较重写前后可见记录、事件键、NULL、业务金额和日期范围,再检查文件分布、扫描成本与实际资源使用。文件压缩和排列变化后,二进制摘要可以不同;业务正确性与性能收益分别验收。预演输出应保存时间和存储清单,过期的候选报告不能直接作为未来清理依据。
停止条件与回退边界
若出现提交结果不明、共享目录、路径匹配异常或业务对账差异,停止进一步重写与清理。重写属于数据布局维护,若已提交,需要确认目标快照及后续并发写入,再评估受支持的恢复操作;不能直接把 catalog 指向猜测的旧 metadata JSON。
快照和文件已被清理后,回到旧快照可能已不可行,恢复依赖保留的完整引用链与独立数据副本。取消 Spark 作业不会撤销已提交的表变化;新业务写入也不能被一起回退掉。交接应记录前后 snapshot、操作范围、候选文件、保留条件和仍未验证的读写引擎。