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)
方案说明
原方案的问题在于:
- 仅处理了
records数组,未解析并处理newRecord字段(原字段为string类型的JSON数组,需先转换为结构化数组) - 使用
transform得到的是timestamp数组,直接agg(max('date'))只能获取每行数组内的最大值,而非全局所有元素的最大值
调整后的两种方案均解决了上述问题,可根据数据规模选择:
- 方案1适合数据量较小的场景,无需展开数据,效率较高
- 方案2适合数据量较大或需要对时间做进一步处理的场景,逻辑更直观
内容的提问来源于stack exchange,提问作者saraherceg
相关产品推荐
相关产品推荐

