如何从结构体数组DataFrame中按最新时间戳提取字段值并新增列
提取数组中最新时间戳对应的colA值并添加为新列
DataFrame结构(Schema)
root |-- parentColumn: array | |-- element: struct | | |-- colA: string | | |-- colB: string | | |-- colTimestamp: string
数据示例
"parentColumn": [ { "colA": "LatestValueA", "colB": "LatestValueB", "colTimestamp": "2020-08-18T04:00:44.986000" }, { "colA": "OldValueA", "colB": "OldValueB", "colTimestamp": "2020-08-17T03:28:44.986000" } ]
需求说明
需要从parentColumn数组中,找到colTimestamp最新的元素,提取其colA值(示例中应返回LatestValueA),并将该值作为新列newColumn添加到DataFrame中,补充df.withColumn("newColumn", ?)的实现逻辑。
解决方案
Scala版本实现
import org.apache.spark.sql.functions._ val dfWithNewColumn = df.withColumn( "newColumn", element_at( array_sort( $"parentColumn", (x, y) => to_timestamp(x.getField("colTimestamp")).gt(to_timestamp(y.getField("colTimestamp"))) ), 1 ).getField("colA") )
Python版本实现
from pyspark.sql import functions as F df_with_new_column = df.withColumn( "newColumn", F.element_at( F.array_sort( F.col("parentColumn"), lambda x, y: F.to_timestamp(x.colTimestamp) > F.to_timestamp(y.colTimestamp) ), 1 ).colA )
逻辑说明
- 时间戳转换:用
to_timestamp把字符串类型的colTimestamp转为Spark可比较的Timestamp类型 - 数组排序:通过
array_sort对parentColumn数组按时间戳降序排列,确保最新元素排在首位 - 提取目标值:用
element_at取出排序后数组的第一个元素,再提取其中的colA字段作为新列值
内容的提问来源于stack exchange,提问作者Yadav
相关产品推荐
相关产品推荐

