Azure Databricks中PySpark过滤后show/collect显示数据异常问题
Azure Databricks过滤结果不一致问题排查与解决
问题场景
加载一张含40+列、700万行的Delta表为PySpark DataFrame,测试分区/版本策略时出现异常:
- 通过UDF生成随机datetime类型
load_date列; - 提取
load_date的年、月为整数列load_year、load_month; - 用
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
相关产品推荐
相关产品推荐

