Iceberg expireSnapshots与deleteOrphanFiles API失效问题排查
问题描述
环境配置(POM依赖)
<!-- spark --> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.12</artifactId> <version>3.1.1</version> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.12</artifactId> <version>3.1.1</version> </dependency> <!-- iceberg --> <dependency> <groupId>org.apache.iceberg</groupId> <artifactId>iceberg-spark-runtime-3.1_2.12</artifactId> <version>1.2.0</version> </dependency>
表目录结构
$ tree db db └── table_name ├── data │ └── year=2024 │ └── month=04 │ └── day=07 │ └── dataType=PD │ └── 00000-0-04211945-446e-4d17-83c4-7e8447f5e4e4-00001.parquet └── metadata ├── d2d68629-eb2f-4582-b2ed-ebb17f923f72-m0.avro ├── snap-8662509695078634239-1-d2d68629-eb2f-4582-b2ed-ebb17f923f72.avro ├── v1.metadata.json ├── v2.metadata.json └── version-hint.text
注:元数据目录下的两个metadata文件是先创建表再插入分区数据导致。
测试代码
import scala.collection.JavaConverters.iterableAsScalaIterableConverter import scala.concurrent.duration.DurationInt import org.apache.hadoop.conf.Configuration import org.apache.iceberg.catalog.TableIdentifier import org.apache.iceberg.hadoop.HadoopCatalog import org.apache.iceberg.spark.actions.SparkActions import org.apache.spark.sql.SparkSession object ExpireSnapshotDemo { def main(args: Array[String]): Unit = { @transient implicit val spark = SparkSession .builder() .master("local[*]") .config("spark.driver.bindAddress", "127.0.0.1") .appName("IcebergTableCreationExample") .config( "spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions" ) .config( "spark.sql.catalog.local", "org.apache.iceberg.spark.SparkCatalog" ) .config("spark.sql.catalog.local.type", "hadoop") .config( "spark.sql.catalog.local.warehouse", "/Users/Code/spark/iceberg" ) .getOrCreate() val catalog = new HadoopCatalog( new Configuration(), "/Users/Code/spark/iceberg" ) val tableIdentifier = TableIdentifier.of("db", "table_name") val table = catalog.loadTable(tableIdentifier) println(table.snapshots().asScala.map(_.snapshotId()).toList) val expireTimeMilliseconds = System.currentTimeMillis() - 2.hours.toMillis println(expireTimeMilliseconds) SparkActions .get(spark) .expireSnapshots(table) .expireOlderThan(expireTimeMilliseconds) .execute() table.refresh() println(table.snapshots().asScala.map(_.snapshotId()).toList) SparkActions .get(spark) .deleteOrphanFiles(table) .olderThan(expireTimeMilliseconds) .execute() table.refresh() println(table.snapshots().asScala.map(_.snapshotId()).toList) } }
预期与实际结果
- 预期:首次打印1个快照ID;调用
expireSnapshots移除元数据中的快照;调用deleteOrphanFiles删除数据文件;最终快照ID列表为空,data目录文件被删除。 - 实际:数据文件仍存在,快照ID依然可查询。已确认无其他Spark作业、分支或标签使用该表,插入数据已过7小时。
原因分析与解决方案
核心原因1:Iceberg快照保留默认规则
Iceberg的expireSnapshots默认强制保留至少1个快照,即使该快照早于过期时间。这是为了避免表进入无快照的不可用状态。当前场景中只有1个快照,所以即使调用expireOlderThan,该快照也不会被移除。
核心原因2:孤儿文件判定逻辑
deleteOrphanFiles仅删除无任何快照引用的文件。由于第一步快照未被移除,数据文件仍被当前快照引用,因此不会被判定为孤儿文件,也就不会被删除。
解决方案
1. 强制过期所有快照(需谨慎,会导致表无快照)
添加retainLast(0)配置,覆盖默认保留1个快照的规则:
SparkActions .get(spark) .expireSnapshots(table) .expireOlderThan(expireTimeMilliseconds) .retainLast(0) // 允许保留0个快照 .execute()
2. 验证快照过期状态
调用expireSnapshots后,可通过以下代码确认快照是否被标记为过期:
println(table.snapshots().asScala.filter(_.isExpired).map(_.snapshotId()).toList)
3. 调整deleteOrphanFiles参数(可选)
若需确保删除所有未被引用的文件,可去掉olderThan参数,扫描所有孤儿文件:
SparkActions .get(spark) .deleteOrphanFiles(table) .execute()
4. 使用表API替代Spark Actions(等效方案)
直接使用Iceberg表API执行快照过期,效果一致:
table.expireSnapshots() .expireOlderThan(expireTimeMilliseconds) .retainLast(0) .commit()
内容的提问来源于stack exchange,提问作者Yuqi Zhang
相关产品推荐
相关产品推荐

