一、整合的基础背景与适用场景
很多做数据处理的朋友,不管是刚入行的新人还是有几年经验的老手,大概率都碰到过一个头疼的问题:一边要实时处理源源不断的用户行为、设备日志这类流式数据,一边还要把这些数据存成能随时查询、不会乱的“数据湖”,既要速度快又要稳定性高。Hudi和Spark Structured Streaming就是专门解决这类问题的组合,接下来先说说什么时候适合用这套组合。
1.1 核心应用场景
这套组合最适合的场景,就是需要“实时入湖+后续批量查询分析”的业务。比如电商平台的实时订单数据:每一笔用户下单、支付、退款的行为都要实时存到数据湖,之后运营人员既能实时看当前的订单成交情况,也能按天、按周批量分析订单的地域分布、品类偏好;再比如物联网场景下的设备状态上报:设备每几秒发一次温度、电压数据,要实时存到湖,之后既能实时监控异常,也能批量做设备寿命预测。
1.2 技术优缺点
先说说优点,第一是实时性够高:Spark Structured Streaming本身就是微批处理,延迟可以控制在几秒以内,搭配Hudi的事务机制,能保证数据不丢不重复;第二是数据湖的特性拉满:Hudi支持数据的更新、删除,还能做版本回溯,比如某天的订单数据录错了,能回滚到错误发生前的版本;第三是成本低:用的都是开源技术,不需要额外付费,而且能直接对接S3、OSS这类廉价存储,不用单独买昂贵的数仓存储。 再说说缺点,第一是配置复杂:两个技术的参数特别多,稍微改错一个就可能导致数据重复、延迟变高;第二是调优难度大:不同的数据量、业务场景需要不同的调优策略,比如小文件太多怎么解决,并发太高怎么控制;第三是学习成本高:新人要同时搞懂Spark Structured Streaming的微批机制、Hudi的表类型、事务原理,得花不少时间。
二、整合前的准备工作
要把Hudi和Spark Structured Streaming整合起来,得先把基础环境搭好,不然写代码的时候会碰到各种莫名其妙的错误。
2.1 技术栈说明
所有示例统一使用以下技术栈:
- Spark 3.3.2(稳定版,对Hudi的支持最完善)
- Hudi 0.12.0(和Spark 3.3.2兼容性最好)
- Scala 2.12.15(Spark和Hudi都适配的Scala版本)
- HDFS 3.3.4(作为数据湖的存储介质,也可以换成S3、OSS等)
2.2 基础环境配置
首先要保证Spark和Hudi的版本匹配,版本不匹配会出现类冲突的错误。然后要把Hudi的依赖包放到Spark的jars目录下,或者在提交任务的时候指定依赖。另外,HDFS要配置好权限,保证Spark任务能读写HDFS上的目录。
三、整合的核心代码示例
接下来给大家一个完整的示例,从读取实时数据到写入Hudi,每一步都有注释,方便大家理解。
3.1 示例代码(Scala)
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.streaming.Trigger
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.config.HoodieWriteConfig._
object HudiStructuredStreamingDemo {
def main(args: Array[String]): Unit = {
// 1. 创建SparkSession,配置Hudi相关参数
val spark = SparkSession.builder()
.appName("HudiStructuredStreamingDemo")
// 配置Hudi的依赖,提交任务时如果已经放了jars可以不用加
.config("spark.jars", "/path/to/hudi-spark3-bundle_2.12-0.12.0.jar")
.master("local[*]") // 本地测试用,生产环境去掉
.getOrCreate()
// 导入Spark的隐式转换,方便后续操作
import spark.implicits._
// 2. 读取实时数据源,这里以读取Kafka为例,模拟实时数据流
// 实际生产中可以换成其他数据源,比如Socket、Flume等
val kafkaDF = spark.readStream
.format("kafka")
.option("kafka.bootstrap.servers", "localhost:9092") // Kafka地址
.option("subscribe", "test_topic") // 订阅的Kafka主题
.option("startingOffsets", "earliest") // 从最早的偏移量开始读取
.option("maxOffsetsPerTrigger", "1000") // 每批次最多读取1000条数据,控制批次大小
.load()
// 3. 解析Kafka的value字段,这里假设value是JSON格式,字段有id、name、age、event_time
// 定义JSON的结构
val schema = "id string, name string, age int, event_time long"
val parsedDF = kafkaDF.selectExpr("CAST(value AS STRING) as json")
.selectExpr(s"from_json(json, '$schema') as data")
.select("data.*")
// 4. 把解析后的数据写入Hudi
val query = parsedDF.writeStream
.format("hudi")
// Hudi表的类型,COW(写时复制)适合读多写少的场景,MOR(读时合并)适合写多读少的场景
.option(TABLE_TYPE_OPT_KEY, MOR_TABLE_TYPE_OPT_VAL)
// Hudi表的主键,保证数据的唯一性
.option(RECORDKEY_FIELD_OPT_KEY, "id")
// Hudi表的分区字段,按event_time的日期分区,方便后续查询
.option(PARTITIONPATH_FIELD_OPT_KEY, "event_time")
// Hudi表的预合并字段,这里用event_time,保证最新的数据优先
.option(PRECOMBINE_FIELD_OPT_KEY, "event_time")
// Hudi表的索引类型,BLOOM适合主键查询,SIMPLE适合全表扫描
.option(INDEX_TYPE_OPT_KEY, "BLOOM")
// Hudi表的路径,存储在HDFS上
.option("path", "hdfs://localhost:9000/hudi/test_table")
// 每批次的触发时间,这里设为10秒,控制微批的间隔
.trigger(Trigger.ProcessingTime("10 seconds"))
// 检查点路径,保证任务失败后能从上次的位置继续执行,防止数据重复
.option("checkpointLocation", "/tmp/hudi/checkpoint/test_table")
.start()
// 等待任务结束
query.awaitTermination()
}
}
3.2 代码的关键说明
上面的代码里有几个点需要特别注意,第一个是Hudi的表类型,COW和MOR的区别要搞清楚,COW写的时候会复制整个文件,读的时候是完整的文件,适合读多写少的场景;MOR写的时候只写增量文件,读的时候合并增量和基础文件,适合写多读少的场景。第二个是主键和预合并字段,主键是保证数据唯一性的,预合并字段是用来解决数据冲突的,比如同一主键有两条数据,会保留预合并字段值大的那条。第三个是检查点路径,这个一定要配置,不然任务失败后重启会从最早的位置重新读取数据,导致数据重复。
四、整合过程中的避坑要点
很多人在整合的时候会碰到各种坑,下面把最常见的坑和解决方法说清楚。
4.1 版本不兼容的坑
这是最常见的坑,Spark和Hudi的版本不匹配会出现类冲突的错误,比如Spark 3.2.x要搭配Hudi 0.10.x,Spark 3.3.x要搭配Hudi 0.12.x以上。解决方法是去Hudi的官网查版本兼容表,严格按照兼容表来选版本。
4.2 数据重复的坑
数据重复的原因主要有两个,一个是检查点路径配置错了,或者检查点路径被删了,导致任务重启后重新读取数据;另一个是Kafka的偏移量配置错了,比如把startingOffsets设成了earliest,任务重启后会重新读取所有数据。解决方法是检查点路径要配置在可靠的存储上,比如HDFS,不要配置在本地磁盘;Kafka的偏移量要设成latest,或者配置成从检查点恢复偏移量。
4.3 小文件太多的坑
小文件太多会导致查询速度变慢,甚至会导致HDFS的NameNode压力过大。小文件产生的原因主要是微批的间隔太短,每批次的数据量太小。解决方法是调整微批的间隔,比如从10秒改成30秒,或者调整每批次读取的最大数据量,比如从1000条改成5000条;另外可以配置Hudi的小文件合并参数,比如hoodie.parquet.small.file.limit,把小文件的阈值设成128MB,小于这个阈值的文件会被合并。
4.4 延迟过高的坑
延迟过高的原因主要有两个,一个是微批的间隔太长,比如设成了1分钟;另一个是Hudi的写操作太慢,比如配置了太多的同步操作。解决方法是调整微批的间隔,根据业务的实时性要求来设,比如业务要求延迟在10秒以内,就设成10秒;另外可以配置Hudi的异步写参数,比如hoodie.write.async,开启异步写,提高写的速度。
五、整合后的验证与维护
整合完成后,要验证数据是否正确写入,还要定期维护,保证系统的稳定运行。
5.1 数据验证的方法
验证数据的方法有两种,一种是用Spark SQL直接查询Hudi表,看数据的条数、内容是否正确;另一种是用Hudi的命令行工具,比如hudi-cli,来查看表的元数据、文件分布等。
5.2 日常维护的要点
日常维护主要包括几个方面,第一是监控数据的延迟,比如用Prometheus+Grafana来监控Spark任务的延迟、Hudi表的写延迟;第二是监控小文件的数量,定期合并小文件;第三是备份检查点路径,防止检查点丢失;第四是定期清理过期的Hudi版本,节省存储空间。
六、文章总结
Hudi和Spark Structured Streaming的整合,是实现实时流式数据入湖的常用方案,能解决实时处理、数据存储、后续分析的需求。在整合的过程中,要注意版本的匹配、参数的配置、避坑要点的处理,还要定期维护,保证系统的稳定运行。只要掌握了这些要点,就能顺利实现实时流式数据入湖,为后续的数据分析、业务决策提供支撑。
Comments