Spark高阶函数reduce未识别报错求助
问题原因与解决方案
原因分析
Spark SQL的reduce高阶函数是Spark 3.1.0及以上版本才新增支持的,你当前使用的Spark版本低于3.1.0,因此会报函数未定义错误。而transform、filter、exists这类高阶函数在更早的Spark版本(如Spark 2.4+)就已支持,所以能正常运行。
解决方案
方案1:升级Spark版本
将Spark升级到3.1.0或更高版本,你的原有代码即可直接正常运行。
方案2:用aggregate函数替代(无需升级)
如果无法升级Spark版本,可以使用aggregate函数实现相同逻辑,该函数在Spark 2.4及以上版本支持:
spark.sql("SELECT" + " celsius," + " aggregate(celsius, (sum_acc => 0, count_acc => 0)," + " (acc, t) => (acc.sum_acc + t, acc.count_acc + 1)," + " acc => (acc.sum_acc div acc.count_acc * 9 div 5) + 32 ) as avgFahrenheit" + " FROM celsiusView")
逻辑说明:
- 初始值
(sum_acc => 0, count_acc => 0):定义两个累加器,分别存储温度总和与数组元素个数 - 第一个匿名函数:遍历数组时,同步累加温度总和与元素计数
- 第二个匿名函数:用总和除以计数得到平均摄氏度,再转换为华氏度
内容的提问来源于stack exchange,提问作者Mandroid
相关产品推荐
相关产品推荐

