You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将Spark DataFrame导出为指定结构的嵌套JSON文档?

Solution to Convert Spark DataFrame to Nested 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 B values per A, this solution will automatically add each B value 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 given A share the same B value. If that's not the case in your full dataset, you may need to adjust the grouping logic to handle multiple B values per A.

内容的提问来源于stack exchange,提问作者user1124702

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 03:40:55