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

Spark Streaming解析Kafka JSON数据并提取值保存至HDFS问题

Fixing Spark Streaming Extraction & Saving of Kafka JSON Values to HDFS

Hey there! Let's work through this issue together. You mentioned you can parse the JSON when printing but run into problems saving the extracted values to HDFS—let's break down what's going on and fix it.

First, looking at your Kafka data format:

{"group_city":""Washington"","group_country":""us"","event_name":""Outdoor Afro Goes Ziplining""}

Notice how each field value has extra escaped quotes (like \"Washington\")? When you parse this JSON, the extracted values will still have those surrounding quotes, which might be causing issues when saving. Plus, we need to make sure we're using Spark's structured streaming correctly to extract and persist the data.

Here's a step-by-step solution with code:

1. Set Up Dependencies & Schema

First, make sure you're using Spark Structured Streaming (it's more robust than the old DStream API) and define a schema for your JSON data—this ensures consistent parsing.

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._

// Define the schema matching your Kafka JSON structure
val kafkaJsonSchema = StructType(Seq(
  StructField("group_city", StringType, nullable = true),
  StructField("group_country", StringType, nullable = true),
  StructField("event_name", StringType, nullable = true)
))

2. Read Kafka Stream & Extract JSON

Next, read the Kafka stream, extract the value field (since that's where your JSON lives), and parse it using the schema we defined. We'll also clean up those extra escaped quotes.

// Initialize Spark Session
val spark = SparkSession.builder()
  .appName("KafkaJsonToHDFS")
  .getOrCreate()

import spark.implicits._

// Read from Kafka
val kafkaStream = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "your-bootstrap-server:9092") // Replace with your brokers
  .option("subscribe", "your-topic-name") // Replace with your topic
  .load()

// Extract Kafka's value as a string (since it's JSON)
val jsonStringStream = kafkaStream.selectExpr("CAST(value AS STRING) AS json_str")

// Parse JSON and clean up escaped quotes
val cleanedStream = jsonStringStream
  // Parse the JSON string into a structured row
  .select(from_json(col("json_str"), kafkaJsonSchema).alias("parsed_data"))
  // Filter out any rows that failed to parse (optional but recommended)
  .filter(col("parsed_data").isNotNull)
  // Extract fields and remove surrounding quotes from values
  .select(
    regexp_replace(col("parsed_data.group_city"), "^\"|\"$", "").alias("group_city"),
    regexp_replace(col("parsed_data.group_country"), "^\"|\"$", "").alias("group_country"),
    regexp_replace(col("parsed_data.event_name"), "^\"|\"$", "").alias("event_name")
  )

3. Save to HDFS

Now, save the cleaned data to HDFS. We'll use structured streaming's writeStream API, which handles fault tolerance with a checkpoint location.

Option 1: Save as Text (one field per line)

If you want to save a specific field (e.g., event_name) as plain text:

val textQuery = cleanedStream.select("event_name")
  .writeStream
  .format("text")
  .option("path", "hdfs://your-hdfs-path/kafka-text-output") // Replace with your HDFS path
  .option("checkpointLocation", "hdfs://your-hdfs-path/kafka-checkpoint") // Required for fault tolerance
  .outputMode("append")
  .start()

textQuery.awaitTermination()

Option 2: Save as Structured Format (CSV/Parquet)

For better data management, use a structured format like Parquet (columnar, efficient) or CSV:

val parquetQuery = cleanedStream
  .writeStream
  .format("parquet")
  .option("path", "hdfs://your-hdfs-path/kafka-parquet-output")
  .option("checkpointLocation", "hdfs://your-hdfs-path/kafka-parquet-checkpoint")
  .outputMode("append")
  .start()

parquetQuery.awaitTermination()

Key Notes to Avoid Issues

  • Checkpoint Location: Never skip this—Spark uses it to recover from failures. Make sure the HDFS path exists and is writable.
  • Cleaning Escaped Quotes: The regexp_replace removes leading/trailing quotes from each value. If your data has other escaped characters, you might need additional cleaning.
  • Bad Data Handling: Add mode("PERMISSIVE") to from_json if you want to keep rows with parsing errors (they'll have null values) or mode("DROPMALFORMED") to discard them entirely:
    from_json(col("json_str"), kafkaJsonSchema, Map("mode" -> "DROPMALFORMED")).alias("parsed_data")
    
  • Production Environment: Remove master("local[*]") if you added it—let your cluster manager (YARN/Kubernetes) handle resource allocation.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:18:20