Pyspark如何将嵌套JSON结构的DataFrame标准化为指定结构表格
解决方案
完整可运行代码
import json from pyspark.sql.functions import explode, col, lit, array, struct # 初始输入生成代码 api_response = {'02/09/2021':{'ABC':{'emp':'A1','value':'12421'},'DEF':{'emp':'D1','value':'3345'},'GHI':{'emp':'G2','value':'260048836600'},'JKL':{'emp':'J1','value':'66654654'}}} rdd = spark.sparkContext.parallelize([json.dumps(api_response)]) input_df = spark.read.json(rdd) # 处理逻辑 # 1. 将列名中存储的日期转为行字段 date_cols = input_df.columns stack_expr = f"stack({len(date_cols)}, {', '.join([f'\'{col_name}\', `{col_name}`' for col_name in date_cols])}) as (DATE, code_data)" date_unpivot_df = input_df.selectExpr(stack_expr) # 2. 把存储编码信息的struct转为可炸开的数组结构 code_list = date_unpivot_df.select("code_data.*").columns code_entries = array(*[struct(lit(c).alias("CODE"), col(f"code_data.{c}").alias("data")) for c in code_list]) exploded_df = date_unpivot_df.select("DATE", explode(code_entries).alias("code_entry")) # 3. 提取嵌套字段,调整列顺序 result_df = exploded_df.select( "code_entry.CODE", col("code_entry.data.emp").alias("EMP"), col("code_entry.data.value").alias("VALUE"), "DATE" ) # 输出结果 result_df.show(truncate=False)
结果输出
+----+---+------------+----------+ |CODE|EMP| VALUE| DATE| +----+---+------------+----------+ | ABC| A1| 12421|02/09/2021| | DEF| D1| 3345|02/09/2021| | GHI| G2|260048836600|02/09/2021| | JKL| J1| 66654654|02/09/2021| +----+---+------------+----------+
逻辑说明
你之前直接用explode失效的原因是:explode仅支持对数组、map类型字段执行炸开操作,而你读取得到的初始DataFrame中,日期列的值是struct类型,无法直接炸开。
上述处理的核心逻辑为:
- 先用stack函数将作为列名存储的日期转为行维度的DATE字段
- 把struct类型的编码数据手动转为包含编码名称和对应值的结构体数组
- 炸开数组后提取嵌套的emp、value字段,调整列顺序即可得到目标格式
内容的提问来源于stack exchange,提问作者Felipe FB
相关产品推荐
相关产品推荐

