Sparksql参数调优
异常调优
兼容性问题与规避方案
spark.sql.hive.convertMetastoreParquet
Parquet 是一种列式存储格式,可用于 Spark SQL 和 Hive 的存储格式。在 Spark 中,如果使用 USING PARQUET 的形式创建表,则创建的是 Spark 的 DataSource 表;而如果使用 STORED AS PARQUET 则创建的是 Hive 表。
默认设置为 true,表示使用 Spark SQL 内置的 Parquet 读写器(即进行反序列化和序列化),具有更好的性能。如果设置为 false,则使用 Hive 的序列化方式。
然而,有时当设置为 true 时,可能会出现使用 Hive 查询表有数据,而使用 Spark 查询为空的情况。
但在某些情况下,将 spark.sql.hive.convertMetastoreParquet 设为 false 时,可能会发生以下异常(Spark 2.3.2):

这是因为当设置为 false 时,使用 Hive Metastore 的元数据读取数据,而如果此表是使用 Spark SQL DataSource 创建的 Parquet 表,其数据类型可能出现不一致的情况。例如,通过 MetaStore 读取到的是 IntWritable 类型,创建了一个 WritableIntObjectInspector 来解析数据,而实际上值是 LongWritable 类型,因此出现了类型转换异常。
与该参数相关的另一个参数是 spark.sql.hive.convertMetastoreParquet.mergeSchema,如果也是 true,那么将会尝试合并各个 Parquet 文件的 schema,以产生一个兼容所有 Parquet 文件的 schema。
精度丢失问题
spark.sql.decimalOperations.allowPrecisionLoss
当该参数为 true(默认),表示允许丢失精度,会根据 Hive 行为和 SQL ANSI 2011 规范来决定结果类型,即如果无法精确表示,则舍入结果的小数部分。
当该参数为 false 时,表示不允许丢失精度,这样会将数据表示得更加精确。
spark.sql.decimalOperations.nullOnOverflow
对于 Decimal 类型,由于其整数部分位数是 (precision - scale),因此该类型能表示的范围是有限的,一旦超出这个范围,就会发生溢出。在 Spark 中,如果 Decimal 计算发生溢出,默认会返回 NULL 值。
引入了参数 spark.sql.decimalOperations.nullOnOverflow 用来控制在 Decimal 操作发生溢出时的处理方式。
遇到表路径下的文件缺失/损坏异常
spark.sql.files.ignoreMissingFiles 和 spark.sql.files.ignoreCorruptFiles
这两个参数只有在进行 Spark DataSource 表查询时才有效,如果是对 Hive 表进行操作则无效。
在进行 Spark DataSource 表查询时,可能会遇到非分区表中的文件缺失/损坏或分区表分区路径下的文件缺失/损坏异常,这时设置这两个参数会忽略这些异常。这两个参数默认都是 false,建议在线上可以都设为 true。
其源码逻辑如下,简单描述就是如果遇到 FileNotFoundException,如果设置了 ignoreMissingFiles=true 则忽略异常,否则抛出异常;如果不是 FileNotFoundException 而是 IOException(FileNotFoundException 的父类)或 RuntimeException,则认为文件损坏,如果设置了 ignoreCorruptFiles=true 则忽略异常。

分区路径下的文件不存在或损坏的处理
上面的两个参数在分区表情况下是针对分区路径存在的情况下,分区路径下的文件不存在或损坏的处理。而有另一种情况是这个分区路径都不存在了。这时异常信息如下:

spark.sql.hive.verifyPartitionPath
参数默认是 false,当设置为 true 时会在获得分区路径时对分区路径是否存在做一个校验,过滤掉不存在的分区路径,这样就会避免上面的错误。
spark.files.ignoreCorruptFiles 和 spark.files.ignoreMissingFiles
这两个参数和上面的 spark.sql.files.ignoreCorruptFiles 很像,但区别很大。在 Spark 进行 DataSource 表查询时 spark.sql.files.* 才会生效,而如果查询的是一张 Hive 表,则会走 HadoopRDD 这条执行路线。
所以即使设置了 spark.sql.files.ignoreMissingFiles,仍然可能报 FileNotFoundException 的情况,异常栈如下:

此时可以将 spark.files.ignoreCorruptFiles 和 spark.files.ignoreMissingFiles 设为 true,其代码逻辑和上面的 spark.sql.file.* 逻辑没有明显区别,此处不再赘述。
Spark Shuffle 几个常用相关参数

spark.shuffle.file.buffer
- 默认值: 32k
- 描述: Shuffle write 端写磁盘文件时缓冲区大小,适量增大可以减少磁盘 I/O 次数,进而提升性能。
spark.reducer.maxSizeInFlight
- 默认值: 48M
- 描述: Shuffle read 端拉取对应分区数据缓冲区大小,适量增大可以减少网络传输次数,进而提升性能。
spark.shuffle.io.maxRetries
- 默认值: 3
- 描述: Shuffle read 端拉取对应数据时,因网络异常拉取失败重新尝试的最大次数。针对超大数据量的应用,可以增大重试次数,大幅度提升稳定性。
spark.maxRemoteBlockSizeFetchToMem
- 描述: 代表可以从远端拉取数据放入内存的最大 size。这个参数作用是在 reduce task 读取 map task block 时放入内存中的最大值;默认是没有限制全放内存。
spark.shuffle.sort.bypassMergeThreshold
- 描述: map 端不进行排序的分区阈值。
spark.shuffle.io.retryWait
- 描述: 重试 2 次的最大等待时间。
spark.shuffle.spill.compress
- 默认值: true
- 描述: 默认使用
spark.io.compression.codec。
注: 根据不同的 Spark 版本有不同的个别 Shuffle 配置如 spark.maxRemoteBlockSizeFetchToMem 等,根据不同的 Spark 版本,查询对应的功能,如果详细查看逻辑,查看源码。
Spark SQL CBO 相关几个参数

