用Spark做大数据处理有几年了,从最开始的入门到现在的深入,踩了不少坑。本文总结Spark大数据处理中常见的坑和对应的解决方案,包括数据倾斜、Shuffle调优、内存管理、序列化、性能调优等方面的实战经验。这些坑很多人都会遇到,希望我的总结能帮你少走弯路,提升Spark作业的性能和稳定性。
一、数据倾斜的坑
数据倾斜是Spark中最常见也最头疼的问题。当某个key的数据量特别大的时候,对应的task处理时间会特别长,成为整个作业的瓶颈。
1. 数据倾斜的表现
数据倾斜的表现很明显:
- 大部分task很快就完成了,但有几个task一直跑不完
- Spark UI上看到某个stage的task执行时间严重不均,有的几秒,有的几小时
- 作业经常因为某个task OOM而失败
- 数据量不大,但作业跑得很慢
如果你遇到了这些情况,大概率是数据倾斜了。
2. 数据倾斜的原因
数据倾斜的根本原因是key的分布不均匀。常见的原因:
- 业务数据本身就不均匀,比如某个热门商品的订单量是其他商品的几十倍
- 数据中有大量的null值或空字符串,这些值会被分到同一个partition
- join的时候,某个key的数据量特别大
- 自定义的partitioner分配不均匀
3. 数据倾斜的解决方案
针对不同的情况,有不同的解决方案。
方案一:过滤掉不需要的key
如果导致倾斜的key是不需要的(比如null值、无效数据),直接过滤掉就行。这是最简单的解决方案。
比如,在做join的时候,如果某个key是null,而且null值的join没有意义,直接filter掉null值。
方案二:给key加随机前缀
如果导致倾斜的key是需要的,可以给这个key加随机前缀,把一个大key拆分成多个小key。
具体做法:
- 把导致倾斜的key单独拿出来
- 给这些key加上随机前缀(比如0到N的随机数)
- 另一边的对应key也做同样的处理(扩大N倍)
- join之后再去掉前缀
这种方法能把一个大key的压力分散到多个task上,有效解决数据倾斜。但实现起来比较复杂,需要写额外的代码。
方案三:使用salting技术
salting技术和加随机前缀类似,是一种更通用的解决数据倾斜的方法。具体做法是在shuffle之前给key加上随机前缀,shuffle之后再去掉。
Spark 3.0之后,内置了AQE(Adaptive Query Execution),可以自动检测和处理数据倾斜。开启AQE之后,Spark会自动把倾斜的partition拆分成多个小task处理。
开启AQE的方式:
spark.sql.adaptive.enabled=true
spark.sql.adaptive.skewJoin.enabled=true如果用的是Spark 3.0以上版本,建议开启AQE,能自动解决大部分数据倾斜问题。
方案四:调整并行度
有时候数据倾斜是因为并行度不够,导致每个partition的数据量太大。可以适当增加并行度,让数据更均匀地分布。
但这种方法只能缓解轻度的数据倾斜,对于严重的数据倾斜效果有限。
方案五:广播小表
如果是join导致的数据倾斜,而且其中一张表比较小,可以用broadcast join,把小表广播到每个节点,避免shuffle,从根本上解决数据倾斜。
broadcast join的方式:
spark.sql("SELECT /*+ BROADCAST(t1) */ * FROM t1 JOIN t2 ON t1.key = t2.key")或者用DataFrame的broadcast函数:
t1.join(broadcast(t2), Seq("key"))小表广播是解决join数据倾斜的有效方法,但只适用于小表(一般在GB以内)。
二、Shuffle的坑
Shuffle是Spark中最耗性能的操作。Shuffle涉及到数据的读写、序列化、网络传输,很容易成为性能瓶颈。
1. Shuffle数据量太大
最常见的问题是Shuffle数据量太大,导致Shuffle时间很长。
原因和解决方案:
- 不需要的列没有提前过滤。在shuffle之前,用select只保留需要的列,减少数据量
- 不需要的行没有提前过滤。在shuffle之前,用filter过滤掉不需要的行
- 可以先做局部聚合再做全局聚合。比如,先在每个partition内做一次聚合,减少shuffle的数据量,再做全局聚合
- 使用map-side pre-aggregation。对于reduceByKey、aggregateByKey等操作,Spark会自动在map端做局部聚合,尽量用这些操作,而不是groupByKey之后再做聚合
2. Shuffle分区数不合理
Shuffle的分区数(spark.sql.shuffle.partitions)默认是200。这个值对于小数据量来说太多了,对于大数据量来说又太少了。
- 数据量小的时候,分区数太多会导致每个task处理的数据量太小,task调度开销大
- 数据量大的时候,分区数太少会导致每个task处理的数据量太大,容易OOM,处理时间长
建议根据数据量调整分区数。一般来说,每个task处理128MB-256MB的数据比较合适。可以用这个公式估算:分区数 = 数据量 / 256MB。
另外,开启AQE之后,Spark会自动调整shuffle分区数,根据数据量动态合并小分区。建议开启AQE,减少手动调整的麻烦。
3. Shuffle文件过多
Shuffle会产生大量的中间文件。每个map task会为每个reduce task生成一个文件,如果map task和reduce task都很多,shuffle文件的数量会非常大(m * n),导致文件系统压力大,甚至inode耗尽。
解决方案:
- 开启shuffle consolidation(spark.shuffle.consolidateFiles=true),让多个map task共享同一个shuffle文件,减少文件数量
- 适当减少分区数,减少shuffle文件的数量
- 使用更好的文件系统(比如HDFS),处理大量小文件的能力更强
4. Shuffle读取失败
有时候会出现shuffle读取失败的情况,常见的原因:
- Executor丢失,shuffle文件也跟着丢失了
- 网络问题,数据传输失败
- 磁盘满了,shuffle文件写不进去
解决方案:
- 开启shuffle文件的备份,或者开启外部shuffle服务(ExternalShuffleService),即使Executor丢失了,shuffle文件也不会丢
- 增加网络重试次数和超时时间
- 确保磁盘有足够的空间,监控磁盘使用率
- 对于重要的作业,可以开启checkpoint,把中间结果持久化,避免从头重算
三、内存管理的坑
Spark的内存管理比较复杂,配置不当很容易出问题。
1. OOM(内存溢出)
OOM是Spark中最常见的问题之一。OOM可能发生在Driver端,也可能发生在Executor端。
Driver OOM的原因和解决方案:
- collect()把大量数据拉到Driver端。避免对大数据集使用collect,只在数据量小的时候用
- Driver端维护了大量的状态(比如广播变量、累加器)。增加Driver内存(spark.driver.memory)
- 广播变量太大。把大的广播变量拆分,或者用join代替广播
Executor OOM的原因和解决方案:
- 每个task处理的数据量太大。增加分区数,减少每个task的数据量
- 数据倾斜,某个task处理的数据量特别大。解决数据倾斜
- 聚合操作在内存中维护了大量数据。用map-side预聚合减少内存压力,或者增加Executor内存
- 广播变量占用了太多内存。减少广播变量的大小,或者增加Executor内存
- UDF中创建了大量对象,导致内存泄漏。优化UDF,复用对象,避免内存泄漏
增加Executor内存(spark.executor.memory)是最直接的解决方案,但要注意,内存不是越大越好,太大的内存会导致GC时间长。一般来说,每个Executor分配4-8GB内存比较合适。
2. GC频繁或GC时间长
Spark运行在JVM上,GC对性能有很大影响。如果GC频繁或GC时间长,会导致task停顿,影响整体性能。
原因和解决方案:
- 内存太小,对象存活率高,导致频繁GC。增加Executor内存
- 创建了大量短生命周期的对象。优化代码,复用对象,减少对象创建
- 大对象过多,导致老年代频繁Full GC。优化数据结构,减少大对象
- GC算法不合适。JDK 8用G1 GC,JDK 11+可以用ZGC,低延迟
- 堆外内存使用不当。合理设置堆内内存和堆外内存的比例
监控GC的方式:在Spark UI的Executors页面,可以看到每个Executor的GC时间。如果GC时间占task时间的比例超过10%,就需要优化了。
3. 堆外内存溢出
Spark除了堆内内存,还会使用堆外内存(比如shuffle、序列化、网络传输)。如果堆外内存不够,也会OOM。
常见原因:
- shuffle数据量太大,堆外内存不够
- 序列化缓冲区太大
- 直接使用了堆外内存(比如Netty)
解决方案:
- 增加堆外内存(spark.executor.memoryOverhead)。默认是executor内存的10%,如果shuffle数据量大,可以适当增加
- 减少每个task处理的数据量
- 优化序列化,减少序列化缓冲区的大小
四、序列化的坑
序列化在Spark中无处不在,shuffle、缓存、广播都需要序列化。序列化的性能对整体性能有很大影响。
1. 默认序列化器性能差
Spark默认用Java序列化器(JavaSerializer),但Java序列化器的性能很差,序列化后的字节数组大,序列化和反序列化的速度慢。
建议用Kryo序列化器,性能比Java序列化器好很多。Kryo序列化后的字节数组更小,序列化和反序列化的速度更快。
开启Kryo的方式:
spark.serializer=org.apache.spark.serializer.KryoSerializer使用Kryo的时候,最好注册自定义的类,这样Kryo可以用更高效的方式序列化:
spark.kryo.classesToRegister=com.example.MyClass1,com.example.MyClass2如果不注册,Kryo也能序列化,但会把类名存进去,增加序列化后的大小。
2. 自定义类没有实现Serializable
如果用Java序列化器,自定义的类必须实现Serializable接口,否则会报NotSerializableException。
用Kryo序列化器的话,不需要实现Serializable接口,但建议还是实现,因为有些地方(比如广播变量)可能还是用Java序列化。
如果遇到NotSerializableException,检查一下:
- 自定义的类有没有实现Serializable
- 有没有在闭包中引用了不可序列化的对象(比如数据库连接、文件句柄)
- 有没有在RDD操作中使用了不可序列化的函数
解决方案:
- 让自定义类实现Serializable
- 对于不可序列化的对象,在map/filter等操作中创建,不要在闭包中引用
- 用transient标记不需要序列化的字段
3. 缓存的数据序列化方式不当
Spark缓存RDD的时候,可以选择序列化或不序列化。不序列化的话,访问速度快,但占用内存多。序列化的话,占用内存少,但访问的时候需要反序列化,速度慢一些。
如果内存充足,用不序列化的缓存(MEMORYONLY),访问速度快。如果内存不够,用序列化的缓存(MEMORYONLY_SER),节省内存。
另外,缓存的时候要注意数据量,如果缓存的数据量太大,会导致频繁GC甚至OOM。只缓存需要反复使用的数据,不需要的不要缓存。
五、Spark SQL的坑
Spark SQL是最常用的API之一,但也有很多坑。
1. 小文件问题
用Spark SQL写数据的时候,如果分区数太多,会产生大量小文件。小文件会导致:
- 文件系统压力大,inode消耗快
- 读取的时候需要打开很多文件,性能差
- 元数据管理压力大
解决方案:
- 写入之前repartition或coalesce,减少分区数
- 用INSERT OVERWRITE的时候,设置合理的分区数
- 开启AQE的自动合并小文件功能(spark.sql.adaptive.coalescePartitions.enabled=true)
- 定期合并小文件,用OPTIMIZE命令(如果用Delta Lake)或自己写合并脚本
2. 类型转换问题
Spark SQL的类型转换有时候会出问题。比如,字符串转数字的时候,如果字符串不是合法的数字,会返回null而不是报错。这可能导致数据丢失而不自知。
另外,不同数据源的类型映射也可能有问题。比如,从MySQL读取的decimal类型,在Spark中可能变成了double,导致精度丢失。
解决方案:
- 显式指定schema,不要依赖自动推断
- 类型转换之后检查null值,确认转换是否成功
- 对于decimal等精确类型,用正确的类型,不要用double
- 写入的时候指定正确的类型,避免隐式转换
3. UDF性能差
Spark SQL的UDF(用户自定义函数)很方便,但性能比内置函数差很多。因为UDF不能被Spark的优化器优化,而且每次调用都有序列化和反序列化的开销。
如果UDF用在大数据量的场景,性能会很差。
解决方案:
- 尽量用内置函数,不要用UDF。Spark SQL的内置函数很丰富,大部分需求都能满足
- 如果必须用UDF,用Scala或Java写,不要用Python写(Python UDF更慢)
- 用更高阶的函数(比如aggregate、transform、filter)代替UDF
- 对于复杂的逻辑,可以考虑用DataSet API,类型安全,性能也比UDF好
4. Join顺序和策略
Spark SQL的join优化很重要。join的顺序和策略会影响性能。
常见问题:
- 大表join大表,shuffle数据量大
- join顺序不合理,导致中间结果太大
- 没有用broadcast join,导致不必要的shuffle
优化建议:
- 小表join大表,用broadcast join,避免shuffle
- 多张表join的时候,先join小表,减少中间结果,再join大表
- 开启AQE,Spark会自动选择join策略和调整join顺序
- 对于大表join大表,确保join key分布均匀,避免数据倾斜
六、其他常见的坑
1. 分区数和并行度不合理
Spark的并行度由分区数决定。分区数太少,并行度不够,资源利用不充分;分区数太多,task调度开销大,每个task处理的数据量太小。
建议:
- 根据数据量和集群资源设置合理的并行度
- 读取数据的时候,根据文件大小和数量设置分区数
- shuffle之后,用spark.sql.shuffle.partitions设置分区数
- 开启AQE,自动调整分区数
2. 数据本地性差
Spark的数据本地性(PROCESSLOCAL、NODELOCAL、RACK_LOCAL、ANY)会影响性能。如果数据和计算不在同一个节点,就需要网络传输,性能差。
常见原因:
- 数据存储在HDFS上,但Spark任务没有调度到数据所在的节点
- 动态分配资源,Executor启动慢,任务等不及就调度到了其他节点
- 数据倾斜,某些节点的数据量特别大
优化建议:
- 增加数据本地性等待时间(spark.locality.wait),让任务等一下数据所在的节点
- 关闭动态分配,固定Executor数量,确保数据本地性
- 确保数据分布均匀,避免数据倾斜
- 用HDFS的时候,确保副本数足够,提高数据本地性的概率
3. 动态资源分配的问题
Spark的动态资源分配(Dynamic Allocation)可以根据负载自动调整Executor数量,节省资源。但在某些场景下,动态分配会导致问题:
- Executor频繁启动和销毁,开销大
- Executor丢失导致shuffle文件丢失,任务失败
- 数据本地性差,因为Executor是动态启动的
如果遇到这些问题,可以考虑关闭动态分配,固定Executor数量。或者调整动态分配的参数,比如设置最小和最大Executor数,增加Executor空闲超时时间。
4. checkpoint和持久化的选择
Spark的cache/persist和checkpoint都可以持久化中间结果,但有区别:
- cache/persist把数据存在内存或磁盘,速度快,但生命周期和application绑定,application结束就没了
- checkpoint把数据存在可靠的存储(比如HDFS),速度慢一些,但可以跨application使用,而且可以截断RDD血缘,用于故障恢复
建议:
- 对于需要在同一个application中反复使用的数据,用cache/persist
- 对于需要跨application使用的数据,或者血缘很长需要截断的,用checkpoint
- 重要的中间结果,用checkpoint持久化,避免故障后从头重算
- cache和checkpoint可以结合使用,先cache再checkpoint,兼顾速度和可靠性
七、性能调优的一般步骤
最后,总结一下Spark性能调优的一般步骤。
1. 监控和定位瓶颈
调优的第一步是找到瓶颈。用Spark UI监控:
- 哪个stage最慢?
- 哪个task最慢?是不是数据倾斜?
- GC时间占比高不高?
- Shuffle读写数据量大不大?
- 数据本地性怎么样?
- 有没有task失败或重试?
找到瓶颈之后,再有针对性地优化。
2. 数据层面优化
- 过滤不需要的数据,减少处理的数据量
- 只选择需要的列,避免select *
- 提前过滤,减少shuffle数据量
- 确保数据分布均匀,避免数据倾斜
3. 资源层面优化
- 合理设置Executor数量、内存、CPU核数
- 合理设置并行度和分区数
- 确保数据本地性
- 开启动态资源分配(如果适用)
4. 代码层面优化
- 避免在闭包中引用大对象或不可序列化的对象
- 用map-side预聚合减少shuffle数据量
- 用broadcast join代替普通join
- 尽量用内置函数,少用UDF
- 合理使用cache和checkpoint
5. 配置层面优化
- 用Kryo序列化器
- 开启AQE,自动优化
- 合理设置shuffle分区数
- 优化JVM和GC参数
- 开启shuffle consolidation
八、写在最后
Spark是一个强大的大数据计算框架,但要用好它并不容易。从数据倾斜到Shuffle调优,从内存管理到序列化优化,每一个环节都有很多需要注意的地方。
本文总结的这些坑,都是我在实际项目中踩过的。可能不是最全面的,但都是实战中验证过有效的。希望能帮助正在用Spark的你,少踩一些坑,提升作业的性能和稳定性。
当然,Spark的技术在不断发展,新的版本会带来新的功能和优化。比如Spark 3.0的AQE,就能自动解决很多以前需要手动调优的问题。建议尽量用新的版本,开启AQE,能省去很多手动调优的麻烦。
最后,我想说,性能调优是一个持续的过程。没有一劳永逸的优化方案,需要不断地监控、分析、调整。但只要掌握了正确的方法和思路,就能在遇到问题的时候快速定位和解决。
用一句话结束本文:"Spark调优没有银弹,只有深入理解原理,结合实际场景,不断监控和优化,才能发挥它的最大潜力。"愿每一个用Spark的开发者,都能驾驭好这个强大的大数据引擎。
评论(0)
暂无评论,快来抢沙发~
评论功能仅对会员开放,请先登录
登录