You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.08 01:45:35