Spark 3.5 结构化流(Structured Streaming):开源实时数据处理(Kafka→Spark→HDFS)实战
·
以下是一个基于Spark 3.5的结构化流处理实战方案,实现Kafka→Spark→HDFS的实时数据处理流程,包含完整代码和关键配置:
系统架构
Kafka (数据源) → Spark Structured Streaming (实时处理) → HDFS (持久化存储)
环境配置要求
- 组件版本:
- Kafka 2.8.0+
- Spark 3.5.0
- Hadoop 3.x (HDFS)
- 依赖包:
<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()
}
}
关键配置解析
-
Kafka连接参数:
.option("kafka.bootstrap.servers", "server1:9092,server2:9092") .option("failOnDataLoss", "false") // 容忍数据丢失 -
输出优化:
.option("maxRecordsPerFile", 1000000) // 控制文件大小 .partitionBy("year", "month", "day") // 三级分区目录 -
容错机制:
.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))
性能优化策略
-
小文件合并:
.option("parquet.block.size", 128 * 1024 * 1024) // 128MB块大小 -
触发间隔:
.trigger(Trigger.ProcessingTime("5 minutes")) // 根据数据量调整 -
资源分配:
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}")
// 添加告警逻辑
}
注意事项:
- HDFS路径需预先创建并有写入权限
- Checkpoint目录不可复用(重启时需清理)
- Kafka偏移量管理通过
checkpointLocation自动维护- 使用$$ \text{Watermark} = \text{event_time} - 10\text{min} $$处理延迟数据
此方案支持每秒处理万级事件,端到端延迟控制在分钟级,适用于日志收集、实时ETL等场景。
更多推荐

所有评论(0)