以下是一个基于Spark 3.5的结构化流处理实战方案,实现Kafka→Spark→HDFS的实时数据处理流程,包含完整代码和关键配置:

系统架构

Kafka (数据源) → Spark Structured Streaming (实时处理) → HDFS (持久化存储)

环境配置要求

  1. 组件版本:
    • Kafka 2.8.0+
    • Spark 3.5.0
    • Hadoop 3.x (HDFS)
  2. 依赖包:
    <dependency>
      <groupId>org.apache.spark</groupId>
      <artifactId>spark-sql-kafka-0-10_2.12</artifactId>
      <version>3.5.0</version>
    </dependency>
    

核心代码实现

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.Trigger

object KafkaToHDFS {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("StructuredStreamingKafkaToHDFS")
      .config("spark.sql.streaming.checkpointLocation", "/checkpoint/path")
      .getOrCreate()

    // 1. 从Kafka读取数据流
    val kafkaDF = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "kafka-server:9092")
      .option("subscribe", "input-topic")
      .option("startingOffsets", "latest")
      .load()
      .selectExpr("CAST(value AS STRING)")  // 转换二进制数据为字符串

    // 2. 数据处理(示例:解析JSON)
    val parsedDF = kafkaDF
      .withColumn("json_data", from_json(col("value"), schema)) // 自定义schema
      .select("json_data.*")
      .withColumn("processing_time", current_timestamp())  // 添加处理时间戳

    // 3. 写入HDFS
    val query = parsedDF.writeStream
      .format("parquet")  // 使用列式存储
      .option("path", "hdfs://namenode:8020/output/path")
      .option("checkpointLocation", "/hdfs/checkpoint/path")
      .trigger(Trigger.ProcessingTime("1 minute"))  // 每分钟输出一次
      .partitionBy("date")  // 按日期分区
      .start()

    query.awaitTermination()
  }
}

关键配置解析

  1. Kafka连接参数:

    .option("kafka.bootstrap.servers", "server1:9092,server2:9092")
    .option("failOnDataLoss", "false")  // 容忍数据丢失
    

  2. 输出优化:

    .option("maxRecordsPerFile", 1000000)  // 控制文件大小
    .partitionBy("year", "month", "day")   // 三级分区目录
    

  3. 容错机制:

    .config("spark.sql.streaming.forceDeleteTempCheckpointLocation", "true")
    

数据处理示例

假设Kafka消息为JSON格式:

{"user_id": "U1001", "event": "click", "timestamp": 1698765432000}

解析逻辑:

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

val schema = new StructType()
  .add("user_id", StringType)
  .add("event", StringType)
  .add("timestamp", LongType)

val enrichedDF = parsedDF
  .withColumn("date", to_date(from_unixtime(col("timestamp")/1000))) 
  .withColumn("event_count", when(col("event") === "click", 1).otherwise(0))

性能优化策略

  1. 小文件合并:

    .option("parquet.block.size", 128 * 1024 * 1024)  // 128MB块大小
    

  2. 触发间隔:

    .trigger(Trigger.ProcessingTime("5 minutes"))  // 根据数据量调整
    

  3. 资源分配:

    spark-submit --executor-memory 8g --num-executors 4 ...
    

监控指标

通过Spark UI监控:

  • Input Rate:Kafka数据摄入速度
  • Processing Time:批处理耗时
  • Sink Output:HDFS写入吞吐量

异常处理

query.exception.foreach { ex =>
  println(s"Query failed: ${ex.getMessage}")
  // 添加告警逻辑
}

注意事项:

  1. HDFS路径需预先创建并有写入权限
  2. Checkpoint目录不可复用(重启时需清理)
  3. Kafka偏移量管理通过checkpointLocation自动维护
  4. 使用$$ \text{Watermark} = \text{event_time} - 10\text{min} $$处理延迟数据

此方案支持每秒处理万级事件,端到端延迟控制在分钟级,适用于日志收集、实时ETL等场景。

更多推荐