You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

不使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 08:45:08