上一篇文章,我推荐了一些Spark大数据工具,很多朋友反馈说,很有帮助。但也有朋友说,工具是有了,但Spark的性能优化和进阶用法,还是不太清楚。今天,就来聊聊Spark大数据的进阶技巧,这些技巧,可能很多人都不知道,但掌握了之后,能让你的Spark程序,跑得更快、更稳、更高效。

先说说背景。我做Spark大数据开发,已经好几年了,从最开始的Spark 1.x,到现在的Spark 2.x、3.x,踩了很多坑,也积累了很多经验。最开始,我写Spark程序,也只是会用基础的API,写个WordCount,写个简单的ETL,性能差的时候,就知道加内存、加CPU,不知道怎么优化。

后来,随着项目越来越复杂,数据量越来越大,基础的写法,根本满足不了需求,程序经常跑几个小时,甚至跑挂。于是,我开始深入研究Spark的原理和优化技巧,踩了很多坑,也总结了很多经验。慢慢地,我写的Spark程序,性能越来越好,稳定性也越来越高。

今天,就把这些进阶技巧,分享出来,涵盖了Spark SQL、DataFrame、RDD、性能优化、稳定性、运维等各个方面。

一、Spark SQL进阶技巧

Spark SQL,是Spark最常用的模块,大部分人,用Spark,都是用Spark SQL。但很多人,用Spark SQL,只是写个简单的SELECT、FROM、WHERE,不知道很多进阶的用法和优化技巧。

1. 用好Catalyst优化器

Spark SQL的核心,是Catalyst优化器,它会对你写的SQL,进行一系列的优化,比如,谓词下推、列裁剪、常量折叠、投影合并等。很多人,不知道Catalyst的存在,写SQL的时候,不注意,导致Catalyst无法优化,性能很差。

用好Catalyst,有几个技巧:

  • 尽量用DataFrame/Dataset API或者SQL,不要用RDD:因为RDD的操作,Catalyst是无法优化的,而DataFrame/Dataset和SQL,Catalyst可以进行全方位的优化。所以,除非是特殊场景,否则,尽量用DataFrame/Dataset或者SQL,不要用RDD。
  • 避免在WHERE子句里,对列做函数操作:比如,WHERE dateformat(createtime, 'yyyy-MM-dd') = '2020-01-01',这样写,Catalyst无法对createtime做谓词下推和索引优化,因为列被函数包裹了。应该改成,WHERE createtime >= '2020-01-01' AND create_time < '2020-01-02',这样,Catalyst就能优化了。
  • 用EXPLAIN查看执行计划:写完SQL之后,用EXPLAIN命令,查看执行计划,看看Catalyst是怎么优化你的SQL的,有没有全表扫描,有没有数据倾斜,有没有不必要的Shuffle。如果发现执行计划,不符合预期,就调整SQL,让Catalyst能更好地优化。

2. 用好窗口函数

窗口函数,是Spark SQL非常强大的功能,可以解决很多复杂的统计问题,比如,分组排名、Top N、累计求和、移动平均等。很多人,遇到这些问题,不知道用窗口函数,而是用复杂的子查询、自连接,性能很差,代码也很难懂。

窗口函数的基本语法是:

函数名(列) OVER (
  PARTITION BY 分组列
  ORDER BY 排序列
  ROWS BETWEEN 起始行 AND 结束行
)

比如,计算每个部门,工资最高的前3名员工:

SELECT * FROM (
  SELECT 
    name, department, salary,
    ROW_NUMBER() OVER (PARTITION BY department ORDER BY salary DESC) as rn
  FROM employees
) t WHERE rn <= 3

用窗口函数,几行SQL,就能解决复杂的统计问题,性能也很好,因为Spark对窗口函数,做了专门的优化。

常用的窗口函数有:ROWNUMBER、RANK、DENSERANK(排名),SUM、AVG、MAX、MIN(聚合),LAG、LEAD(取前后行),FIRSTVALUE、LASTVALUE(取首尾值)等。

掌握了窗口函数,很多复杂的SQL问题,都能迎刃而解。

3. 用好Broadcast Join

Join,是Spark SQL最常见的操作,也是最容易出性能问题的操作。Spark的Join,有几种实现方式:Broadcast Join、Sort Merge Join、Shuffle Hash Join等。其中,Broadcast Join,是性能最好的,因为它不需要Shuffle,把小表广播到每个节点,然后,在每个节点,本地Join。

用好Broadcast Join,有几个技巧:

  • 小表Join大表,用Broadcast Join:如果Join的两个表,一个很小(比如,小于10MB),一个很大,一定要用Broadcast Join,把小表广播出去。可以用broadcast()函数,显式指定广播小表:

