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

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.

Parsing Nested JSON from Kafka in Spark Structured Streaming 2.2.1

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 TimestampType and must be used with a window or groupBy.
  • Error handling: If some messages have malformed JSON, add a filter to drop bad records using isnull on the parsed fields, or use try_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:01:22