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

如何将DataFrame中JSON字符串数组列展开为多行多列并提取Schema

问题描述

我的表中有一列存储JSON字符串数组,每个JSON对象对应一个时间戳。数据以sheet为维度记录,param_value列是包含各时间戳参数值的JSON数组,需要将其转换为按sheet、equipment、point展示的扁平结构。

之前尝试过展开JSON,但无法用*选择全部Schema,且ETL任务无法提前确定Schema,需要动态构建StructType。

原始数据

sheetequipparam_value
a1E1[{'point':'1','status':'no','log':'no'},{'point':'2','status':'ok','log':'no'},{'point':'3', 'status':'ok','log':'ok'}]
a2E1[{'point':'1','status':'no','log':'no'},{'point':'2','status':'ok','log':'no'},{'point':'3', 'status':'ok','log':'ok'}]
a3E1[{'point':'1','status':'no','log':'no'},{'point':'2','status':'ok','log':'no'},{'point':'3', 'status':'ok','log':'ok'}]

预期结果

sheetequipmentpointstatuslog
a1E11nono
a1E12okno
a1E13okok
a2E11nono
a2E12okno
a2E13okok
a3E11nono
a3E12okno
a3E13okok

当前实现与疑问

已通过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并完成数据展开:

  1. 提取第一条非空的param_value数据,解析为JSON数组
  2. 从数组中取第一个JSON对象,自动推断其StructType
  3. 构建包含该StructType的ArrayType Schema
  4. 使用该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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 07:20:34