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

如何从结构体数组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
)

逻辑说明

  1. 时间戳转换:用to_timestamp把字符串类型的colTimestamp转为Spark可比较的Timestamp类型
  2. 数组排序:通过array_sort对parentColumn数组按时间戳降序排列,确保最新元素排在首位
  3. 提取目标值:用element_at取出排序后数组的第一个元素,再提取其中的colA字段作为新列值

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 03:00:51