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

Azure Databricks中PySpark过滤后show/collect显示数据异常问题

Azure Databricks过滤结果不一致问题排查与解决

问题场景

加载一张含40+列、700万行的Delta表为PySpark DataFrame,测试分区/版本策略时出现异常:

  1. 通过UDF生成随机datetime类型load_date列;
  2. 提取load_date的年、月为整数列load_year、load_month;
  3. 用load_date或load_year过滤数据时,count()结果符合预期,但show()/collect()/display()返回数据包含不符合过滤规则的行(如2021年数据出现在load_year>2022的结果中)。

已尝试重置笔记本、重启集群、更换浏览器,问题仍存在。技术环境:DB runtime 10.4 LTS(Apache Spark 3.2.1、Scala 2.12),启用Photon,使用Spot实例,Worker为Standard_DS3_v2。

排查与解决步骤

1. 修复非确定性UDF问题

自定义UDF若未标记为确定性,会导致Spark执行计划出现不一致——count()和show()可能触发不同计算路径,结果偏差。

  • 优先用Spark内置函数生成随机日期,确保逻辑可重复:
    from pyspark.sql.functions import expr, rand
    
    df = df.withColumn("load_date", expr(
        "date_add('2020-01-01', cast(rand() * 1460 as int))"
    ))
    
  • 必须用UDF时,标记为确定性并固定随机种子:
    from pyspark.sql.functions import udf
    from datetime import datetime
    import random
    
    @udf("timestamp", deterministic=True)
    def generate_random_date():
        random.seed(42)  # 固定种子保证结果一致
        start = datetime(2020, 1, 1)
        end = datetime(2023, 12, 31)
        return start + (end - start) * random.random()
    
    df = df.withColumn("load_date", generate_random_date())
    

2. 校验load_year/load_month生成逻辑

确保提取操作基于正确的类型转换,避免类型不匹配导致的过滤失效:

  • 用Spark内置函数替代自定义提取逻辑:
    from pyspark.sql.functions import year, month
    
    df = df.withColumn("load_year", year("load_date"))
           .withColumn("load_month", month("load_date"))
    
  • 验证列类型,确保为整数:
    df.select("load_year", "load_month").printSchema()
    
    若显示为字符串类型,需转换为整数后再过滤:
    from pyspark.sql.functions import col
    
    df = df.withColumn("load_year", col("load_year").cast("int"))
    

3. 修复Delta表元数据与数据不一致

Delta表分区元数据或旧版本快照可能导致查询结果异常:

  • 清理旧版本数据:
    VACUUM delta.`/path/to/your/delta/table` RETAIN 0 HOURS;
    
  • 优化表结构并重写分区元数据:
    OPTIMIZE delta.`/path/to/your/delta/table` ZORDER BY (load_year, load_month);
    
  • 加载表时指定最新版本:
    latest_version = spark.sql("DESCRIBE HISTORY delta.`/path/to/table`") \
        .select("version") \
        .orderBy("version", ascending=False) \
        .first()[0]
    
    df = spark.read.format("delta") \
        .option("versionAsOf", latest_version) \
        .load("/path/to/table")
    

4. 临时禁用Photon引擎排查兼容性

Photon引擎在特定场景下可能存在兼容性问题,临时关闭验证:

spark.conf.set("spark.databricks.photon.enabled", "false")

重新执行过滤和查询操作,观察结果是否一致。

5. 清除DataFrame缓存

若DataFrame被缓存,可能存在缓存数据与实际计算结果不一致的情况:

df.unpersist()  # 清除缓存
filtered_df = df.filter(df.load_year > 2022)
filtered_df.count()
filtered_df.show()

内容的提问来源于stack exchange,提问作者T. Doe

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 17:53:18