如何将Spark DataFrame导出为指定结构的嵌套JSON文档?
Here's how you can transform your DataFrame into the desired nested JSON structure using Spark (Scala and PySpark examples included):
Step-by-Step Explanation & Code
1. Sample DataFrame Setup
First, let's recreate your sample DataFrame to test the solution:
Scala
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ val spark = SparkSession.builder().appName("NestedJsonExport").master("local[*]").getOrCreate() import spark.implicits._ val df = Seq( ("c1", "MCC", "cd1", 1, 2, 1.1, 1.05), ("c1", "MCC", "cd2", 2, 3, 1.1, 1.05), ("c1", "MCC", "cd3", 3, 4, 1.1, 1.05) ).toDF("A", "B", "val_of_B", "val1", "val2", "val3", "val4")
PySpark
from pyspark.sql import SparkSession from pyspark.sql.functions import array, collect_list, first, lit, create_map, to_json spark = SparkSession.builder.appName("NestedJsonExport").master("local[*]").getOrCreate() df = spark.createDataFrame([ ("c1", "MCC", "cd1", 1, 2, 1.1, 1.05), ("c1", "MCC", "cd2", 2, 3, 1.1, 1.05), ("c1", "MCC", "cd3", 3, 4, 1.1, 1.05) ], ["A", "B", "val_of_B", "val1", "val2", "val3", "val4"])
2. Create Individual MCC Entries
First, we combine val_of_B, val1, and val2 into a single array that represents each entry in the "MCC" list:
Scala
val dfWithEntries = df.withColumn("mcc_entry", array($"val_of_B", $"val1", $"val2"))
PySpark
df_with_entries = df.withColumn("mcc_entry", array("val_of_B", "val1", "val2"))
3. Group by A and Aggregate Data
Next, we group the data by A (along with val3 and val4 since they're A-level attributes) and collect all the MCC entries into a list. We also capture the B value to use as the JSON key:
Scala
val groupedDF = dfWithEntries.groupBy("A", "val3", "val4") .agg( collect_list("mcc_entry").alias("mcc_list"), first("B").alias("b_key") )
PySpark
grouped_df = df_with_entries.groupBy("A", "val3", "val4")\ .agg( collect_list("mcc_entry").alias("mcc_list"), first("B").alias("b_key") )
4. Build the Nested JSON Structure
We create a map that combines the A-level attributes (val3, val4) with the MCC list, using the B value as the key for the list. Then we convert this map to a JSON string:
Scala
val finalDF = groupedDF.withColumn("nested_json", to_json( map( lit("val3"), $"val3", lit("val4"), $"val4", $"b_key", $"mcc_list" ) ) ).select("A", "nested_json")
PySpark
final_df = grouped_df.withColumn("nested_json", to_json( create_map( lit("val3"), "val3", lit("val4"), "val4", "b_key", "mcc_list" ) ) ).select("A", "nested_json")
5. Export the JSON
To write the JSON to a file (one line per A):
Scala
finalDF.write.mode("overwrite").text("/path/to/your/output/directory")
PySpark
final_df.write.mode("overwrite").text("/path/to/your/output/directory")
Output Preview
Running finalDF.show(false) will give you this result:
+---+-----------------------------------------------------------------------------+ |A |nested_json | +---+-----------------------------------------------------------------------------+ |c1 |{"val3":1.1,"val4":1.05,"MCC":[["cd1",1,2],["cd2",2,3],["cd3",3,4]]} | +---+-----------------------------------------------------------------------------+
This matches exactly the nested JSON structure you requested!
Key Notes
- If your data has multiple
Bvalues perA, this solution will automatically add eachBvalue as a separate key in the JSON object with its corresponding list of entries. - The
first("B")call works here because all rows for a givenAshare the sameBvalue. If that's not the case in your full dataset, you may need to adjust the grouping logic to handle multipleBvalues perA.
内容的提问来源于stack exchange,提问作者user1124702

