Apache Spark CPU性能无法线性扩展问题排查
为什么Spark并行素数测试无法实现CPU性能线性扩展?
测试场景说明
我正在编写Apache Spark基准测试对比不同服务器的CPU性能,使用并行素数生成脚本,可设置1、64、128或256个分区,并通过spark-submit的--total-executor-cores参数限制Executor核心数,执行命令如下:
./bin/spark-submit --total-executor-cores x --master spark://server:7077 generate_primes_x.py
素数生成逻辑简单,预期性能接近线性扩展,但实际结果并未达到预期。测试脚本如下(numSlices对应设置的分区数):
from pyspark.sql import SparkSession def is_prime(n): if n <= 1: return False if n <= 3: return True if n % 2 == 0 or n % 3 == 0: return False i = 5 while i * i <= n: if n % i == 0 or n % (i + 2) == 0: return False i += 6 return True if __name__ == "__main__": spark = SparkSession.builder \ .appName("ParallelPrimeGenerator_256") \ .getOrCreate() sc = spark.sparkContext upper_limit = 10000000 # Create an RDD with numbers from 2 to upper_limit numbers_rdd = sc.parallelize(range(2, upper_limit + 1), numSlices=256) # You can adjust the number of slices # Filter the prime numbers using the is_prime function prime_numbers_rdd = numbers_rdd.filter(is_prime) # Collect the prime numbers and print them prime_numbers = prime_numbers_rdd.collect() print("Prime numbers up to", upper_limit, ":", prime_numbers) spark.stop()
无法线性扩展的原因分析
- 计算负载严重不均衡:素数判断的计算量随数字大小呈非线性增长,大数字的素数验证耗时远高于小数字。而
sc.parallelize默认按连续区间划分分区,导致前序分区集中小数字,计算很快完成,后序分区全是大数字,耗时极长。任务的整体完成时间由最慢的分区决定,额外核心无法加速这部分瓶颈任务。 - PySpark的跨进程序列化开销:PySpark中Python UDF与JVM之间需要频繁进行数据序列化/反序列化,每个数字在JVM和Python进程间传递都会产生额外开销。核心数越多,这种跨进程通信的总开销占比越高,会逐步抵消并行计算带来的性能收益。
- 分区与核心的匹配度问题:若分区数远大于核心数,会导致Spark频繁进行任务调度切换,增加调度开销;若分区数少于核心数,则会有部分核心处于闲置状态。此外,过多的分区会让任务调度的时间成本上升,进一步拉低整体性能。
- Collect操作的串行瓶颈:脚本最后调用
collect()将所有结果拉取到Driver节点,这一步是串行执行的。当核心数增加、生成的素数数量增多时,该步骤的IO和内存压力会急剧上升,成为整体性能的瓶颈,无法体现CPU的线性扩展能力。 - CPU缓存效率下降:大数字的素数判断需要更多计算步骤,且内存访问模式不利于CPU缓存利用。当核心数增多时,缓存竞争加剧,缓存命中率下降,导致每个核心的实际计算效率降低,无法保持线性增长。
内容的提问来源于stack exchange,提问作者Jamesla
相关产品推荐
相关产品推荐

