Apache Spark是目前最流行的大数据计算框架之一,我用Spark做大数据处理也有三年了。这三年来,我用Spark做过数据清洗、特征工程、机器学习、实时计算等各种任务,踩了很多坑,也总结了很多经验。本文是我使用Spark三年的一些感悟和总结,包括性能调优、内存管理、数据倾斜处理、稳定性保障等方面的经验和教训。如果你也在用Spark做大数据开发,希望这篇文章能帮你少走弯路。

一、我和Spark的缘分

先说说我是怎么开始用Spark的吧。

三年前,公司的业务快速增长,数据量越来越大,原来用MySQL和Python脚本处理数据已经跟不上了。一个几亿行的表,用Python处理要跑好几个小时,而且经常因为内存不足而失败。

这时候我们开始调研大数据处理框架。当时Hadoop MapReduce已经比较成熟了,但是Spark更火,因为它基于内存计算,速度比MapReduce快很多,而且API更友好,支持Scala、Java、Python、R等多种语言。

我们花了一个月的时间搭建Spark集群,学习Spark的用法,然后把原来的Python数据处理脚本迁移到Spark上。迁移完成之后,原来要跑几个小时的任务,现在十几分钟就能跑完,效率提升了几十倍。

从那以后,Spark就成了我们大数据处理的主力框架。这三年来,我们的Spark集群从最初的3个节点发展到现在的几十个节点,数据量从TB级增长到PB级,应用场景也从简单的数据清洗扩展到机器学习、实时计算、图计算等。

在这个过程中,我踩了无数的坑。有因为配置不当导致任务跑几天几夜跑不完的,有因为数据倾斜导致某个Task跑几个小时的,有因为内存管理不当导致OOM的,有因为数据倾斜导致整个任务失败的。每一个坑都让我对Spark有了更深的理解。

下面就把这些经验和教训分享给大家。

二、性能调优:Spark的永恒主题

用Spark,性能调优是永远绕不开的话题。同样的逻辑,不同的写法,性能可能差几十倍甚至上百倍。

1. 合理设置并行度

并行度是Spark性能调优最基础也是最重要的参数。并行度太低,数据分布不均匀,有的Task数据量很大,有的Task很闲,资源利用率低。并行度太高,Task太多,调度开销大,反而影响性能。

Spark的并行度由两个参数控制:spark.default.parallelism(控制RDD的默认分区数)和spark.sql.shuffle.partitions(控制Spark SQL shuffle之后的分区数,默认是200)。

我的经验是,并行度一般设置为集群总CPU核数的2到3倍比较合适。比如集群有100个CPU核,那并行度设置在200到300之间比较合适。

但是也要根据数据量来调整。数据量小的时候,并行度不要太高,否则Task太多,调度开销大。数据量大的时候,并行度要足够大,否则每个Task处理的数据量太大,容易OOM。

还有一个技巧是,在shuffle操作之后,如果数据量变小了,可以用coalesce或者repartition来减少分区数,避免小文件过多和调度开销过大。

2. 用好缓存

Spark的缓存(cache/persist)是一个非常重要的功能。如果一个RDD或者DataFrame会被多次使用,一定要缓存起来,避免重复计算。

我刚开始用Spark的时候,不知道缓存的重要性,一个DataFrame用了好几次,每次都从头计算,浪费了很多时间。后来加了cache之后,性能提升了好几倍。

但是缓存也不能滥用。缓存太多数据会占用大量内存,导致GC压力大,甚至OOM。只缓存那些会被多次使用、计算成本高的数据。

缓存的存储级别也要选对。默认的MEMORYONLY是只存在内存里,如果内存不够就会丢失,下次使用的时候重新计算。如果数据比较大,可以用MEMORYANDDISK,内存不够的时候溢写到磁盘。如果数据序列化之后比较小,可以用MEMORYONLY_SER,序列化之后占用的内存更少,但是读取的时候需要反序列化,会有一点CPU开销。

还有一点,缓存之后要记得在不用的时候unpersist,释放内存。否则缓存的数据会一直占用内存,影响其他任务的执行。

3. 避免shuffle

shuffle是Spark中最昂贵的操作,因为它涉及到数据的跨节点传输,还有磁盘读写和网络传输。性能调优的一个重要方向就是尽量减少shuffle。

会产生shuffle的操作包括:reduceByKey、groupByKey、join、distinct、repartition等。这些操作在使用的时候要特别注意。

比如,能用reduceByKey就不要用groupByKey。reduceByKey会在map端先做一次聚合,减少shuffle的数据量。groupByKey不会在map端聚合,所有数据都要shuffle,性能差很多,而且容易导致内存溢出。

再比如,join的时候,如果一个表很小,可以用broadcast join,把小表广播到每个节点,避免shuffle。Spark SQL在小表join的时候会自动选择broadcast join,但是阈值可以通过spark.sql.autoBroadcastJoinThreshold来调整。

