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
相关产品推荐
相关产品推荐

