Spark Structured Streaming 2.2.1对接Kafka:嵌套JSON数据解析处理问题
Hey Martin, let's walk through exactly how to parse that nested JSON from Kafka into a usable Spark Structured Streaming dataset with version 2.2.1—this is a super common use case, so we'll break it down step by step.
The core workflow here is: Pull raw Kafka messages → Extract JSON string → Define nested schema → Parse JSON to structured data → Unpack nested fields for calculations. Let's dive in.
1. Receive Raw Kafka Data
First, set up your Spark session and read from Kafka. Remember, Kafka messages are stored in the value field as binary data, so we'll cast that to a string first.
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // Initialize Spark Session val spark = SparkSession.builder() .appName("SensorStreamProcessing") .master("local[*]") // Remove this for production clusters .getOrCreate() // Read stream from Kafka val kafkaRawDF = spark.readStream .format("kafka") .option("kafka.bootstrap.servers", "your-broker-1:9092,your-broker-2:9092") .option("subscribe", "your-sensor-topic") .option("startingOffsets", "latest") // Use "earliest" to backfill old data .load() // Extract JSON string from Kafka's binary value field val jsonStringDF = kafkaRawDF.selectExpr("CAST(value AS STRING) AS json_payload")
2. Define Your Nested JSON Schema
Since your data has nested structures (like the pump object), we need to explicitly define a schema using StructType—Spark 2.2.1 doesn't support automatic schema inference for streaming JSON, so this step is mandatory.
Based on your sample data, here's how to define the schema (adjust fields to match your actual JSON):
// Define the nested schema for the `pump` object val pumpSubSchema = StructType(Seq( StructField("current", DoubleType, nullable = true), StructField("pressure", DoubleType, nullable = true) // Add any other fields inside `pump` here )) // Define the top-level schema for your sensor data val sensorSchema = StructType(Seq( StructField("id", IntegerType, nullable = false), StructField("timestamp", StringType, nullable = false), StructField("pump", pumpSubSchema, nullable = true) ))
Pro tip: If your timestamp is in a standard format (e.g., 2024-05-20T14:30:00), define it as TimestampType directly—this makes window calculations way easier later.
3. Parse JSON to a Structured Dataset
Use Spark's from_json function to convert the JSON string into a structured row. Then we'll unpack the top-level fields to make them easy to work with:
// Parse the JSON string using our defined schema val parsedSensorDF = jsonStringDF.select( from_json(col("json_payload"), sensorSchema).alias("sensor_data") ).select("sensor_data.*") // Unpack top-level fields (id, timestamp, pump) // Flatten nested `pump` fields to top-level columns (optional but recommended for calculations) val flattenedDF = parsedSensorDF.select( col("id"), col("timestamp").cast(TimestampType).alias("event_time"), // Convert to timestamp if needed col("pump.current").alias("pump_current"), col("pump.pressure").alias("pump_pressure") )
Now flattenedDF has clean, top-level columns that you can use for any calculations you need.
4. Add 60-Second Window Calculations (Per Your Requirement)
Since you're processing data every 60 seconds, we can set up a tumbling window (non-overlapping) with Spark's window functions. We'll also add a watermark to clean up old state and avoid memory bloat:
import org.apache.spark.sql.streaming.Trigger // Calculate 60-second aggregates per sensor val windowAggDF = flattenedDF .withWatermark("event_time", "1 minute") // Drop data older than 1 minute to manage state .groupBy( col("id"), window(col("event_time"), "60 seconds") // Tumbling window every 60 seconds ) .agg( avg("pump_current").alias("avg_pump_current"), max("pump_pressure").alias("max_pump_pressure") // Add any other calculations here ) // Write results to console (replace with your sink of choice: Kafka, Parquet, etc.) val streamingQuery = windowAggDF.writeStream .outputMode("update") // Use "append" for immutable results, "complete" for full aggregates .format("console") .trigger(Trigger.ProcessingTime("60 seconds")) // Trigger calculation every 60 seconds .start() streamingQuery.awaitTermination()
Key Notes for Spark 2.2.1
- No automatic schema inference: You must define the full nested schema manually—skip this and you'll get null values or parsing errors.
- Watermark limitations: Watermarks in 2.2.1 require the time column to be a
TimestampTypeand must be used with a window or groupBy. - Error handling: If some messages have malformed JSON, add a filter to drop bad records using
isnullon the parsed fields, or usetry_cast(though it's limited in 2.2.1).
If your JSON has deeper nesting, just extend the schema with additional StructType layers—this pattern scales to any nested structure.
内容的提问来源于stack exchange,提问作者Martin

