如何从含数组结构体的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
相关产品推荐
相关产品推荐

