PySpark explode嵌套列后Schema与底层嵌套结构不匹配问题
在Azure Synapse环境中使用PySpark,将多个结构一致的嵌套JSON文件读取为DataFrame,使用的样例JSON如下:
{ "AmountOfOrders": 2, "TotalEarnings": 1800, "OrderDetails": [ { "OrderNumber": 1, "OrderDate": "2022-7-06", "OrderLine": [ { "LineNumber": 1, "Product": "Laptop", "Price": 1000 }, { "LineNumber": 2, "Product": "Tablet", "Price": 500 }, { "LineNumber": 3, "Product": "Mobilephone", "Price": 300 } ] }, { "OrderNumber": 2, "OrderDate": "2022-7-06", "OrderLine": [ { "LineNumber": 1, "Product": "Printer", "Price": 100, "Discount": 0 }, { "LineNumber": 2, "Product": "Paper", "Price": 50, "Discount": 0 }, { "LineNumber": 3, "Product": "Toner", "Price": 30, "Discount": 0 } ] } ] }
现有自定义函数用于提取DataFrame中的数组、Struct类型字段,目标是提取OrderNumber为1的订单对应的OrderLine生成独立DataFrame,使用的代码如下:
def read_nested_structure(df,excludeList,messageType,coll): display(df.limit(10)) print('read_nested_structure') cols =[] match = 0 match_field = "" print(df.schema[coll].dataType.fields) for field in df.schema[coll].dataType.fields: for c in excludeList: if c == field.name: print('Match = ' + field.name) match = 1 if match == 0: cols.append(col(coll + "." + field.name).alias(field.name)) match = 0 print(cols) df = df.select(cols) return df def read_nested_structure_2(df,excludeList,messageType): cols =[] match = 0 for coll in df.schema.names: if isinstance(df.schema[coll].dataType, ArrayType): print( coll + "-- : Array") df = df.withColumn(coll, explode(coll).alias(coll)) cols.append(coll) elif isinstance(df.schema[coll].dataType, StructType): if messageType == 'Header': for field in df.schema[coll].dataType.fields: cols.append(col(coll + "." + field.name).alias(coll + "_" + field.name)) elif messageType == 'Content': print('Struct - Content') for field in df.schema[coll].dataType.fields: cols.append(col(coll + "." + field.name).alias(field.name)) else: for c in excludeList: if c == coll: match = 1 if match == 0: cols.append(coll) df = df.select(cols) return df df = spark.read.load(datalakelocation + '/test.json', format='json') df = unpack_to_content_dataframe_simple_2(df,exclude) df = df.filter(df.OrderNumber == 1) df = unpack_to_content_dataframe_simple_2(df,exclude) display(df.limit(10))
代码执行后,结果DataFrame中出现了不属于OrderNumber=1订单的Discount属性(该字段仅存在于OrderNumber=2的OrderLine子结构中),需要实现:过滤DataFrame行数据后同步更新Schema,移除筛选结果中实际不存在的字段(本例中即为Discount属性)。
Spark读取JSON文件时,会扫描全量数据合并推导出全局统一Schema,行级过滤操作不会触发Schema的自动更新。只要某个字段存在于全局Schema中,哪怕过滤后所有行的该字段值都是null,字段也会被保留。你遇到的Discount字段就是这种情况:它仅在OrderNumber=2的OrderLine结构中存在,被Spark合并到了全局OrderLine的Struct结构定义里,过滤OrderNumber=1之后该字段全为null,但不会被自动删除。
方法1:过滤后自动移除全null列(通用型,适配现有代码逻辑)
该方法不需要修改你已有的嵌套结构展开函数,只需要新增一个工具函数,在所有展开、过滤操作完成后,自动扫描并删除所有值全为null的列即可,侵入性最低,完全适配通用处理场景。
from pyspark.sql.functions import col, count def drop_all_null_columns(df): # 统计每列的非null值数量 non_null_stat = df.select([count(column).alias(column) for column in df.columns]).collect()[0] # 仅保留非null值数量大于0的列 keep_cols = [column for column in df.columns if non_null_stat[column] > 0] return df.select(keep_cols)
将原有处理流程最后一步加入该函数调用即可,同时注意修正原代码的函数名笔误(你定义的展开函数名为read_nested_structure_2,原调用写的是不存在的unpack_to_content_dataframe_simple_2):
df = spark.read.load(datalakelocation + '/test.json', format='json') df = read_nested_structure_2(df, exclude, 'Content') df = df.filter(col("OrderNumber") == 1) df = read_nested_structure_2(df, exclude, 'Content') # 移除全null的Discount字段 df = drop_all_null_columns(df) display(df.limit(10))
方法2:过滤后再解析JSON(性能更优,适合大数据量场景)
如果明确只需要处理OrderNumber=1的订单数据,可以先将JSON文件读为纯文本,过滤出目标数据片段后再执行Schema推导,这样生成的Schema天然不会包含仅存在于其他订单中的Discount字段,同时避免了解析无用字段的性能开销。
使用该方法需要提前明确目标数据的结构,灵活性不如方法1。
内容的提问来源于stack exchange,提问作者Erik hoeven

