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

基于时间戳关联的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

SensorLocation_CityLocation_StateTimestampRain_valueSun_valueWind_value
seda_01Los AngelesCA15640735210.020.110.21
seda_01Los AngelesCA15640735220.010.100.21
seda_01Los AngelesCA15640735230.030.130.23
seda_02San FranciscoCA15640735210.050.15NULL
seda_02San FranciscoCA1564073522NULL0.120.25
seda_02San FranciscoCA15640735230.04NULL0.24

1. PySpark Implementation

The core approach here is:

  1. Flatten the nested Location struct.
  2. Explode each timestamp-value array into separate DataFrames.
  3. Collect all unique timestamps per sensor to ensure no entries are missed.
  4. 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 explode to break down nested arrays into rows, and explode_outer to 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 (or double if needed) makes them usable for numerical operations.
  • array_distinct removes duplicate timestamps from the combined list to avoid redundant rows.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:36:09