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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:17:23