``scala largeTable.join(broadcast(smallTable), "id") ` 或者,在SQL里,用广播提示: `sql SELECT /+ BROADCAST(t2) / * FROM t1 JOIN t2 ON t1.id = t2.id ``

  • 调整spark.sql.autoBroadcastJoinThreshold:Spark默认,会自动广播小于10MB的表,可以通过调整这个参数,让更多的小表,自动广播。比如,设置成100MB,这样,小于100MB的表,都会自动广播。但要注意,不要设置得太大,不然,广播大表,会占用很多内存,反而影响性能。
  • 避免大表Join大表:如果两个表,都很大,就没法用Broadcast Join了,只能用Sort Merge Join,需要Shuffle,性能会差很多。所以,尽量避免大表Join大表,可以考虑,在数据导入的时候,就把两个表关联好,存成宽表,查询的时候,直接查宽表,不需要Join。

4. 用好缓存

如果一个DataFrame或者表,会被多次使用,一定要缓存起来,避免重复计算。Spark的缓存,有几种级别:MEMORYONLY、MEMORYANDDISK、MEMORYONLYSER、MEMORYANDDISKSER等。

用好缓存,有几个技巧:

  • 缓存会被多次使用的数据:比如,一个基础的DataFrame,会被多次查询、多次Join,一定要缓存起来,这样,第一次计算之后,后面的查询,直接用缓存的数据,不需要重新计算,性能会提升很多。
  • 选择合适的缓存级别:如果数据量不大,内存够,用MEMORYONLY,性能最好;如果数据量比较大,内存不够,用MEMORYANDDISK,内存不够的时候,溢写到磁盘;如果想节省内存,可以用序列化的缓存级别(SER),数据会被序列化,占用更少的内存,但读取的时候,需要反序列化,有一点CPU开销。
  • 及时释放不用的缓存:缓存的数据,不用了之后,要及时unpersist(),释放内存,避免占用太多内存,影响其他计算。
  • 缓存之后,先触发一次Action:缓存是懒执行的,调用cache()之后,不会立即缓存,要等第一次Action之后,才会真正缓存。所以,缓存之后,最好先触发一次Action(比如,count()),确保数据被缓存了,后面的查询,才能用到缓存。

二、性能优化进阶技巧

性能优化,是Spark开发中,最重要、最有技术含量的部分。很多人,写Spark程序,性能差,就知道加资源,不知道怎么优化。其实,大部分性能问题,都是可以通过优化代码和参数,来解决的。

1. 解决数据倾斜

数据倾斜,是Spark最常见、最头疼的性能问题。数据倾斜,就是某个Task处理的数据量,远大于其他Task,导致这个Task特别慢,整个Stage,都在等它。

解决数据倾斜,有几个常用的方法:

  • 增加并行度:如果是分区数太少,导致的数据倾斜,可以增加分区数,让数据更分散。比如,用repartition(),增加分区数,或者,调整spark.sql.shuffle.partitions参数,增加Shuffle之后的分区数。
  • 两阶段聚合:对于聚合操作导致的数据倾斜,可以用两阶段聚合,先给key加随机前缀,做局部聚合,然后,去掉前缀,再做全局聚合。这样,把倾斜的key,打散到多个Task,避免单个Task,处理太多数据。
  • Broadcast Join:对于Join操作导致的数据倾斜,如果小表不大,可以用Broadcast Join,避免Shuffle,也就避免了数据倾斜。
  • 拆分倾斜key:如果是少数几个key,导致的数据倾斜,可以把这几个key,单独拿出来处理,然后,再和其他key的结果,合并起来。
  • 加盐:给倾斜的key,加上随机前缀,把数据打散到多个Task,处理完之后,再去掉前缀,聚合结果。

我之前,处理过很多次数据倾斜,这些方法,都用过,效果都不错。关键是,要先通过Spark Web UI,找到是哪个Stage、哪个key导致的倾斜,然后,根据具体情况,选择合适的方法。

2. 优化Shuffle

Shuffle,是Spark性能的瓶颈,因为Shuffle需要,把数据在节点之间传输,开销很大。大部分Spark性能问题,都和Shuffle有关。

