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

如何将Spark DataFrame两列转换为List[(Long,Int)]传入自定义UDF

解决方案

核心实现方法

你可以直接使用Spark内置的arrays_zip函数配对两个数组的同位置元素,Spark会自动完成类型映射转换为UDF需要的List[(Long, Int)]格式,具体操作步骤如下:

  1. 首先导入依赖的内置函数
import org.apache.spark.sql.functions.{arrays_zip, lit, size}
  1. 调用UDF时传入配对后的数组列即可,以下是示例写法:
// 可选:提前过滤两个数组长度不一致的脏数据,避免arrays_zip按短数组截断导致数据异常
val filteredDf = df.filter(size($"lane_path.coordinates") === size($"lane_path.diffs"))

val resultDf = filteredDf.withColumn("generated_result", generateValues(
  arrays_zip($"lane_path.coordinates", $"lane_path.diffs"),
  lit(100L), // 此处替换为你实际的x参数值,如果x是df的字段直接写列名即可
  lit(10) // 此处替换为你实际的y参数值,如果y是df的字段直接写列名即可
))

兼容优化方案(避免版本类型映射异常)

如果不同Spark版本的自动类型映射不稳定,可以修改UDF入参类型手动转换,兼容性更强:

val generateValues = spark.udf.register("my_udf",
       (laneAttrs: Seq[Row], x: Long, y: Int) => {
           // 手动将结构体序列转为目标List[(Long, Int)]
           val laneAttrList = laneAttrs.map(row => (row.getAs[Long](0), row.getAs[Int](1))).toList
           Converter.convert(laneAttrList, x, y)
          .map {
            case coord@Left(f) => throw new InterruptedException(s"$coord: $f")
            case Right(result) => Map("a" -> result.a,
              "b" -> result.b)
          }
  })

内容的提问来源于stack exchange,提问作者N. Surendran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 05:48:04