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

Spark/Scala中如何将字符串类型的嵌套数组转换为结构体数组?

解决Spark 2.3.2中转换数组为结构体数组的问题

针对你的需求,Spark 2.3.2因为不支持transform这类高阶数组函数,我们可以通过拆分-转换-聚合的方式来实现,或者用自定义UDF,这里推荐前者(内置函数更高效且避免序列化问题):

方法1:使用内置函数(推荐)

这个方法通过explode把数组拆分成单行,处理每个字符串字段后再聚合回结构体数组:

import org.apache.spark.sql.functions._

// 第一步:拆分ar数组为单个字符串行
val explodedDF = fam_array.select($"id", explode($"ar").alias("family_str"))

// 第二步:把每个字符串拆分成字段数组,提取对应位置的元素生成结构体字段
val structDF = explodedDF
  // 拆分字符串,处理逗号前后的空格,避免空元素干扰
  .withColumn("parts", split(regexp_replace($"family_str", ",\\s*", ","), ","))
  // 按位置提取字段,不足的位置自动填充null
  .withColumn("f_name", $"parts".getItem(0))
  .withColumn("l_name", $"parts".getItem(1))
  .withColumn("status", $"parts".getItem(2))
  .withColumn("ph_no", $"parts".getItem(3))
  // 生成结构体
  .withColumn("family_struct", struct($"f_name", $"l_name", $"status", $"ph_no"))

// 第三步:按id聚合,把结构体收集成数组
val finalDF = structDF
  .groupBy($"id")
  .agg(collect_list($"family_struct").alias("family_name"))
  // 可选:如果需要保留原始的family_name列,可以join回去
  .join(fam_array.select($"id", $"family_name".alias("original_family_name")), "id")

运行后,finalDF的family_name列就是你需要的array<struct<f_name:string,l_name:string,status:string,ph_no:string>>类型。

方法2:使用自定义UDF

如果你更倾向于用UDF直接处理数组,可以定义一个转换数组的UDF:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

// 定义结构体Schema
val familyStructSchema = StructType(Seq(
  StructField("f_name", StringType),
  StructField("l_name", StringType),
  StructField("status", StringType),
  StructField("ph_no", StringType)
))

// 定义UDF:把字符串数组转换成结构体数组
val arrayToStructArray = udf((arr: Seq[String]) => {
  arr.map { str =>
    // 拆分字符串并清理空格、空元素
    val parts = str.split(",").map(_.trim).filter(_.nonEmpty)
    // 提取对应字段,不足则设为null
    val fName = if (parts.length >= 1) parts(0) else null
    val lName = if (parts.length >= 2) parts(1) else null
    val status = if (parts.length >= 3) parts(2) else null
    val phNo = if (parts.length >= 4) parts(3) else null
    Row(fName, lName, status, phNo)
  }
}, ArrayType(familyStructSchema))

// 应用UDF生成目标列
val finalDF = fam_array.withColumn("family_name", arrayToStructArray($"ar"))

结果验证

两种方法都能得到符合要求的Schema,你可以通过finalDF.printSchema()查看:

root
 |-- id: string (nullable = true)
 |-- family_name: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- f_name: string (nullable = true)
 |    |    |-- l_name: string (nullable = true)
 |    |    |-- status: string (nullable = true)
 |    |    |-- ph_no: string (nullable = true)

内容的提问来源于stack exchange,提问作者underwood

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 22:52:32