Spark Streaming解析Kafka JSON数据并提取值保存至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_replaceremoves leading/trailing quotes from each value. If your data has other escaped characters, you might need additional cleaning. - Bad Data Handling: Add
mode("PERMISSIVE")tofrom_jsonif you want to keep rows with parsing errors (they'll have null values) ormode("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

