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

如何从含数组结构体的Spark Scala DataFrame中获取指定Key的值

提取Spark DataFrame数组结构体中的指定字段值

问题背景

现有Spark DataFrame结构如下:

root
|-- recs: array (nullable = true)
| |-- element: struct (containsNull = true)
| | |-- key: string (nullable = true)
| | |-- value: binary (nullable = true)

数据示例(每行对应一个recs数组):

[{id, 123}, {source, test1}, {type, aa}, {subType, xy}, {createTs, 2022-10-31T10:36:34.951Z}, {requestTrackingId, trackId}, {locationId, 5560}]
[{id, 345}, {source, test2}, {type, ab}, {subType, xyz}, {createTs, 2022-10-29T10:36:34.951Z}, {requestTrackingId, trackId2}, {locationId, 55603}]

需要提取每个recs数组中key为createTs对应的value值,输出结果如下:

2022-10-31T10:36:34.951Z 
2022-10-29T10:36:34.951Z

Scala实现代码

方法1:高阶函数直接处理数组(性能更优)

无需展开数组,直接在数组层面过滤提取:

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

// 假设原DataFrame名为df
val resultDF = df.select(
  // 过滤出key为createTs的结构体,取第一个匹配项的value并转字符串
  element_at(filter(col("recs"), x => x.getField("key") === lit("createTs")), 1)
    .getField("value")
    .cast("string")
    .alias("createTs")
)

// 打印结果
resultDF.show(false)

方法2:展开数组后过滤提取

逻辑更直观,适合需要处理多匹配项的场景:

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

val resultDF = df
  // 展开recs数组为单行结构体
  .select(explode(col("recs")).alias("rec"))
  // 过滤出key为createTs的记录
  .filter(col("rec.key") === lit("createTs"))
  // 提取value并转换为字符串类型
  .select(col("rec.value").cast("string").alias("createTs"))

// 打印结果
resultDF.show(false)

说明

  • 原value字段类型为binary,必须通过cast("string")转换为字符串才能得到可读的时间格式输出。
  • 方法1中element_at(...,1)用于取数组中第一个匹配createTs的元素,适用于每个数组仅含一个createTs的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 09:55:15