如何用PySpark高效展平多层嵌套的API响应JSON结构?
多层嵌套JSON展平为PySpark DataFrame的优化方案
问题背景
我有一个来自API的多层嵌套JSON数组,需要将其展平为包含所有字段的单一PySpark DataFrame,目前已通过多次explode操作实现需求,但作为Spark新手,想了解是否有更优的实现方式。
API响应JSON示例
[ { "id": 1, "collection_name": "gym_equipment", "total_price": 5400, "lineitems": [ { "item_no": 1, "item_name": "dumbell", "quantity": 5, "price": 200, "splitaccountings": [ { "internalorder": "yes", "percentage": 50, "test": "no" }, { "internalorder": "no", "percentage": 50, "test": "yes" } ] }, { "item_no": 2, "item_name": "kettlebell", "quantity": 5, "price": 300, "splitaccountings": [ { "internalorder": "yes", "percentage": 50, "test": "no" }, { "internalorder": "no", "percentage": 50, "test": "yes" } ] }, { "item_no": 2, "item_name": "weight-set", "quantity": 15, "price": 420, "splitaccountings": [ { "internalorder": "yes", "percentage": 50, "test": "no" }, { "internalorder": "no", "percentage": 50, "test": "yes" } ] } ] }, { "id": 2, "collection_name": "holiday_equipment", "total_price": 5400, "lineitems": [ { "item_no": 1, "item_name": "suncream", "quantity": 5, "price": 200, "splitaccountings": [ { "internalorder": "yes", "percentage": 50, "test": "no" }, { "internalorder": "no", "percentage": 50, "test": "yes" } ] }, { "item_no": 2, "item_name": "beer", "quantity": 15, "price": 420, "splitaccountings": [ { "internalorder": "yes", "percentage": 100, "test": "no" } ] } ] }, { "id": 3, "collection_name": "hiking_equipment", "total_price": 5400, "lineitems": [ { "item_no": 1, "item_name": "100", "quantity": 5, "price": 200, "splitaccountings": [ { "internalorder": "yes", "percentage": 50, "test": "no" }, { "internalorder": "no", "percentage": 50, "test": "yes" } ] } ] } ]
期望输出DataFrame
+--------+-----------------+-----------+------------------+--------------------+-------------------+----------------+-------------------------------+----------------------------+----------------------+ |order_id|collection_name |total_price|**lineitem_item_no|**lineitem_item_name|**lineitem_quantity|**lineitem_price|**splitaccounting_internalorder|**splitaccounting_percentage|**splitaccounting_test| +--------+-----------------+-----------+------------------+--------------------+-------------------+----------------+-------------------------------+----------------------------+----------------------+ |1 |gym_equipment |5400 |1 |dumbell |5 |200 |yes |50 |no | |1 |gym_equipment |5400 |1 |dumbell |5 |200 |no |50 |yes | |1 |gym_equipment |5400 |2 |kettlebell |5 |300 |yes |50 |no | |1 |gym_equipment |5400 |2 |kettlebell |5 |300 |no |50 |yes | |1 |gym_equipment |5400 |2 |weight-set |15 |420 |yes |50 |no | |1 |gym_equipment |5400 |2 |weight-set |15 |420 |no |50 |yes | |2 |holiday_equipment|5400 |1 |suncream |5 |200 |yes |50 |no | |2 |holiday_equipment|5400 |1 |suncream |5 |200 |no |50 |yes | |2 |holiday_equipment|5400 |2 |beer |15 |420 |yes |100 |no | |3 |hiking_equipment |5400 |1 |100 |5 |200 |yes |50 |no | |3 |hiking_equipment |5400 |1 |100 |5 |200 |no |50 |yes | +--------+-----------------+-----------+------------------+--------------------+-------------------+----------------+-------------------------------+----------------------------+----------------------+
当前实现代码
from pyspark.sql.functions import explode, col spark = SparkSession.builder.appName("nested_json_processing").getOrCreate() df = spark.read.format("json").load("/content/test.json") df.show(truncate=False) # Explode 'lineitems' lineitems_exploded = df.select( col("id").alias("order_id"), col("collection_name"), col("total_price"), explode(col("lineitems")).alias("lineitem") ) lineitems_exploded.show() # Further explode 'splitaccountings' from the lineitems splitaccountings_exploded = lineitems_exploded.select( col("order_id"), col("collection_name"), col("total_price"), col("lineitem.item_no").alias("**lineitem_item_no"), col("lineitem.item_name").alias("**lineitem_item_name"), col("lineitem.quantity").alias("**lineitem_quantity"), col("lineitem.price").alias("**lineitem_price"), explode(col("lineitem.splitaccountings")).alias("splitaccounting") ) splitaccountings_exploded.show() # Flatten all fields flattened_df = splitaccountings_exploded.select( col("order_id"), col("collection_name"), col("total_price"), col("**lineitem_item_no"), col("**lineitem_item_name"), col("**lineitem_quantity"), col("**lineitem_price"), col("splitaccounting.internalorder").alias("**splitaccounting_internalorder"), col("splitaccounting.percentage").alias("**splitaccounting_percentage"), col("splitaccounting.test").alias("**splitaccounting_test") ) # Show the flattened DataFrame flattened_df.show(truncate=False) # Stop the Spark session if no longer needed spark.stop()
优化方案
你的现有实现逻辑是正确的,多次explode是处理嵌套数组的标准方式,但可以通过不同写法简化代码结构,提升可读性,以下是几种更优的实现方式:
方案1:链式调用简化代码
把多次select和explode操作串联起来,减少中间变量的创建,让代码更紧凑:
from pyspark.sql.functions import explode, col from pyspark.sql import SparkSession spark = SparkSession.builder.appName("nested_json_processing").getOrCreate() df = spark.read.format("json").load("/content/test.json") flattened_df = df.select( col("id").alias("order_id"), col("collection_name"), col("total_price"), explode(col("lineitems")).alias("lineitem") ).select( col("order_id"), col("collection_name"), col("total_price"), col("lineitem.item_no").alias("**lineitem_item_no"), col("lineitem.item_name").alias("**lineitem_item_name"), col("lineitem.quantity").alias("**lineitem_quantity"), col("lineitem.price").alias("**lineitem_price"), explode(col("lineitem.splitaccountings")).alias("splitaccounting") ).select( col("order_id"), col("collection_name"), col("total_price"), col("**lineitem_item_no"), col("**lineitem_item_name"), col("**lineitem_quantity"), col("**lineitem_price"), col("splitaccounting.internalorder").alias("**splitaccounting_internalorder"), col("splitaccounting.percentage").alias("**splitaccounting_percentage"), col("splitaccounting.test").alias("**splitaccounting_test") ) flattened_df.show(truncate=False) spark.stop()
方案2:使用Spark SQL语法
如果更熟悉SQL,可以注册临时视图后用SQL语句实现展平,逻辑更直观:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("nested_json_processing").getOrCreate() df = spark.read.format("json").load("/content/test.json") df.createOrReplaceTempView("orders") flattened_df = spark.sql(""" SELECT o.id AS order_id, o.collection_name, o.total_price, li.item_no AS `**lineitem_item_no`, li.item_name AS `**lineitem_item_name`, li.quantity AS `**lineitem_quantity`, li.price AS `**lineitem_price`, sa.internalorder AS `**splitaccounting_internalorder`, sa.percentage AS `**splitaccounting_percentage`, sa.test AS `**splitaccounting_test` FROM orders o LATERAL VIEW EXPLODE(o.lineitems) AS li LATERAL VIEW EXPLODE(li.splitaccountings) AS sa """) flattened_df.show(truncate=False) spark.stop()
方案3:使用withColumn逐步处理
这种方式步骤清晰,每一步操作明确,便于调试和修改:
from pyspark.sql.functions import explode, col from pyspark.sql import SparkSession spark = SparkSession.builder.appName("nested_json_processing").getOrCreate() df = spark.read.format("json").load("/content/test.json") # 重命名id为order_id df = df.withColumnRenamed("id", "order_id") # 展开lineitems数组 df = df.withColumn("lineitem", explode(col("lineitems"))).drop("lineitems") # 提取lineitem中的字段 df = df.select( "order_id", "collection_name", "total_price", col("lineitem.item_no").alias("**lineitem_item_no"), col("lineitem.item_name").alias("**lineitem_item_name"), col("lineitem.quantity").alias("**lineitem_quantity"), col("lineitem.price").alias("**lineitem_price"), col("lineitem.splitaccountings").alias("splitaccountings") ) # 展开splitaccountings数组 df = df.withColumn("splitaccounting", explode(col("splitaccountings"))).drop("splitaccountings") # 提取splitaccounting中的字段 flattened_df = df.select( "order_id", "collection_name", "total_price", "**lineitem_item_no", "**lineitem_item_name", "**lineitem_quantity", "**lineitem_price", col("splitaccounting.internalorder").alias("**splitaccounting_internalorder"), col("splitaccounting.percentage").alias("**splitaccounting_percentage"), col("splitaccounting.test").alias("**splitaccounting_test") ) flattened_df.show(truncate=False) spark.stop()
方案对比
- 原有代码:逻辑清晰,但中间变量较多,代码稍显冗长。
- 链式调用:代码简洁,减少中间变量,适合逻辑简单的嵌套展平。
- Spark SQL:对于熟悉SQL的开发者更友好,逻辑直观,复杂嵌套场景下可读性更强。
withColumn方式:步骤清晰,每一步操作明确,便于调试和修改。
这几种方案在性能上没有本质差异,Spark优化器会对执行计划统一优化,选择哪种方式主要取决于你的代码习惯和场景复杂度。
内容的提问来源于stack exchange,提问作者Looz
相关产品推荐
相关产品推荐