还有,在join之前,如果可以先过滤数据,就先过滤,减少参与shuffle的数据量。

4. 数据序列化

序列化对Spark性能的影响也很大。Spark默认用的是Java序列化,但是Java序列化的性能比较差,序列化之后的体积也比较大。

推荐使用Kryo序列化,它比Java序列化快很多,序列化之后的体积也小很多。配置方法是设置spark.serializerorg.apache.spark.serializer.KryoSerializer

使用Kryo的时候,最好注册你自己的类,这样性能会更好。可以通过spark.kryo.classesToRegister来注册类。

三、数据倾斜:Spark最头疼的问题

数据倾斜是Spark开发中最常见也最头疼的问题。数据倾斜就是说,数据分布不均匀,某个key的数据量特别大,导致对应的Task处理的数据量远大于其他Task,这个Task会成为整个任务的瓶颈,拖慢整个任务的执行速度,甚至导致OOM。

1. 如何发现数据倾斜

发现数据倾斜的方法很简单。看Spark UI的Stage页面,如果某个Stage的大部分Task都很快执行完了,但是有几个Task执行时间特别长,或者数据量特别大,那就是数据倾斜了。

也可以在代码中统计每个key的数据量,看看是不是有某些key的数据量远大于其他key。

2. 数据倾斜的解决方法

数据倾斜的解决方法有几种,根据不同的场景选择不同的方法。

第一种方法是过滤掉导致倾斜的key。如果导致倾斜的key是无效数据(比如null、空字符串),可以直接过滤掉,这些数据对结果没有影响。

第二种方法是增加随机前缀。对于导致倾斜的key,可以给它加上随机前缀,把一个大key拆分成多个小key,这样数据就分散到多个Task上了。聚合完成之后,再去掉前缀,做一次二次聚合。这种方法适合聚合类的操作。

第三种方法是扩容维度表。如果是join的时候发生数据倾斜,而且是小表join大表,可以给小表扩容,比如给小表的每一行都加上1到N的随机前缀,大表也加上对应的前缀,这样join的时候数据就分散了。

第四种方法是广播小表。如果join的两个表一个很大一个很小,可以把小表广播到每个节点,避免shuffle,从根本上解决数据倾斜。

第五种方法是单独处理倾斜的key。如果只有少数几个key导致倾斜,可以把这些key单独拿出来处理,和其他key分开计算,最后再合并结果。

数据倾斜没有万能的解法,需要根据具体的场景选择合适的方法。但是只要掌握了基本的思路,大部分数据倾斜问题都能解决。

四、内存管理:OOM是永远的痛

用Spark,OOM(内存溢出)是经常遇到的问题。Spark的内存管理比较复杂,理解它的内存模型才能有效避免OOM。

1. Spark的内存模型

Spark的内存主要分为两部分:执行内存和存储内存。执行内存用于shuffle、join、排序、聚合等计算过程中的临时数据存储。存储内存用于缓存数据和广播变量。

这两部分内存共享一个统一的区域,一方空闲的时候另一方可以借用。但是执行内存的优先级更高,存储内存被执行内存抢占之后,缓存的数据会被淘汰。

Spark的内存配置主要有两个参数:spark.executor.memory(每个Executor的总内存)和spark.memory.fraction(执行内存和存储内存占总内存的比例,默认是0.6)。

2. 常见的OOM原因和解决方法

第一种OOM是Driver端的OOM。Driver端OOM通常是因为收集了太多数据到Driver端,比如用collect把大量数据拉到Driver,或者用了太大的广播变量。解决方法是,不要随便用collect,如果需要查看数据用take或者show。广播变量不要太大,如果太大可以考虑用其他方式。

第二种OOM是Executor端的OOM。Executor端OOM的原因比较多。可能是单个Task处理的数据量太大,这时候需要增加并行度,让每个Task处理的数据量变小。可能是数据倾斜,某个Task处理的数据量远大于其他Task,这时候需要解决数据倾斜。可能是shuffle的时候数据量太大,这时候需要优化shuffle,减少shuffle的数据量。可能是缓存了太多数据,这时候需要减少缓存,或者用MEMORYANDDISK存储级别。

第三种OOM是shuffle的时候OOM。shuffle的时候,如果map端或者reduce端的数据量太大,就可能OOM。可以通过增加并行度、使用Kryo序列化、调整shuffle内存比例等方式来解决。

3. 内存调优的建议

内存调优的几个建议:第一,合理设置Executor内存,不要太小也不要太大。太小容易OOM,太大会导致GC时间长。一般每个Executor分配4到8G内存比较合适。第二,合理设置Executor的CPU核数,一般2到4个核比较合适。核数太多,多个Task共享内存,容易OOM。第三,使用Kryo序列化,减少内存占用。第四,及时释放不需要的数据,用完就unpersist。第五,避免在Executor端做大数据量的本地集合操作。

