如何在Spark DataFrame中通过GroupBy提取时序JSON数据
在Spark中展开JSON格式的时序数据
问题背景
现有如下结构的Spark DataFrame:
| id | series |
|---|---|
| 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:
| id | timestamp | value |
|---|---|---|
| 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内置函数,性能最优,步骤如下:
- 解析
series列的JSON字符串为Map类型 - 将Map转换为(key, value)的条目数组
- 展开数组为多行
- 提取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
相关产品推荐
相关产品推荐

