Databricks中unionAll已持久化数据集后卡“Filtering files for query”求助
解决PySpark unionAll后卡在“Filtering files for query”的问题
可能的原因
- 小文件泛滥:即便执行了
count()触发持久化,若用默认的MEMORY_AND_DISK存储,磁盘上可能生成大量小分区文件。unionAll时需要扫描所有小文件的元数据,这会显著拖慢文件过滤阶段。 - 分区数过多:三个数据集的分区总数过大(比如每个几百上千分区),union时要遍历每个分区的文件元数据,累加起来耗时剧增。
- 元数据缓存失效或未命中:集群的文件系统元数据缓存没更新,或者持久化的数据集缓存未被执行计划识别,导致需要重新从底层存储读取文件信息,而非直接用内存缓存。
- 持久化未完全完成:
count()返回结果不代表磁盘写入已经结束,若后台还在写数据,union操作会等待写入完成,表现为过滤阶段卡住。 - 旧API的执行计划问题:
unionAll是PySpark旧API,部分场景下执行计划优化器可能未正确识别持久化缓存,仍会触发全量文件扫描。
解决方法
优化持久化存储格式与文件大小
改用Delta或Parquet格式存储,并合并小文件,这类格式的元数据查询效率更高:# 替换原持久化方式,用Delta表存储 df1.repartition(20).write.mode("overwrite").format("delta").saveAsTable("temp_df1") df2.repartition(20).write.mode("overwrite").format("delta").saveAsTable("temp_df2") df3.repartition(20).write.mode("overwrite").format("delta").saveAsTable("temp_df3") # 读取Delta表进行合并 combined_df = spark.table("temp_df1").union(spark.table("temp_df2")).union(spark.table("temp_df3"))Delta格式会自动管理文件大小,避免小文件问题,同时元数据查询更高效。
减少分区数量
合并前对每个数据集执行coalesce减少分区数,降低文件扫描的总量:# 根据数据量调整分区数,比如合并到5个分区 df1 = df1.coalesce(5) df2 = df2.coalesce(5) df3 = df3.coalesce(5) combined_df = df1.union(df2).union(df3)确认缓存状态并刷新元数据
清理旧缓存后重新持久化,确保数据集被正确缓存:from pyspark.storagelevel import StorageLevel spark.catalog.clearCache() # 重新持久化 df1.persist(StorageLevel.MEMORY_AND_DISK) df2.persist(StorageLevel.MEMORY_AND_DISK) df3.persist(StorageLevel.MEMORY_AND_DISK) # 触发缓存加载 df1.count() df2.count() df3.count() # 验证缓存是否生效 print(f"df1 cached: {spark.catalog.isCached(df1)}") print(f"df2 cached: {spark.catalog.isCached(df2)}") print(f"df3 cached: {spark.catalog.isCached(df3)}") # 再执行合并 combined_df = df1.union(df2).union(df3)检查执行计划并替换为新API
用union替代unionAll,并查看执行计划确认是否命中缓存:combined_df = df1.union(df2).union(df3) # 查看执行计划,若显示InMemoryTableScan则说明命中缓存 combined_df.explain()如果执行计划显示
FileScan,说明缓存未被使用,需要重新检查持久化步骤。排查集群资源状态
登录Databricks集群监控页面,查看工作节点的CPU、内存使用率:- 若内存不足,持久化数据可能被逐出到磁盘,导致union时需要读取大量磁盘小文件;此时可调整节点内存或增加节点数。
- 确认节点无故障,资源未被其他任务挤占。
内容的提问来源于stack exchange,提问作者João Paulo Andrade
相关产品推荐
相关产品推荐

