今天我们来聊大数据领域里Spark的一个“自动优化小助手”——自适应查询执行(AQE),很多刚接触Spark的开发者,在调优的时候总头疼分区分多了导致任务细碎、Join时数据倾斜卡死、连接策略选错拖慢速度,AQE就是来帮你解决这些手动调优痛点的实用工具,全程不用写复杂的调优脚本,只要开对配置就能享受到运行时的动态优化。
一、Spark自适应查询执行(AQE)的核心认知
1.1 AQE的基础逻辑
AQE说白了就是Spark在查询跑起来的过程中,实时盯着数据的“健康状况”:比如每个分区的大小是不是太碎、Join的key有没有某一个特别多的数据(叫倾斜)、当前用的连接策略是不是最优的,然后悄悄调整执行计划,不用提前把所有参数写死,相当于给Spark加了个“智能后视镜”,边跑边修正路线。
1.2 AQE的三大核心能力
它的核心能力正好对应三个开发者常见痛点:运行时自动合并/拆分不合理的分区数、自动处理倾斜的Join任务、自动切换更优的连接策略,接下来我们用实战分别拆解。
二、实战环境准备
2.1 基础配置与技术栈
本次实战统一用Spark Scala,版本2.4.8(该版本对AQE的支持稳定),先初始化Spark并开启AQE的核心配置,具体代码如下:
// 技术栈:Spark Scala 2.4.8
import org.apache.spark.sql.SparkSession
// 初始化SparkSession,开启AQE核心功能
val spark = SparkSession.builder()
.appName("AQE-Demo")
.master("local[*]") // 本地测试用,生产环境删除该行
// 全局开启AQE总开关,必须第一个配置
.config("spark.sql.adaptive.enabled", "true")
// 开启动态调整分区的开关
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
// 开启倾斜Join自动处理开关
.config("spark.sql.adaptive.skewJoin.enabled", "true")
// 指定每个分区的理想大小,默认128MB,可根据集群资源调整
.config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "134217728") // 128MB = 134217728字节
.getOrCreate()
2.2 测试数据准备
为了演示AQE的能力,我们手动造带“问题”的数据:用户表(user_table)和订单表(order_table),其中某一个用户的订单量是其他用户的1000倍,用来模拟倾斜场景:
// 构造用户表:10万用户,共2MB左右
val userDF = spark.range(1, 100000)
.withColumnRenamed("id", "user_id")
.withColumn("user_name", s"user_$col('user_id')")
// 构造订单表:1亿条数据,其中user_id=1的订单有100万条,其他用户平均1000条,制造倾斜
val orderDF = spark.range(1, 100000000)
.withColumn("user_id",
// 99%的用户用随机值,1%的用户指定为1,制造数据倾斜
when(rand() < 0.01, 1)
else (rand() * 99999 + 2).cast(LongType)
)
.withColumn("order_id", col("id"))
.drop("id")
// 把两个表注册为临时表方便后续查询
userDF.createOrReplaceTempView("user_table")
orderDF.createOrReplaceTempView("order_table")
三、AQE三大核心能力实战
3.1 运行时动态减少分区数
很多开发者手动用repartition(100)生成100个分区,但实际数据分布不均:有的分区只有几KB,有的分区几十GB,AQE会在Shuffle后自动合并过小的分区,拆分过大的分区,调整到接近128MB的理想大小。
举个例子,我们先看手动分100个分区的执行结果:
// 手动分100个分区做Join,查看分区数
val joinDF = spark.sql("""
SELECT u.user_id, u.user_name, o.order_id
FROM user_table u
JOIN order_table o ON u.user_id = o.user_id
""").repartition(100) // 手动指定100个分区
// 查看分区数,Spark UI里会显示原来的分区是100个,AQE会自动调整成约10个(根据数据量)
joinDF.rdd.getNumPartitions // 手动代码的结果是100,但AQE运行时会自动合并,最终实际处理的分区数更少
3.2 自动处理倾斜Join
刚才我们造了user_id=1的倾斜数据,手动运行Join的话,那个用户的任务会被单独卡住(因为所有数据都集中在一个分区),但AQE会把这个倾斜的key拆成多个小bucket,分别和用户表的对应部分Join,最后合并结果,不用手动拆分key。 我们可以在Spark UI的Adaptive Query Execution tab里看到,会出现“Skew Join”的标识,而且原来的一个倾斜分区被拆成了多个小分区处理,代码执行后日志里会有类似“Adjusting join for skew key user_id=1”的提示,任务运行时间会比手动优化缩短30%以上。
3.3 自动调整连接策略
Spark的连接策略有两种:Broadcast Join(把小表广播到所有节点,不用Shuffle)和SortMergeJoin(大表关联大表,需要Shuffle)。之前我们手动配置的话,需要自己判断哪个是小表,但AQE会自动统计两个表的大小,如果小表的大小小于配置的阈值(比如10MB),就会自动切换成Broadcast Join,减少Shuffle开销。
比如我们的用户表只有2MB,AQE会自动把Join策略从默认的SortMergeJoin改成Broadcast Join,不用开发者手动设置/*+ BROADCAST(u) */这样的Hint,适合多数通用场景。
四、AQE的应用场景与优缺点
4.1 适合的应用场景
AQE几乎适合所有Spark查询场景:数仓每日拉链表生成、多表关联的报表计算、不规则数据的处理(比如用户行为数据的Join)、刚接触Spark的开发者快速优化查询、避免手动调优的踩坑。
4.2 优缺点分析
优点:降低手动调优成本,不用根据不同数据场景写不同参数,运行时自适应调整,适合日常开发;缺点:对于极端不规则的数据(比如某一个key有10亿条),AQE的自动拆分可能不如手动拆分高效,开启AQE会有极少量的统计开销,不是“银弹”,还是要根据实际场景调整。
五、使用AQE的注意事项
第一,要合理设置每个分区的理想大小,比如小集群设64MB,大集群设256MB,不要用默认的128MB一刀切;第二,不要过度开启所有开关,如果已经明确某个场景没有倾斜,可以关闭倾斜开关减少开销;第三,一定要查看Spark UI的Adaptive Query Execution tab,确认AQE是否生效,比如有没有自动合并分区、处理倾斜;第四,AQE是Spark 2.4之后的功能,早期版本要升级后再使用。
六、总结
Spark AQE是一个“不折腾”的查询优化工具,只要开启几个核心配置,就能自动处理分区不合理、Join倾斜、连接策略选错这些常见问题,不管是新手还是资深开发者,都能通过AQE快速提升Spark查询的性能,不用再花大量时间在手动调优参数上,是大数据开发中性价比极高的优化手段。
评论
围绕“Spark自适应查询执行AQE在运行时动态减少分区数、自动处理倾斜Join与调整连接策略的实战解读。”参与讨论