优化Shuffle,有几个技巧:

  • 减少Shuffle的次数:尽量减少会导致Shuffle的操作,比如,reduceByKey、groupByKey、join、distinct、orderBy等。能合并的操作,尽量合并,减少Shuffle的次数。比如,多次聚合,可以合并成一次聚合;多次Join,可以考虑,先把数据合并,再一次Join。
  • 减少Shuffle的数据量:在Shuffle之前,尽量过滤掉不需要的数据,只保留需要的列和行,减少Shuffle的数据量。比如,用filter(),先过滤掉不需要的行;用select(),只选择需要的列。这样,Shuffle的时候,传输的数据量,就小很多,性能会提升很多。
  • 调整Shuffle的并行度:Shuffle之后的分区数,由spark.sql.shuffle.partitions参数控制,默认是200。如果数据量很大,200个分区,每个分区的数据量,就会很大,Task处理起来,就会很慢。可以适当调大这个参数,比如,设置成500、1000,让每个分区的数据量,小一些,Task处理起来,更快。但也不要设置得太大,不然,分区太多,调度开销大,反而影响性能。
  • 使用高效的Shuffle实现:Spark的Shuffle,有几种实现方式,比如,Sort Shuffle、Hash Shuffle等。现在,默认用的是Sort Shuffle,性能比较好,稳定性也高。一般,不需要改,但可以了解一下,根据具体情况,选择合适的Shuffle实现。
  • 调整Shuffle的内存比例:Shuffle的时候,需要内存来存储和排序数据,如果内存不够,就会溢写到磁盘,影响性能。可以调整spark.shuffle.memoryFraction参数,增加Shuffle的内存比例,减少溢写。但要注意,不要设置得太大,不然,其他计算的内存,就不够了。

3. 优化内存使用

Spark是内存计算框架,内存的使用,对性能影响很大。很多Spark程序,性能差,甚至OOM,都是因为内存使用不合理。

优化内存使用,有几个技巧:

  • 合理设置Executor内存:根据集群的资源和任务的需求,设置合适的Executor内存,不要太小,也不要太大。太小,容易OOM;太大,浪费资源,而且,GC时间长。一般,每个Executor,设置4-8GB内存,比较合适,具体,根据任务的需求调整。
  • 调整内存比例:Spark的内存,分为存储内存和执行内存,可以通过spark.memory.fraction和spark.memory.storageFraction参数,调整两者的比例。如果任务,缓存的数据多,就增加存储内存的比例;如果任务,计算复杂,Shuffle多,就增加执行内存的比例。
  • 使用序列化:序列化的数据,比反序列化的数据,占用更少的内存。可以用Kryo序列化,比Java序列化,更省内存,更快。注册需要序列化的类,让Kryo更高效。
  • 避免数据倾斜导致的内存压力:数据倾斜,会导致某个Task,处理的数据量太大,内存不够,OOM。解决了数据倾斜,也就解决了内存压力。
  • 及时释放不用的数据:不用的DataFrame和缓存,及时unpersist(),释放内存;不用的对象,及时置为null,让GC回收。避免内存里,存了很多不用的数据,占用空间。

4. 优化GC

Spark是JVM程序,GC(垃圾回收)对性能影响很大。如果GC时间太长,会导致Task执行慢,甚至,Executor失联。

优化GC,有几个技巧:

  • 使用G1垃圾回收器:G1垃圾回收器,比默认的Parallel GC,更适合大内存、多核心的场景,GC停顿时间更短。可以通过spark.executor.extraJavaOptions参数,设置使用G1垃圾回收器。
  • 调整年轻代和老年代的比例:如果任务,创建的临时对象多,年轻代可以设置大一些,减少对象晋升到老年代,减少Full GC。
  • 减少对象创建:在代码里,尽量减少不必要的对象创建,比如,循环里,不要创建大量的临时对象;用基本类型,代替包装类型;用数组,代替集合等。对象创建少了,GC的压力,就小了。
  • 使用序列化缓存:缓存数据的时候,用序列化的缓存级别,数据以序列化的形式,存在内存里,占用的空间小,而且,不需要GC的时候,扫描大量的对象,GC效率更高。
  • 监控GC时间:通过Spark Web UI,查看Task的GC时间,如果GC时间,占Task执行时间的比例很高(比如,超过10%),就说明,GC有问题,需要优化。

三、稳定性进阶技巧

除了性能,稳定性,也是Spark开发中,非常重要的部分。很多Spark程序,性能还行,但经常跑挂,或者,结果不准确,这就是稳定性的问题。

1. 处理数据质量问题

大数据,数据质量,经常有问题,比如,缺失值、异常值、格式错误、重复数据等。如果不处理这些问题,会导致程序报错,或者,结果不准确。

