Spark使用repartitionByRange处理5000万行数据时Driver内存过高问题咨询
内存占用过高的核心影响因素
- Range分区采样的基础内存开销:Spark执行
repartitionByRange前需要拉取所有上游分区的采样数据来生成分区边界,默认配置spark.sql.execution.rangeExchange.sampleSizePerPartition为100,即每个上游分区采样100条数据。若上游RDD分区数达到数万级别,总采样数据量可达数百万条。每条采样数据包含hash(Int)、salt(Double)、采样权重三个字段,加上Java对象头、对齐填充的额外开销,单个采样记录的内存占用是原始数据类型的3~5倍,仅采样数据本身就会占用数GB内存。 - Driver端排序的临时内存放大:所有采样数据拉取到Driver后,需要全量排序计算分区边界。排序过程需要至少2倍于采样数据大小的临时内存,用于数据交换、临时数组存储,会进一步放大内存占用。
- 高分区数的边界存储开销:若你设置的
partitionsNumber过大(比如超过10万),生成的分区边界数组会包含partitionsNumber - 1个条目,每个条目存储两个排序键,仅边界数组的内存占用就可达数百MB到数GB。 - 无重复键的内存无优化空间:你用
rand()生成salt作为第二排序键,所有采样键几乎完全唯一,没有重复值的压缩空间,进一步推高了内存占用。
相关参考代码
preExportRdd.toDF .withColumn("provider", $"provider_id") .withColumn("hash", hash($"provider", $"qk10")) .withColumn("salt", rand()) .repartitionByRange(partitionsNumber, $"hash", $"salt") .drop("hash", "salt") .write .partitionBy(partitionColumns:_*) .format("parquet") .option("compression", "gzip") .mode(SaveMode.Append) .save(exportUrl)
优化建议
- 降低采样比例:调整
spark.sql.execution.rangeExchange.sampleSizePerPartition配置到20~50,只要数据分布相对均匀,降低采样量不会显著影响分区均衡性。 - 替换分区逻辑:如果仅需要打散数据避免小文件,用
repartition(partitionsNumber, $"hash", $"salt")替代repartitionByRange,哈希分区不需要Driver端采样排序,不会产生大额Driver内存开销。 - 减少上游分区数:如果上游RDD分区数过高,先执行
coalesce减少分区数后再做范围分区,可大幅降低总采样数据量。
内容的提问来源于stack exchange,提问作者jk1
相关产品推荐
相关产品推荐

