AWS EMR 6.4.0运行Spark应用仅单任务执行,如何实现并行?
问题根因
你代码中硬编码指定了setMaster("local"),这是导致所有任务始终串行执行的核心原因:
local模式的含义是强制Spark仅使用单线程在本地运行,会完全忽略集群资源、你手动指定的分片数、以及spark-submit提交时的集群配置参数,所以不管你是本地测试还是提交到EMR集群,都只会有1个task执行,你设置的100个分片也只会在同一个线程里串行处理。
修复方案
1. 移除硬编码的Master配置
删除代码里写死的master指定,修改SparkConf初始化代码:
// 修改前(错误写法) SparkConf sparkConf = new SparkConf().setAppName("MyAppName").setMaster("local"); // 修改后(正确写法) SparkConf sparkConf = new SparkConf().setAppName("MyAppName");
Master地址不需要在业务代码里硬编码,提交到EMR集群时spark-submit会自动从集群环境获取对应配置。如果需要本地测试,可以在本地执行时通过启动参数传入--master local[*]来启用本地多线程运行即可。
2. (可选)调整提交参数优化并行效率
如果要最大化利用10台主机的集群资源,可以在提交命令中增加资源和并行度配置:
aws emr add-steps --profile myprofile --region us-east-1 --cluster-id j-SOMEID --steps Type=CUSTOM_JAR,Name=test-downloader,ActionOnFailure=CONTINUE,Jar=command-runner.jar,Args=spark-submit,--class,com.my.main.MyClass,--executor-cores,4,--num-executors,10,--conf spark.default.parallelism=200,s3://somebucket/my.packaged.app-1.0.jar,-arg1,some,more,cli,args
新增参数说明:
--executor-cores 4:每个Executor分配4核CPU,可同时并行处理4个task--num-executors 10:集群总共启动10个Executor,匹配10台主机的集群规模--conf spark.default.parallelism=200:设置全局默认并行度,避免后续RDD转换操作自动将并行度缩到过小
3. 验证效果
重新打包提交任务后,在Spark UI的Stages页面可以看到对应Stage的task数量会和你设置的分片数/并行度匹配,多个task会分散到不同Executor上并行执行,处理速度会有明显提升。
内容的提问来源于stack exchange,提问作者user361676
相关产品推荐
相关产品推荐

