Spark DF转Pandas用Sklearn算MAE每次结果不同的解决方法
问题产生原因
结果不一致的核心是Spark的惰性执行机制与当前代码写法的共同影响,具体分为两点:
- 现有代码分别对
qty列、pred列单独调用toPandas(),相当于触发了两次独立的Spark作业。Spark默认不会缓存中间计算结果,每次触发action都会从头重算整个DataFrame的血缘链路。如果testDataFrame的上游存在任何非确定性逻辑——比如使用了rand()等随机函数、读取的是正在并发写入的表、未做全局排序的分区扫描、开启动态自适应优化导致分区策略动态调整,两次作业拉取到的数据行顺序、甚至数据内容本身都可能出现差异,最终导致传入sklearn的真实值和预测值行没有一一对应,计算结果自然每次都不同。 - 即使把两列合并到一次
toPandas()调用,如果testDataFrame没有做持久化,上游的非确定性逻辑依然会导致每次拉取的数据集存在差异,结果同样会波动。
可行解决方案
按落地优先级从高到低排列:
- 优先使用Spark原生内置函数直接计算指标,完全绕开
toPandas()和mllib依赖,性能最好,也不存在Driver端内存溢出风险,不需要额外调整集群配置:
from pyspark.sql import functions as F # 直接用Spark SQL聚合函数计算MAE、MAPE metric_row = test.select( F.avg(F.abs(F.col("qty") - F.col("pred"))).alias("mae"), F.avg(F.abs((F.col("qty") - F.col("pred")) / F.col("qty"))).alias("mape") ).first() mae = metric_row["mae"] mape = metric_row["mape"]
注意:如果qty列存在0值,计算MAPE时需要提前过滤对应行或者增加平滑项,避免出现除零错误。
- 如果必须使用sklearn的计算逻辑,必须保证仅触发一次
toPandas()动作,一次性拉取所有需要的列,从根源上避免两次作业数据不对齐的问题:
from sklearn.metrics import mean_absolute_error, mean_absolute_percentage_error # 单次查询同时取两列,仅触发一次action pdf = test.select("qty", "pred").toPandas() mae = mean_absolute_error(pdf["qty"], pdf["pred"]) mape = mean_absolute_percentage_error(pdf["qty"], pdf["pred"])
- 如果
testDataFrame上游确实存在非确定性逻辑,在计算前先对DataFrame做持久化固化结果,避免每次重算血缘导致数据变动:
# 根据数据量选择合适的存储级别,数据量不大直接用默认内存缓存即可 test.persist() # 触发一次action让缓存生效 test.count() # 后续再执行指标计算逻辑即可得到稳定结果
内容的提问来源于stack exchange,提问作者ranemak
相关产品推荐
相关产品推荐

