Scala Spark中foldLeft()与foreach()应用UDF时表现异常问题咨询
Spark两个列转换函数的差异及问题原因
核心问题
第二个函数所有列转换不生效的根本原因有两点:
- Spark的DataFrame是不可变对象,
withColumn不会修改原DataFrame,只会返回一个修改后的新DataFrame - 错误使用了无返回值的
foreach遍历列,所有withColumn返回的新DataFrame都被直接丢弃,最终返回的还是未修改的原始DataFrame
两个函数的运行逻辑差异
- 函数1逻辑:
内层采用foldLeft遍历待处理列,每一次调用withColumn生成的新DataFrame都会作为下一轮迭代的输入,所有列的修改会逐层叠加,最终返回的是叠加了所有列转换的DataFrame,因此转换可以正常生效。 - 函数2逻辑:
遍历待处理列时使用了无返回值的foreach算子,foreach内部调用withColumn生成的新DataFrame没有被任何变量接收,直接被丢弃,迭代结束后返回的仍然是最开始传入的未做任何修改的accumulator,因此看不到任何列转换效果。
函数2修复方案
把foreach遍历替换为foldLeft即可,修改后代码如下:
def transformColumns(df: DataFrame, columns: Map[Seq[String], TAlgorithm]): DataFrame = { try { columns.foldLeft(df) { (accumulator: DataFrame, sanitization: (Seq[String], TAlgorithm)) => import org.apache.spark.sql.functions.udf val aes: TAlgorithm = new AES256(key, iv) @transient lazy val udfFunction = udf(aes.decrypt(_)) // 替换foreach为foldLeft,叠加每次withColumn的结果 sanitization._1.foldLeft(accumulator) { (innerAcc, elem) => innerAcc.withColumn(elem, when(col(elem).isNotNull, udfFunction(col(elem))).otherwise(lit(null))) } } } }
额外优化建议
当前实现每处理一组列都会新建一个AES256实例,实际使用中可以将AES实例和UDF的初始化放到循环外,减少不必要的对象创建开销。
内容的提问来源于stack exchange,提问作者Aman Vaishya
相关产品推荐
相关产品推荐

