如何将DataFrame中JSON字符串数组列展开为多行多列并提取Schema
问题描述
我的表中有一列存储JSON字符串数组,每个JSON对象对应一个时间戳。数据以sheet为维度记录,param_value列是包含各时间戳参数值的JSON数组,需要将其转换为按sheet、equipment、point展示的扁平结构。
之前尝试过展开JSON,但无法用*选择全部Schema,且ETL任务无法提前确定Schema,需要动态构建StructType。
原始数据
| sheet | equip | param_value |
|---|---|---|
| a1 | E1 | [{'point':'1','status':'no','log':'no'},{'point':'2','status':'ok','log':'no'},{'point':'3', 'status':'ok','log':'ok'}] |
| a2 | E1 | [{'point':'1','status':'no','log':'no'},{'point':'2','status':'ok','log':'no'},{'point':'3', 'status':'ok','log':'ok'}] |
| a3 | E1 | [{'point':'1','status':'no','log':'no'},{'point':'2','status':'ok','log':'no'},{'point':'3', 'status':'ok','log':'ok'}] |
预期结果
| sheet | equipment | point | status | log |
|---|---|---|---|---|
| a1 | E1 | 1 | no | no |
| a1 | E1 | 2 | ok | no |
| a1 | E1 | 3 | ok | ok |
| a2 | E1 | 1 | no | no |
| a2 | E1 | 2 | ok | no |
| a2 | E1 | 3 | ok | ok |
| a3 | E1 | 1 | no | no |
| a3 | E1 | 2 | ok | no |
| a3 | E1 | 3 | ok | ok |
当前实现与疑问
已通过DDL格式字符串定义Schema实现展开,但需要手动指定字段列表schema_list。现在需要从JSON数组中自动提取Schema(所有JSON对象结构一致)。
当前实现代码:
schema_list = ['point', 'status','log'] # 需要从JSON数组中自动提取 schema = 'array<struct<' for c in schema_list: string_to_add = ',' + c +':string' schema = schema + string_to_add schema = schema.replace(",", "", 1)+'>>' s = "'" + schema + "'" print(s) # 'array<struct<point:string,status:string,log:string>>' f = df.selectExpr("sheet", "equip", f"inline(from_json(param_value, {s}))") f.show() # 输出结果: # +-----+-----+-----+------+---+ # |sheet|equip|point|status|log| # +-----+-----+-----+------+---+ # | a1| E1| 1| no| no| # | a1| E1| 2| ok| no| # | a1| E1| 3| ok| ok| # | a2| E1| 1| no| no| # | a2| E1| 2| ok| no| # | a2| E1| 3| ok| ok| # | a3| E1| 1| no| no| # | a3| E1| 2| ok| no| # | a3| E1| 3| ok| ok| # +-----+-----+-----+------+---+
解决方案:自动提取JSON数组的Schema
可以通过以下步骤动态提取Schema并完成数据展开:
- 提取第一条非空的
param_value数据,解析为JSON数组 - 从数组中取第一个JSON对象,自动推断其StructType
- 构建包含该StructType的ArrayType Schema
- 使用该Schema解析整个
param_value列并展开
代码实现(标准双引号JSON场景,生产环境推荐)
import json from pyspark.sql import functions as F from pyspark.sql.types import ArrayType # 提取第一条非空的param_value样本 sample_json = df.filter(F.col("param_value").isNotNull() & (F.col("param_value") != "")).select("param_value").first()[0] # 解析为JSON数组,取第一个对象 sample_obj = json.loads(sample_json)[0] # 自动推断单个JSON对象的Schema temp_df = spark.createDataFrame([sample_obj]) struct_schema = temp_df.schema # 构建数组类型的Schema array_schema = ArrayType(struct_schema) # 解析JSON数组并展开数据 df_parsed = df.withColumn("parsed_param", F.from_json(F.col("param_value"), array_schema)) result = df_parsed.selectExpr("sheet", "equip", "inline(parsed_param)") result.show()
兼容单引号JSON的临时处理(不推荐生产环境使用)
如果原始数据是单引号格式的非标准JSON,可以用eval临时解析(注意安全风险):
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType # 提取样本并解析 sample_json = df.filter(F.col("param_value").isNotNull() & (F.col("param_value") != "")).select("param_value").first()[0] sample_obj = eval(sample_json)[0] # 推断Schema并处理 temp_df = spark.createDataFrame([sample_obj]) struct_schema = temp_df.schema array_schema = ArrayType(struct_schema) # 解析展开 df_parsed = df.withColumn("parsed_param", F.from_json(F.col("param_value"), array_schema)) result = df_parsed.selectExpr("sheet", "equip", "inline(parsed_param)") result.show()
内容的提问来源于stack exchange,提问作者TingWei Chen
相关产品推荐
相关产品推荐

