PySpark如何扁平化DLT表含嵌套JSON的数组列生成扁平表
DLT表嵌套JSON数组扁平化解决方案
问题根因
直接调用select(explode("lorem"))仅会将数组炸开为多行,每行返回的默认列col是完整的嵌套Struct结构,没有显式提取Struct内部的各级字段,因此无法直接拿到field1_1、field4_1这类嵌套属性。
实现方案
分为PySpark API和Spark SQL两种实现方式:
1. PySpark 实现
from pyspark.sql.functions import explode, col # 第一步:炸开lorem数组,每个JSON元素对应一行,别名为item exploded_df = 原始DLT表变量.select(explode("lorem").alias("item")) # 第二步:逐层提取嵌套字段,生成扁平表 flattened_df = exploded_df.select( col("item.field1.field1_1").alias("field1_1"), col("item.field1.field1_2").alias("field1_2"), col("item.field2").alias("field2"), col("item.field3").alias("field3"), col("item.field4.field4_1").alias("field4_1"), col("item.field4.field4_2").alias("field4_2"), col("item.field5").alias("field5") ) # 保存为新DLT表 flattened_df.write.saveAsTable("新表名")
2. Spark SQL 实现
CREATE OR REPLACE TABLE 新表名 AS SELECT item.field1.field1_1 AS field1_1, item.field1.field1_2 AS field1_2, item.field2 AS field2, item.field3 AS field3, item.field4.field4_1 AS field4_1, item.field4.field4_2 AS field4_2, item.field5 AS field5 FROM ( -- 子查询先完成数组炸开 SELECT EXPLODE(lorem) AS item FROM 原始DLT表名 ) t
特殊场景补充:lorem列为JSON字符串的处理
如果你的lorem列存储的是字符串格式的JSON,而非已经解析好的Array<Struct>类型,需要先完成JSON结构解析再执行炸开操作:
from pyspark.sql.functions import from_json from pyspark.sql.types import StructType, StructField, StringType, IntegerType, NullType # 定义数组内单个JSON元素的Schema结构 item_schema = StructType([ StructField("field1", StructType([ StructField("field1_1", NullType()), StructField("field1_2", NullType()) ])), StructField("field2", StringType()), StructField("field3", IntegerType()), StructField("field4", StructType([ StructField("field4_1", NullType()), StructField("field4_2", NullType()) ])), StructField("field5", IntegerType()) ]) array_schema = ArrayType(item_schema) # 先将JSON字符串解析为Struct数组 parsed_df = 原始DLT表变量.withColumn("lorem_parsed", from_json("lorem", array_schema)) # 后续使用lorem_parsed字段执行上述的炸开、提取字段逻辑即可
内容的提问来源于stack exchange,提问作者António Mendes
相关产品推荐
相关产品推荐

