Spark大规模集群中count、distinct及distinct count性能缓慢问题求助
Spark 集群统计任务性能优化指南
嘿,我来帮你捋捋这个Spark性能卡壳的问题——20台m4.xlarge的集群处理20GB、近2亿条记录,总记录数统计都要7分钟,distinct类任务更慢,确实有点离谱。结合我踩过的各种Spark性能坑,给你几个实打实的优化方向:
一、先搞定基础配置的“坑”
- 资源分配要到位:m4.xlarge单台有4vCPU、16GB内存,别用默认的小Executor配置!建议给每个Executor分配
10G内存+3核CPU,每台机器跑3个Executor,这样20台机器就能开出60个Executor。提交任务时可以这么设置:
确保集群资源不闲置,这是性能提升的基础。spark-submit --executor-memory 10G --executor-cores 3 --num-executors 60 ... - 赶紧换列式存储格式:如果你的源文件是CSV、JSON这类文本格式,转成Parquet或ORC绝对是立竿见影的优化。列式存储不仅压缩比高(20GB文本转Parquet可能只剩5-8GB),还能大幅降低IO开销,count和distinct的速度直接翻倍都不奇怪。转换代码示例:
// Scala版本 spark.read.csv("/path/to/your/source/file").write.parquet("/path/to/parquet/output") - 调整分区数匹配并行度:默认分区数大概率不适合你的集群,导致任务跑不起来。读文件时可以用
repartition把初始分区调整到每个分区128MB-256MB(Spark最优分区大小);shuffle分区数(spark.sql.shuffle.partitions)建议设为Executor总核心数的2-3倍,比如60个Executor×3核=180,设成300左右就很合适。
二、针对三个统计任务的精准优化
1. 总记录数(count)优化
- 如果你用的是
df.count(),要是已经转成Parquet/ORC,记得开启元数据统计:
这样Spark直接从文件元数据拿count值,不用扫全量数据,速度快到离谱。--conf spark.sql.parquet.enableVectorizedReader=true \ --conf spark.sql.statistics.histogram.enabled=true - 别画蛇添足!如果只是单纯统计总记录数,别在count之前加没必要的过滤、映射操作,直接读文件就count。
2. 总唯一记录数优化
- 别直接用
df.distinct().count()!全列distinct会把所有数据都shuffle一遍,开销超大。改成只对整行的哈希值做distinct:
只传输哈希值,数据量直接砍到原来的几十分之一,shuffle压力骤降。df.selectExpr("hash(*) as row_hash").distinct().count() - 要是数据有主键/唯一标识列,直接用该列的distinct count代替全列distinct,性能提升更明显。
3. FundamentalSeriesId列的唯一记录数优化
- 业务允许的话,用
approx_count_distinct代替countDistinct!这个函数用HyperLogLog算法,不需要全量shuffle,误差率大概5%,但速度能提升10倍以上。代码示例:df.select(approx_count_distinct("FundamentalSeriesId")).show() - 必须要精确统计的话,除了调整shuffle分区数,记得开启自适应执行(Spark 2.4+支持):
它会根据实际数据量自动调整shuffle分区数和任务并行度,避免资源浪费。--conf spark.sql.adaptive.enabled=true
三、其他锦上添花的优化
- 缓存复用数据:如果三个统计任务是连续执行的,先把DataFrame缓存到内存:
df.cache(),这样后面的任务不用重复读文件,省不少IO时间。20GB的文件缓存后大概占10-15GB内存,你的集群总内存完全够撑。 - 检查硬件瓶颈:m4.xlarge如果用的是EBS磁盘,换成gp3高IO磁盘,避免磁盘拖后腿;另外确保集群在同一个VPC内,网络带宽足够,shuffle时的网络传输别卡壳。
先从文件格式转换和资源配置调整入手,这两个是最容易见效的,然后再针对具体任务做细节优化,应该能把耗时降到合理范围。
内容的提问来源于stack exchange,提问作者Atharv Thakur
相关产品推荐
相关产品推荐

