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

如何在Spark DataFrame中通过GroupBy提取时序JSON数据

在Spark中展开JSON格式的时序数据

问题背景

现有如下结构的Spark DataFrame:

idseries
1{"2016-01-31T00:00:00.000Z": null, "2016-06-30T00:00:00.000Z": 6394317.0, "2016-07-31T00:00:00.000Z": 6550781.0, "2016-08-31T00:00:00.000Z": 7107308.0}
2{"2016-01-31T00:00:00.000Z": null, "2016-06-30T00:00:00.000Z": 6394317.0}

需要将其转换为如下格式的DataFrame:

idtimestampvalue
1"2016-01-31T00:00:00.000Z"null
1"2016-06-30T00:00:00.000Z"6394317.0
1"2016-07-31T00:00:00.000Z"6550781.0
1"2016-08-31T00:00:00.000Z"7107308.0
2"2016-01-31T00:00:00.000Z"null
2"2016-06-30T00:00:00.000Z"6394317.0

用户已通过Pandas遍历分组对象实现需求,希望用Spark分布式操作完成相同功能。

解决方案

Spark无需遍历分组,直接利用内置JSON处理和行展开函数即可实现,推荐以下两种方法:

方法一:使用from_json+map_entries+explode(Spark 2.3+)

这种方法完全依赖Spark内置函数,性能最优,步骤如下:

  1. 解析series列的JSON字符串为Map类型
  2. 将Map转换为(key, value)的条目数组
  3. 展开数组为多行
  4. 提取key和value并重命名列
from pyspark.sql import functions as F
from pyspark.sql.types import MapType, StringType, DoubleType

# 定义Map类型Schema:键为字符串格式时间戳,值为可空Double类型
map_schema = MapType(StringType(), DoubleType())

# 处理数据
result_df = (
    df
    # 解析JSON字符串为Map结构
    .withColumn("series_map", F.from_json(F.col("series"), map_schema))
    # 将Map转换为(key, value)条目数组
    .withColumn("series_entries", F.map_entries(F.col("series_map")))
    # 展开数组为多行
    .withColumn("entry", F.explode(F.col("series_entries")))
    # 提取时间戳和值,并重命名列
    .select(
        F.col("id"),
        F.col("entry.key").alias("timestamp"),
        F.col("entry.value").alias("value")
    )
)

result_df.show(truncate=False)

方法二:自定义UDF+explode(兼容旧版本Spark)

如果你的Spark版本低于2.3,没有map_entries函数,可通过自定义UDF解析JSON并生成条目数组:

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

# 自定义UDF:将JSON字符串转换为(timestamp, value)的数组
@F.udf(returnType=ArrayType(StructType([
    StructField("timestamp", StringType()),
    StructField("value", DoubleType())
])))
def parse_series_json(json_str):
    if not json_str:
        return []
    data = json.loads(json_str)
    return [(k, v) for k, v in data.items()]

# 处理数据
result_df = (
    df
    # 解析JSON为条目数组
    .withColumn("series_entries", parse_series_json(F.col("series")))
    # 展开数组为多行
    .withColumn("entry", F.explode(F.col("series_entries")))
    # 提取列
    .select(
        F.col("id"),
        F.col("entry.timestamp"),
        F.col("entry.value")
    )
)

result_df.show(truncate=False)

说明

  • 优先选择方法一,内置函数的分布式执行效率远高于遍历分组
  • 两种方法都会保留原始数据中的null值,完全匹配需求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 12:50:22