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

PySpark提取JSON数组所有time值并求最大日期的问题

问题:提取所有数组中的time值并求最大日期

需要从JSON文件的records数组以及newRecord字段对应的数组中提取time值,转换为日期列后求该列的最大日期。但当前方案仅能处理首个records数组的time值,无法获取所有数组中的对应值,需调整方案。

JSON文件结构

root
 |-- newRecord: string (nullable = true)
 |-- records: array (nullable = true)
 |    |-- element: struct (containsNull = true)
 |    |    |-- val1: long (nullable = true)
 |    |    |-- val2: long (nullable = true)
 |    |    |-- time: string(nullable = true)

当前解决方案(存在局限)

import pyspark.sql.functions as F
df = spark.read.json("....json")
df.select (
    "newRecord",
    F.transform('records', lambda x: to_timestamp(x['time'])).alias('date')).agg(max('date')).first()

输入示例

{
    "records": [
        {
             "val1": 37.99711,
             "val2": 231.7571,
             "time":  "2023-01-17T14:39:09.207Z"
        },
        {
             "val1": 37.99711,
             "val2": 231.7571,
             "time":  "2023-01-16T14:37:09.207Z"
        }
    ],
    "newRecord": "[{\"val1\": 123, \"val2\": 456, \"time\": \"2023-01-18T10:00:00.000Z\"}]"
}

调整后的方案

方案1:合并数组后取全局最大值

该方案通过合并records和newRecord的时间数组,直接计算所有时间的最大值:

import pyspark.sql.functions as F
from pyspark.sql.types import StructType, StructField, LongType, StringType, ArrayType

# 定义数组元素的结构,与records一致
element_schema = StructType([
    StructField("val1", LongType(), nullable=True),
    StructField("val2", LongType(), nullable=True),
    StructField("time", StringType(), nullable=True)
])

# 读取JSON文件
df = spark.read.json("....json")

# 解析newRecord字符串为数组(原字段是string类型的JSON数组)
df = df.withColumn("newRecord_array", F.from_json(F.col("newRecord"), ArrayType(element_schema)))

# 提取并转换两个数组中的time为timestamp
df = df.withColumn("records_times", F.transform("records", lambda x: F.to_timestamp(x["time"])))
df = df.withColumn("newRecord_times", F.transform("newRecord_array", lambda x: F.to_timestamp(x["time"])))

# 合并两个时间数组
df = df.withColumn("all_times", F.concat("records_times", "newRecord_times"))

# 计算全局最大日期
max_time = df.select(F.array_max("all_times").alias("max_time")).first()["max_time"]
print(max_time)

方案2:展开所有时间元素后聚合

该方案将数组展开为单行记录,再通过聚合函数求最大值:

import pyspark.sql.functions as F
from pyspark.sql.types import StructType, StructField, LongType, StringType, ArrayType

element_schema = StructType([
    StructField("val1", LongType(), nullable=True),
    StructField("val2", LongType(), nullable=True),
    StructField("time", StringType(), nullable=True)
])

df = spark.read.json("....json")

# 解析newRecord为数组
df = df.withColumn("newRecord_array", F.from_json(F.col("newRecord"), ArrayType(element_schema)))

# 展开records中的所有time并转换为timestamp
records_time_df = df.select(F.explode(F.transform("records", lambda x: F.to_timestamp(x["time"]))).alias("time"))

# 展开newRecord中的所有time并转换为timestamp
newRecord_time_df = df.select(F.explode(F.transform("newRecord_array", lambda x: F.to_timestamp(x["time"]))).alias("time"))

# 合并数据集并求最大日期
max_time = records_time_df.union(newRecord_time_df).agg(F.max("time")).first()[0]
print(max_time)

方案说明

原方案的问题在于:

  1. 仅处理了records数组,未解析并处理newRecord字段(原字段为string类型的JSON数组,需先转换为结构化数组)
  2. 使用transform得到的是timestamp数组,直接agg(max('date'))只能获取每行数组内的最大值,而非全局所有元素的最大值

调整后的两种方案均解决了上述问题,可根据数据规模选择:

  • 方案1适合数据量较小的场景,无需展开数据,效率较高
  • 方案2适合数据量较大或需要对时间做进一步处理的场景,逻辑更直观

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 06:20:27