不使用zip_with实现Spark DataFrame数组按位置拼接的方案
Spark 3.x 按位置拼接两个数组的替代实现方案
无UDF的最简原生实现(无需zip_with)
Spark 2.4及以上版本可以直接用arrays_zip+transform的组合实现完全相同的效果,代码如下:
import org.apache.spark.sql.functions._ val resDF = df.withColumn( "jCols", transform( arrays_zip(col("prop1"), col("values")), elem => array(elem.getField("prop1").cast(DoubleType), elem.getField("values")) ) )
运行后输出和你原有的zip_with实现完全一致,比如输入[3, 2]和[0.0, 0.1]时,输出结果为[[3.0, 0.0], [2.0, 0.1]]。
是否需要使用UDF?
不需要,除非你需要兼容Spark 2.4以下的极旧版本。原生函数已经经过Catalyst优化,性能远高于自定义UDF。如果有兼容需求可以用如下UDF实现:
val zipMergeUDF = udf((arr1: Seq[Int], arr2: Seq[Double]) => { arr1.zip(arr2).map(t => Seq(t._1.toDouble, t._2)) }) val resDF = df.withColumn("jCols", zipMergeUDF(col("prop1"), col("values")))
底层实现逻辑
zip_with的底层会先对齐两个输入数组的长度(默认按较短数组长度截断),然后按索引遍历数组元素,对每一组同位置元素执行传入的表达式逻辑,最终汇总为新数组返回。arrays_zip+transform的组合底层逻辑和zip_with几乎等价:arrays_zip负责将两个数组的同位置元素封装为结构体,transform负责遍历结构体数组,将每个结构体转换为目标子数组,两者的性能差异可以忽略。
内容的提问来源于stack exchange,提问作者Ged
相关产品推荐
相关产品推荐