五、稳定性保障:让任务跑得稳

性能很重要,但是稳定性更重要。一个任务跑得再快,如果经常失败,那也没用。

1. 数据质量监控

大数据任务,数据质量是基础。如果输入数据有问题,后面的计算结果肯定也有问题。

所以要在任务的关键节点加入数据质量检查。比如检查数据量是否在合理范围内,检查关键字段是否有空值,检查数据格式是否正确,检查主键是否有重复等。如果发现数据异常,及时报警,避免错误的数据影响下游。

我们就遇到过一次,上游系统出了问题,数据量突然减少了90%,但是我们的任务没有检查,直接跑了,结果下游的报表数据全错了,影响了业务决策。从那以后,我们在所有关键任务中都加入了数据质量检查。

2. 任务重试和容错

Spark本身有一定的容错能力,Task失败了会自动重试。但是有些错误不是重试能解决的,比如数据倾斜导致的OOM,重试多少次都会失败。

对于重要的任务,要设置合理的重试次数和重试间隔。有时候因为网络抖动、节点故障等原因导致的临时失败,重试一次就能成功。但是重试次数不要太多,否则失败的时候会浪费很多时间。

还要做好任务的幂等性设计。任务失败重跑的时候,不会产生重复数据,不会影响结果的正确性。比如写入数据的时候用覆盖写,或者用主键去重。

3. 资源隔离

如果多个任务共享一个集群,要做好资源隔离。否则一个任务占用了太多资源,会影响其他任务的执行。

可以用YARN的队列来做资源隔离,不同的任务用不同的队列,每个队列分配一定的资源。重要的任务用独立的队列,避免被其他任务影响。

也可以用Spark的动态资源分配,让任务根据需要动态申请和释放资源,提高资源利用率。

4. 监控和告警

一定要有完善的监控和告警。监控任务的执行时间、数据量、成功率等指标,发现异常及时告警。

我们用的是Prometheus + Grafana来监控Spark任务,通过Spark的REST API采集指标。任务失败、执行时间过长、数据量异常等情况都会触发告警,通过企业微信通知到负责人。

有了监控和告警,问题就能及时发现,及时处理,不会等到业务方反馈才知道出了问题。

六、一些实用的技巧

最后分享一些实用的小技巧。

1. 用Spark SQL而不是RDD

如果是结构化数据处理,优先用Spark SQL(DataFrame/Dataset API),而不是RDD API。Spark SQL有Catalyst优化器,会自动优化执行计划,性能比手写RDD好很多,而且代码更简洁易读。

只有在需要做底层的、非结构化的处理时,才用RDD API。

2. 用好Spark UI

Spark UI是排查问题的利器。任务出问题的时候,第一时间去看Spark UI,看Stage的执行情况,看Task的分布,看GC时间,看shuffle的数据量。大部分问题都能在Spark UI上找到线索。

一定要保留Spark UI的日志,任务结束之后也能查看。可以用Spark History Server来查看历史任务的UI。

3. 小文件合并

Spark任务经常会产生大量小文件,小文件太多会影响后续任务的性能,也会给HDFS带来压力。

在写入数据的时候,可以用coalesce或者repartition来控制输出文件的数量,避免产生太多小文件。也可以定期做小文件合并,把大量小文件合并成少量大文件。

4. 合理使用广播变量

如果有一个比较小的数据集需要在所有Task中使用,可以用广播变量,把它广播到每个节点,避免每个Task都拉取一份。广播变量可以大大减少网络传输和内存占用。

但是广播变量不要太大,一般不超过1G。太大的话广播本身就很慢,而且会占用大量内存。

七、写在最后

用了三年Spark,我最大的感受是:Spark是一个强大但是复杂的框架。

它的强大在于,它能处理海量的数据,能做各种复杂的计算,能大大提高数据处理的效率。它的复杂在于,它有很多参数,很多机制,很多坑。要用好Spark,需要深入理解它的原理,需要在实践中不断踩坑和总结。

这三年来,我踩了很多坑,也解决了很多问题。每解决一个问题,我对Spark的理解就深了一层。现在回头看,那些曾经让我头疼的问题,都成了我最宝贵的经验。

本文分享的这些经验,都是我在实际项目中总结出来的,希望能帮助正在用Spark或者准备用Spark的朋友少走一些弯路。

当然,Spark的技术还在不断发展,新的版本、新的特性、新的优化手段层出不穷。我们也要保持学习的心态,不断更新自己的知识,跟上技术的发展。

最后用一句话结束本文:"大数据没有银弹,只有对原理的深刻理解和对细节的极致追求。"愿每一个大数据开发者都能驾驭Spark,让它为业务创造更大的价值。