性能调优
除了遇到异常需要被动调整参数之外,我们还可以主动调整参数从而对性能进行调优。
spark.hadoopRDD.ignoreEmptySplits
- 默认值: false
- 描述: 如果是
true,则会忽略那些空的 splits,减小 task 的数量。
spark.hadoop.mapreduce.input.fileinputformat.split.minsize
- 描述: 用于聚合 input 的小文件,用于控制每个 mapTask 的输入文件,防止小文件过多时产生太多的 task。
spark.sql.autoBroadcastJoinThreshold 和 spark.sql.broadcastTimeout
- 描述: 用于控制在 Spark SQL 中使用 BroadcastJoin 时表的大小阈值,适当增大可以让一些表走 BroadcastJoin,提升性能,但如果设置太大又会造成 driver 内存压力。
broadcastTimeout是用于控制 Broadcast 的 Future 的超时时间,默认是 300s,可根据需求进行调整。
spark.sql.adaptive.enabled 和 spark.sql.adaptive.shuffle.targetPostShuffleInputSize
- 描述: 该参数用于开启 Spark 的自适应执行,这是 Spark 较老版本的自适应执行,后面的
targetPostShuffleInputSize是用于控制之后的 shuffle 阶段的平均输入数据大小,防止产生过多的 task。
Intel 大数据团队开发的 adaptive-execution 相较于目前 Spark 的 AE 更加实用,该特性也已经加入到社区 3.0 之后的 roadmap 中,令人期待。
spark.sql.parquet.mergeSchema
- 默认值: false
- 描述: 当设为
true,Parquet 会聚合所有 Parquet 文件的 schema,否则是直接读取 Parquet summary 文件,或者在没有 Parquet summary 文件时随机选择一个文件的 schema 作为最终的 schema。
spark.hadoop.mapreduce.fileoutputcommitter.algorithm.version
- 默认值: 1
- 描述: 1 或者 2,默认是 1。MapReduce-4815 详细介绍了
fileoutputcommitter的原理,实践中设置了version=2的比默认version=1的减少了 70% 以上的 commit 时间,但 1 更健壮,能处理一些情况下的异常。
spark.sql.files.maxPartitionBytes
- 默认值: 128MB
- 描述: 单个分区读取的最大文件大小。
Spark AQE 相关
Spark AQE 自动分区合并
spark.sql.adaptive.enabled: Spark 3.0 AQE 开关,默认是false,如果要使用 AQE 功能,得先设置为true。spark.sql.adaptive.coalescePartitions.enabled: 动态缩小分区参数,默认值是true,但得先保证spark.sql.adaptive.enabled为true。spark.sql.adaptive.coalescePartitions.initialPartitionNum: 任务刚启动时的初始分区,此参数可以设置得大点,默认值与spark.sql.shuffle.partition一样为 200。spark.sql.adaptive.coalescePartitions.minPartitionNum: 进行动态缩小分区,最小缩小至多少分区,最终分区数不会小于此参数。spark.sql.adaptive.advisoryPartitionSizeInBytes: 缩小分区或进行拆分分区操作后所期望的每个分区的大小(数据量)。
Spark AQE 动态 Join
spark.sql.adaptive.enabled: 是否开启 AQE 优化。spark.sql.adaptive.join.enabled: 是否开启 AQE 动态 Join 优化。spark.sql.adaptive.localShuffleReader.enabled: 在不需要进行 shuffle 重分区时,尝试使用本地 shuffle 读取器。将 sort-merge join 转换为广播 join。spark.sql.autoBroadcastJoinThreshold: 中间文件尺寸总和小于广播阈值。spark.sql.adaptive.nonEmptyPartitionRatioForBroadcastJoin: 空文件占比小于配置项。
Spark AQE 自动倾斜处理
spark.sql.adaptive.enabled: 开启 AQE 功能,默认关闭。spark.sql.adaptive.skewJoin.enabled: 开启 AQE 倾斜 Join,需要先将spark.sql.adaptive.enabled设置为true。spark.sql.adaptive.skewJoin.skewedPartitionFactor: 倾斜因子,如果分区的数据量大于此因子乘以分区的中位数,并且也大于spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes,那么认为是数据倾斜的,默认值为 5。spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes: 每个分区的阈值,默认 256MB,此参数应该大于spark.sql.adaptive.advisoryPartitionSizeInBytes。spark.sql.adaptive.advisoryPartitionSizeInBytes: 缩小分区或进行拆分分区操作后所期望的每个分区的大小(数据量)。
Spark AQE 动态申请资源
spark.sql.adaptive.enabled: 是否开启 AQE 优化。spark.dynamicAllocation.enabled: 是否开启动态资源申请。spark.dynamicAllocation.shuffleTracking.enabled: 是否开启 shuffle 状态跟踪。为执行程序启用随机文件跟踪,从而无需外部随机服务即可动态分配。此选项将尝试保持为活动作业存储随机数据的执行程序。spark.dynamicAllocation.shuffleTracking.timeout: 启用随机跟踪时,控制保存随机数据的执行程序的超时。默认值意味着 Spark 将依靠垃圾回收中的 shuffle 来释放执行程序。如果由于某种原因垃圾回收无法足够快地清理随机数据,则此选项可用于控制执行者何时超时,即使它们正在存储随机数据。
Spark SQL 参数表(Spark 2.3.2)


end
