SQL Server如何展平嵌套数组并合并值?万级JSON数据的处理方案
Alright, let's work through this problem since you've got a ton of dynamic name_* array columns and can't hardcode each one. The goal is to flatten all those nested arrays, align values by id and date, and fill missing values with 0. Below are two solid solutions using PySpark (perfect for your Spark connector setup with SQL Server) and Spark SQL.
Solution 1: PySpark Programmatic Approach
This method dynamically handles all name_* columns without hardcoding, making it scalable even for 10k+ columns.
Step 1: Load Your JSON Data
First, load all your JSON files into a Spark DataFrame:
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col, when, max as spark_max, lit # Initialize Spark session spark = SparkSession.builder.appName("FlattenNestedArrays").getOrCreate() # Load all JSON files (adjust the path to your file location) raw_df = spark.read.json("path/to/your/json/directory/*.json")
Step 2: Identify Dynamic Columns
Filter out the fixed id column and collect all name_* array columns:
# Get all columns starting with "name_" name_columns = [col for col in raw_df.columns if col.startswith("name_")] id_column = "id"
Step 3: Explode & Unify All Arrays
We'll explode each name_* array, reshape the data into a long format, and combine all results:
exploded_dfs = [] for col_name in name_columns: # Explode the array and extract date/val, plus track the original column name exploded_df = raw_df.select( col(id_column), explode(col(col_name)).alias("array_element") ).select( id_column, col("array_element.date").alias("date"), lit(col_name).alias("metric"), col("array_element.val").alias("val") ) exploded_dfs.append(exploded_df) # Combine all exploded DataFrames into one unified_long_df = exploded_dfs[0] for df_segment in exploded_dfs[1:]: unified_long_df = unified_long_df.unionByName(df_segment)
Step 4: Pivot to Wide Format & Fill Missing Values
Convert the long format back to a wide table, replacing nulls (missing dates for a column) with 0:
# Pivot the data to get one column per name_* metric pivoted_df = unified_long_df.groupBy(id_column, "date").pivot("metric").agg(spark_max("val")) # Replace null values with 0 for all name columns for col_name in name_columns: pivoted_df = pivoted_df.withColumn( col_name, when(col(col_name).isNull(), lit(0)).otherwise(col(col_name)) ) # Sort the final result by id and date final_df = pivoted_df.orderBy(id_column, "date") # View the output final_df.show()
Solution 2: Spark SQL Approach
If you prefer writing SQL, we can generate dynamic SQL to handle all name_* columns automatically.
Step 1: Register a Temporary View
First, register your raw DataFrame as a SQL view:
raw_df.createOrReplaceTempView("raw_nested_data")
Step 2: Generate & Execute Dynamic SQL
We'll build SQL fragments to explode each array, union the results, then pivot and fill missing values:
# Generate SQL fragments for exploding each name_* column explode_fragments = [] for col_name in name_columns: fragment = f""" SELECT id, array_element.date AS date, '{col_name}' AS metric, array_element.val AS val FROM raw_nested_data LATERAL VIEW EXPLODE({col_name}) exploded_view AS array_element """ explode_fragments.append(fragment) # Combine fragments into a UNION ALL query union_query = " UNION ALL ".join(explode_fragments) # Build the final pivot query with COALESCE to fill nulls with 0 pivot_query = f""" SELECT id, date, {", ".join([f"COALESCE({col_name}, 0) AS {col_name}" for col_name in name_columns])} FROM ( {union_query} ) unified_data PIVOT ( MAX(val) FOR metric IN ({", ".join([f"'{col_name}' AS {col_name}" for col_name in name_columns])}) ) pivoted_data ORDER BY id, date """ # Run the SQL query final_df = spark.sql(pivot_query) final_df.show()
Key Notes
- Dynamic Column Handling: Both solutions automatically detect all
name_*columns, so you don't need to update code if new columns are added. - Missing Value Handling:
COALESCE(SQL) andwhen/otherwise(PySpark) ensure any date missing from aname_*column gets a 0 value. - Performance: For large datasets, consider repartitioning by
idbefore grouping/pivoting to optimize Spark's parallel processing:raw_df.repartition("id").
内容的提问来源于stack exchange,提问作者Dan