处理数据质量问题,有几个技巧:

  • 做好数据校验:在数据处理的流程里,加上数据校验的步骤,检查数据的完整性、一致性、准确性、唯一性等。比如,检查必填字段,有没有空值;检查数值字段,有没有异常值;检查日期字段,格式对不对;检查ID,有没有重复。
  • 处理缺失值:对于缺失值,根据具体情况,处理掉,比如,用默认值填充、用平均值/中位数填充、删除缺失的行、单独标记等。不要让缺失值,参与计算,导致结果错误。
  • 处理异常值:对于异常值,根据具体情况,处理掉,比如,删除、用最大值/最小值替换、单独标记等。不要让异常值,影响统计结果。
  • 去重:对于重复数据,要去重,避免重复计算。可以用dropDuplicates(),或者,用窗口函数,去重。
  • 记录脏数据:处理掉的脏数据,最好记录下来,存到一个单独的地方,方便后续排查和分析,看看数据质量,有没有问题,是数据源的问题,还是处理流程的问题。

2. 处理数据倾斜导致的稳定性问题

前面说了,数据倾斜,会导致性能问题,严重的时候,还会导致稳定性问题,比如,某个Task,处理的数据量太大,内存不够,OOM,程序跑挂。

所以,解决数据倾斜,不仅是性能优化,也是稳定性优化。具体的方法,前面已经说了,这里就不重复了。

3. 合理设置重试和超时

Spark程序,运行的时候,可能会遇到各种瞬时错误,比如,网络抖动、节点故障、磁盘满等。这些错误,可能会导致Task失败,甚至,程序跑挂。

合理设置重试和超时,可以提高程序的稳定性:

  • 设置Task重试次数:通过spark.task.maxFailures参数,设置Task的最大失败次数,默认是4。如果Task失败了,Spark会自动重试,重试几次,都失败了,才会宣告程序失败。这样,一些瞬时错误,重试一下,就成功了,不会导致整个程序失败。
  • 设置Executor重试次数:通过spark.yarn.maxAppAttempts(YARN模式)参数,设置Application的最大重试次数,如果整个Application失败了,会自动重试。
  • 设置网络超时:通过spark.network.timeout参数,设置网络超时时间,如果网络有问题,超时之后,会自动重试,或者,宣告失败。
  • 设置心跳超时:通过spark.executor.heartbeatInterval参数,设置Executor心跳间隔,如果Executor,长时间没有心跳,会被认为失联,重新调度Task。

合理设置这些参数,可以让程序,更健壮,遇到瞬时错误,能自动恢复,不会轻易跑挂。

4. 做好监控和告警

Spark程序,运行的时候,要做好监控和告警,及时发现问题,及时处理。

监控和告警,有几个方面:

  • 监控Spark程序的运行状态:通过Spark Web UI、YARN UI、Prometheus + Grafana等,监控Spark程序的运行状态,比如,Job的进度、Stage的耗时、Task的失败率、数据量、内存使用、GC时间等。
  • 监控数据质量:监控数据的质量,比如,数据量是否正常、有没有缺失值、有没有异常值、有没有重复数据等。如果数据质量,有异常,及时告警。
  • 监控资源使用:监控集群的资源使用,比如,CPU、内存、磁盘、网络等,如果资源使用异常,及时告警。
  • 设置告警:对于关键的指标,设置告警,比如,Task失败率超过阈值、程序运行时间超过阈值、数据量异常、资源使用异常等,及时通知相关人员,及时处理。

做好监控和告警,就能在问题发生的早期,发现问题,及时处理,避免问题扩大,导致严重的后果。

四、写在最后

今天,分享了Spark大数据的进阶技巧,涵盖了Spark SQL、性能优化、稳定性等各个方面,包括用好Catalyst优化器、窗口函数、Broadcast Join、缓存,解决数据倾斜、优化Shuffle、优化内存和GC,处理数据质量问题、合理设置重试和超时、做好监控和告警等。

这些技巧,都是我在实际工作中,踩了很多坑,才总结出来的,每一个,都能实实在在地,提升Spark程序的性能和稳定性。希望能帮助大家,更好地使用Spark,写出高性能、高稳定的Spark程序。

当然,Spark的进阶技巧,远不止这些,还有很多,比如,Spark Streaming的优化、Structured Streaming的优化、MLlib的优化、GraphX的优化、Spark on Kubernetes、Spark 3.x的新特性等,以后有机会,再继续分享。

最后,想说的是,Spark的优化,没有银弹,没有哪个技巧,是万能的。关键是,要深入理解Spark的原理,根据具体的场景和问题,选择合适的优化方法。而且,优化是一个持续的过程,要不断地监控、分析、优化,让程序,跑得越来越快,越来越稳。

希望这篇文章,能对大家有所帮助。如果有什么问题或者不同的看法,欢迎在评论区交流。