PySpark 3.5+Databricks 14.3 ML LTS同代码结果不一致问题咨询
问题1:代码运行结果不一致的原因与解决方法
核心原因
你的代码中窗口函数的orderBy("year_month")存在非确定性排序:当同一outlet_id下存在多个相同year_month的行时,Spark无法保证这些行的固定顺序,并行执行时的任务调度顺序或分区数据分布变化,会导致f.last("distributor_id", True)取到不同的值,最终影响distinct后的组合数统计结果。
解决方案
- 增加确定性排序键:在
orderBy中加入能唯一标识行的字段(如自增ID、时间戳等),确保同一outlet_id和year_month下的行顺序固定:
window_c = ( Window() .partitionBy("outlet_id") .orderBy("year_month", "unique_row_id") # 新增唯一排序键 .rowsBetween(Window.unboundedPreceding, Window.currentRow) )
- 重构逻辑为分组聚合:通过分组替代窗口函数,避免依赖行顺序:
# 先按outlet_id、year_month取最后一个distributor_id,再关联item_id计算唯一组合数 latest_distributor = ( master_keys .groupBy("outlet_id", "year_month") .agg(f.last("distributor_id", True).alias("distributor_id")) ) distinct_combination_count = ( master_keys .join(latest_distributor, on=["outlet_id", "year_month"], how="left") .select("item_id", "distributor_id") .distinct() .count() )
- 校验数据质量:检查
master_keys表中是否存在distributor_id频繁变更、null值异常或重复数据,这些也会导致结果波动。
问题2:Databricks 10.4 ML LTS与14.3 ML LTS的已知结果差异
两个版本基于不同Spark内核(10.4对应Spark 3.2.1,14.3对应Spark 3.5.0),核心差异包括:
- 排序与窗口函数行为:Spark 3.4+优化了排序稳定性,当orderBy键重复时,旧版本可能依赖物理存储顺序,新版本引入更明确的排序规则,导致窗口函数结果变化。
- 聚合函数null处理:部分聚合函数(如
count、sum)对null值的处理逻辑微调,比如count(*)与count(col)的边界场景。 - 数据格式读写:Parquet/ORC默认读写版本升级,旧版本写入的文件在新版本读取时,可能因兼容性问题导致字段值或分区解析差异。
- MLlib算法更新:部分模型(如XGBoost、Random Forest)的默认参数、训练逻辑或评估指标计算方式变更,迁移时需重新校验模型输出。
- 时区与时间处理:默认时区配置可能调整,涉及时间字段转换、分区裁剪的逻辑会出现结果差异。
- 分区与 shuffle 策略:Spark 3.5优化了shuffle分区分配和数据分布,依赖分区顺序的自定义逻辑(如UDF、自定义聚合)可能出现结果波动。
内容的提问来源于stack exchange,提问作者Kavishka Gamage
相关产品推荐
相关产品推荐

