如何将Spark DataFrame两列转换为List[(Long,Int)]传入自定义UDF
解决方案
核心实现方法
你可以直接使用Spark内置的arrays_zip函数配对两个数组的同位置元素,Spark会自动完成类型映射转换为UDF需要的List[(Long, Int)]格式,具体操作步骤如下:
- 首先导入依赖的内置函数
import org.apache.spark.sql.functions.{arrays_zip, lit, size}
- 调用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
相关产品推荐
相关产品推荐

