基于时间戳关联的JSON扁平化:Java/Scala/Spark/PySpark实现
Solution for Timestamp-Based JSON Flattening & Joining
Got it, let's tackle this JSON flattening problem where we need to explode timestamp-based arrays and join them on matching timestamps—filling NULLs if a metric is missing for a given timestamp. Below are implementation solutions across PySpark, Scala (Spark), and Java (Spark).
First, let's recap the problem with sample input and expected output to make sure we're aligned:
Sample Input JSON
[ { "Sensor": "seda_01", "Location": { "City": "Los Angeles", "State": "CA" }, "rain_value": [ [ "1564073521", "0.02" ], [ "1564073522", "0.01" ], [ "1564073523", "0.03" ] ], "sun_value": [ [ "1564073521", "0.11" ], [ "1564073522", "0.10" ], [ "1564073523", "0.13" ] ], "wind_value": [ [ "1564073521", "0.21" ], [ "1564073522", "0.21" ], [ "1564073523", "0.23" ] ] }, { "Sensor": "seda_02", "Location": { "City": "San Francisco", "State": "CA" }, "rain_value": [ [ "1564073521", "0.05" ], [ "1564073523", "0.04" ] ], "sun_value": [ [ "1564073521", "0.15" ], [ "1564073522", "0.12" ] ], "wind_value": [ [ "1564073522", "0.25" ], [ "1564073523", "0.24" ] ] } ]
Expected Output DataFrame
| Sensor | Location_City | Location_State | Timestamp | Rain_value | Sun_value | Wind_value |
|---|---|---|---|---|---|---|
| seda_01 | Los Angeles | CA | 1564073521 | 0.02 | 0.11 | 0.21 |
| seda_01 | Los Angeles | CA | 1564073522 | 0.01 | 0.10 | 0.21 |
| seda_01 | Los Angeles | CA | 1564073523 | 0.03 | 0.13 | 0.23 |
| seda_02 | San Francisco | CA | 1564073521 | 0.05 | 0.15 | NULL |
| seda_02 | San Francisco | CA | 1564073522 | NULL | 0.12 | 0.25 |
| seda_02 | San Francisco | CA | 1564073523 | 0.04 | NULL | 0.24 |
1. PySpark Implementation
The core approach here is:
- Flatten the nested
Locationstruct. - Explode each timestamp-value array into separate DataFrames.
- Collect all unique timestamps per sensor to ensure no entries are missed.
- Left-join all metric DataFrames to the full timestamp list to fill NULLs for missing metrics.
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col, array_distinct, explode_outer spark = SparkSession.builder.appName("JsonTimestampFlatten").getOrCreate() # Read input JSON file df = spark.read.json("path/to/your/input.json") # Step 1: Flatten the Location struct df_flat = df.select( col("Sensor"), col("Location.City").alias("Location_City"), col("Location.State").alias("Location_State"), col("rain_value"), col("sun_value"), col("wind_value") ) # Step 2: Create individual DataFrames for each metric def create_metric_df(base_df, metric_name): return base_df.select( "Sensor", "Location_City", "Location_State", explode(col(f"{metric_name}_value")).alias(f"{metric_name}_entry") ).select( "Sensor", "Location_City", "Location_State", col(f"{metric_name}_entry")[0].alias("Timestamp"), col(f"{metric_name}_entry")[1].cast("float").alias(f"{metric_name.capitalize()}_value") ) rain_df = create_metric_df(df_flat, "rain") sun_df = create_metric_df(df_flat, "sun") wind_df = create_metric_df(df_flat, "wind") # Step 3: Collect all unique timestamps per sensor all_timestamps_df = df_flat.select( "Sensor", "Location_City", "Location_State", array_distinct( col("rain_value").getItem(0) + col("sun_value").getItem(0) + col("wind_value").getItem(0) ).alias("all_timestamps") ).select( "Sensor", "Location_City", "Location_State", explode_outer(col("all_timestamps")).alias("Timestamp") ) # Step 4: Join all metrics to the full timestamp list final_df = all_timestamps_df.join( rain_df, on=["Sensor", "Location_City", "Location_State", "Timestamp"], how="left" ).join( sun_df, on=["Sensor", "Location_City", "Location_State", "Timestamp"], how="left" ).join( wind_df, on=["Sensor", "Location_City", "Location_State", "Timestamp"], how="left" ) # Show the result final_df.show()
2. Scala (Spark) Implementation
Same logic translated to Scala, with helper functions to reduce code duplication:
import org.apache.spark.sql.{SparkSession, DataFrame} import org.apache.spark.sql.functions._ object SensorJsonFlattener { def main(args: Array[String]): Unit = { val spark = SparkSession.builder.appName("JsonTimestampFlatten").getOrCreate() import spark.implicits._ // Read input JSON val df = spark.read.json("path/to/your/input.json") // Flatten Location struct val dfFlat = df.select( $"Sensor", $"Location.City".alias("Location_City"), $"Location.State".alias("Location_State"), $"rain_value", $"sun_value", $"wind_value" ) // Helper to create metric-specific DataFrames def createMetricDF(baseDF: DataFrame, metricName: String): DataFrame = { baseDF.select( $"Sensor", $"Location_City", $"Location_State", explode(col(s"${metricName}_value")).alias(s"${metricName}_entry") ).select( $"Sensor", $"Location_City", $"Location_State", col(s"${metricName}_entry")(0).alias("Timestamp"), col(s"${metricName}_entry")(1).cast("float").alias(s"${metricName.capitalize}_value") ) } val rainDF = createMetricDF(dfFlat, "rain") val sunDF = createMetricDF(dfFlat, "sun") val windDF = createMetricDF(dfFlat, "wind") // Collect all unique timestamps per sensor val allTimestampsDF = dfFlat.select( $"Sensor", $"Location_City", $"Location_State", array_distinct( $"rain_value".getItem(0) ++ $"sun_value".getItem(0) ++ $"wind_value".getItem(0) ).alias("all_timestamps") ).select( $"Sensor", $"Location_City", $"Location_State", explode_outer($"all_timestamps").alias("Timestamp") ) // Join all metrics to the full timestamp list val finalDF = allTimestampsDF .join(rainDF, Seq("Sensor", "Location_City", "Location_State", "Timestamp"), "left") .join(sunDF, Seq("Sensor", "Location_City", "Location_State", "Timestamp"), "left") .join(windDF, Seq("Sensor", "Location_City", "Location_State", "Timestamp"), "left") finalDF.show() } }
3. Java (Spark) Implementation
Java version of the same solution, with explicit dataset transformations:
import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import static org.apache.spark.sql.functions.*; public class SensorJsonFlatteningJava { public static void main(String[] args) { SparkSession spark = SparkSession.builder() .appName("JsonTimestampFlattenJava") .getOrCreate(); // Read input JSON Dataset<Row> df = spark.read().json("path/to/your/input.json"); // Flatten Location struct Dataset<Row> dfFlat = df.select( col("Sensor"), col("Location.City").alias("Location_City"), col("Location.State").alias("Location_State"), col("rain_value"), col("sun_value"), col("wind_value") ); // Create rain DataFrame Dataset<Row> rainDF = dfFlat.select( col("Sensor"), col("Location_City"), col("Location_State"), explode(col("rain_value")).alias("rain_entry") ).select( col("Sensor"), col("Location_City"), col("Location_State"), col("rain_entry").getItem(0).alias("Timestamp"), col("rain_entry").getItem(1).cast("float").alias("Rain_value") ); // Create sun DataFrame Dataset<Row> sunDF = dfFlat.select( col("Sensor"), col("Location_City"), col("Location_State"), explode(col("sun_value")).alias("sun_entry") ).select( col("Sensor"), col("Location_City"), col("Location_State"), col("sun_entry").getItem(0).alias("Timestamp"), col("sun_entry").getItem(1).cast("float").alias("Sun_value") ); // Create wind DataFrame Dataset<Row> windDF = dfFlat.select( col("Sensor"), col("Location_City"), col("Location_State"), explode(col("wind_value")).alias("wind_entry") ).select( col("Sensor"), col("Location_City"), col("Location_State"), col("wind_entry").getItem(0).alias("Timestamp"), col("wind_entry").getItem(1).cast("float").alias("Wind_value") ); // Collect all unique timestamps per sensor Dataset<Row> allTimestampsDF = dfFlat.select( col("Sensor"), col("Location_City"), col("Location_State"), array_distinct( concat( col("rain_value").getItem(0), col("sun_value").getItem(0), col("wind_value").getItem(0) ) ).alias("all_timestamps") ).select( col("Sensor"), col("Location_City"), col("Location_State"), explode_outer(col("all_timestamps")).alias("Timestamp") ); // Join all metrics to the full timestamp list Dataset<Row> finalDF = allTimestampsDF .join(rainDF, new String[]{"Sensor", "Location_City", "Location_State", "Timestamp"}, "left") .join(sunDF, new String[]{"Sensor", "Location_City", "Location_State", "Timestamp"}, "left") .join(windDF, new String[]{"Sensor", "Location_City", "Location_State", "Timestamp"}, "left"); finalDF.show(); } }
Key Notes
- We use
explodeto break down nested arrays into rows, andexplode_outerto retain all timestamps even if some metrics are missing. - Left joins ensure that missing metric values for a timestamp are automatically filled with NULL.
- Casting value columns to
float(ordoubleif needed) makes them usable for numerical operations. array_distinctremoves duplicate timestamps from the combined list to avoid redundant rows.
内容的提问来源于stack exchange,提问作者user2458922
相关产品推荐
相关产品推荐

