如何判断Iceberg查询仅使用元数据/清单文件还是访问数据文件
我创建了如下Iceberg表:
CREATE TABLE IF NOT EXISTS raw_data.civ ( date timestamp, marketplace_id int, ... some more columns ) USING ICEBERG PARTITIONED BY ( marketplace_id, days(date) ) TBLPROPERTIES ( 'write.sort-order' = 'dimension' )
之后将date列重命名为snapshot_day。
我执行了如下查询:
civIbDf.select("snapshot_day").distinct().orderBy(col("snapshot_day").desc).show(100)
该表的snapshot_day约有150个唯一值,但每个值对应大量数据。我认为该查询只需读取清单文件即可获取分区键的唯一值,应在几秒内完成,但在EMR Notebook中耗时2-3分钟,在AWS Athena中耗时超20分钟。
我的spark.sql.catalog.my_catalog.warehouse存储在S3,metadata文件夹仅约112个文件。
以下是查询的执行计划:
civIbDf.select("snapshot_day").distinct().orderBy(col("snapshot_day").desc).explain() == Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- Sort [snapshot_day#12 DESC NULLS LAST], true, 0 +- Exchange rangepartitioning(snapshot_day#12 DESC NULLS LAST, 4800), ENSURE_REQUIREMENTS, [plan_id=244] +- HashAggregate(keys=[snapshot_day#12], functions=[], schema specialized) +- Exchange hashpartitioning(snapshot_day#12, 4800), ENSURE_REQUIREMENTS, [plan_id=241] +- HashAggregate(keys=[snapshot_day#12], functions=[], schema specialized) +- BatchScan my_catalog.raw_data.civ[snapshot_day#12] my_catalog.raw_data.civ (branch=null) [filters=, groupedBy=, pushedLimit=None] RuntimeFilters: []
请问如何确认该查询是否仅使用清单文件和元数据,而未处理数据文件?因为Iceberg的优势之一就是此类基础查询可仅通过元数据返回结果。
一、验证查询是否仅读取元数据/清单文件
1. 检查Spark执行计划标识
当前执行计划显示的是BatchScan,说明Spark正在扫描数据文件而非仅读取元数据。如果查询触发了元数据优化,执行计划中会出现MetadataOnlyScan(不同Iceberg版本可能有细微差异)。
2. 查看Spark日志详情
在EMR Notebook中调整Spark日志级别为INFO,搜索日志内容:
- 若仅出现读取
metadata.json、snapshot-*.json、manifest-*.avro这类元数据/清单文件的记录,说明查询仅依赖元数据; - 若日志中存在大量读取
.parquet/.orc等数据格式文件的记录,证明查询正在处理数据文件。
3. 分析S3访问日志
查看S3存储桶的访问日志,统计查询执行期间的文件访问类型:
- 仅访问
metadata目录及manifest相关文件 → 未读取数据文件; - 存在大量
data目录下文件的访问请求 → 正在扫描数据文件。
二、查询未触发元数据优化的原因
你的表按days(date)分区,后续将date重命名为snapshot_day,但分区定义仍基于原date列的转换结果。Spark Iceberg未识别出查询的snapshot_day与分区键days(date)的关联,因此无法直接从元数据中提取唯一值。
三、优化查询触发元数据读取
1. 直接查询分区列
通过DESCRIBE EXTENDED raw_data.civ查看实际分区列名(通常会以days_date这类格式命名),直接查询该分区列:
civIbDf.select("days_date").distinct().orderBy(col("days_date").desc).show(100)
或者直接使用分区转换函数:
civIbDf.select("days(snapshot_day)").distinct().orderBy(col("days(snapshot_day)").desc).show(100)
2. 更新表分区定义
修改表的分区规则,让分区键与重命名后的列直接关联:
ALTER TABLE raw_data.civ REPLACE PARTITIONING BY (marketplace_id, days(snapshot_day))
3. 启用Iceberg元数据优化配置
确保Spark配置中开启了元数据扫描优化:
spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions spark.sql.catalog.my_catalog=org.apache.iceberg.spark.SparkCatalog spark.sql.catalog.my_catalog.type=hadoop spark.sql.catalog.my_catalog.warehouse=s3://your-bucket/path spark.sql.iceberg.metadata.cache.enabled=true spark.sql.iceberg.scan.metadata-only.enabled=true
内容的提问来源于stack exchange,提问作者shiva

