Scala中如何使用函数获取DataFrame行中Struct类型的元素
Scala/Spark行级计算实操指南
嘿,我来帮你搞定这个嵌套结构DataFrame的行级计算场景,先从创建DataFrame开始,再给你一些实用的计算示例参考~
第一步:创建目标DataFrame
你可以直接用提供的JSON字符串在Scala/Spark中快速构建DataFrame,代码如下:
import org.apache.spark.sql.SparkSession // 初始化SparkSession val spark = SparkSession.builder() .appName("RowLevelCalculationDemo") .master("local[*]") // 本地调试用,生产环境请移除该配置 .getOrCreate() import spark.implicits._ // 你的原始JSON数据 val jsonData = """{"available":false,"createTime":"2016-01-08","dataValue":{"names_source":{"first_names":["abc", "def"],"last_names_id":[123,456]},"another_source_array":[{"first":"1.1","last":"ONE"}],"another_source":"TableSources","location":"GMP", "timestamp":"2018-02-11"},"deleteTime":"2016-01-08"}""" // 从JSON字符串创建DataFrame val df = spark.read.json(Seq(jsonData).toDS()) // 打印Schema确认结构 df.printSchema()
运行后输出的Schema和你提供的一致,完整展示如下:
root |-- available: boolean (nullable = true) |-- createTime: string (nullable = true) |-- dataValue: struct (nullable = true) | |-- another_source: string (nullable = true) | |-- another_source_array: array (nullable = true) | | |-- element: struct (containsNull = true) | | | |-- first: string (nullable = true) | | | |-- last: string (nullable = true) | |-- location: string (nullable = true) | |-- names_source: struct (nullable = true) | | |-- first_names: array (nullable = true) | | | |-- element: string (containsNull = true) | | |-- last_names_id: array (nullable = true) | | | |-- element: long (containsNull = true) | |-- timestamp: string (nullable = true) |-- deleteTime: string (nullable = true)
第二步:行级计算示例
针对这个嵌套+数组混合的结构,我整理了几个常见的行级计算场景,你可以按需参考:
1. 嵌套字段提取与简单逻辑计算
比如从嵌套的时间字段中提取年份,或者根据available状态生成描述字段:
import org.apache.spark.sql.functions._ // 提取dataValue.timestamp中的年份 val dfWithYear = df.withColumn("data_year", year(to_timestamp($"dataValue.timestamp"))) // 根据available字段生成可读状态 val dfWithStatus = dfWithYear.withColumn("status", when($"available", "可用").otherwise("不可用")) // 查看结果 dfWithStatus.show(false)
2. 数组类型字段的行级处理
针对first_names、last_names_id这类数组字段,我们可以做拼接、长度统计等操作:
// 将first_names数组拼接为逗号分隔的字符串 val dfWithJoinedNames = df.withColumn("full_first_names", concat_ws(", ", $"dataValue.names_source.first_names")) // 统计last_names_id数组的元素数量 val dfWithIdCount = dfWithJoinedNames.withColumn("last_id_count", size($"dataValue.names_source.last_names_id")) // 查看结果 dfWithIdCount.show(false)
3. 自定义UDF实现复杂行级计算
如果内置函数满足不了你的需求,可以自定义UDF处理复杂逻辑,比如处理another_source_array中的嵌套数组数据:
// 自定义UDF:提取数组中第一个元素的first字段,转为Double后加1 val addOneToFirst = udf((arr: Seq[Map[String, String]]) => { arr.headOption.map(_("first").toDouble + 1).getOrElse(0.0) }) // 应用UDF生成新字段 val dfWithCustomCalc = df.withColumn("first_plus_one", addOneToFirst($"dataValue.another_source_array")) // 查看结果 dfWithCustomCalc.show(false)
内容的提问来源于stack exchange,提问作者ChikuMiku
相关产品推荐
相关产品推荐

