Spark任务数与spark.default.parallelism设置不符相关问题咨询
问题解答
1. 为何设置spark.default.parallelism=300后实际任务数为147
spark.default.parallelism参数仅在两类场景生效:
- 用
sc.parallelize等接口从本地内存创建RDD,且没有显式指定分区数时 - 执行
reduceByKey、join等触发shuffle的算子,且没有显式指定shuffle分区数时
你当前的业务逻辑使用sc.textFile读取HDFS数据,该接口的分区数默认由HDFS输入文件的总块数决定,和spark.default.parallelism没有直接关系:你的输入文件在HDFS上总共对应147个块,所以读取生成的RDD初始分区数就是147。后续的filter、map都属于窄依赖算子,不会修改分区数,saveAsTextFile也不会触发shuffle,所以整个Stage的任务数始终为147。
如果需要调整为300个分区,可选择两种方案: - 给
textFile传入第二个参数指定最小分区数:sc.textFile("hdfs://master:9000/path-to-input", 300) - 读取完成后手动重分区:
lines.repartition(300)
2. 任务编号后的.1后缀含义
该后缀代表任务的重试次数,计数从0开始:
.0代表该分区的第一次执行.1代表该分区的第二次重试执行.2代表该分区的第三次重试执行
Spark默认开启任务失败自动重试机制,由参数spark.task.maxFailures控制最大重试次数(默认值为4),出现重试通常是因为对应分区的前序执行出现了异常(如executor故障、网络超时、读取数据失败等)。你日志中出现的133.1、146.2分别对应对应分区的第二次、第三次重试执行。
内容的提问来源于stack exchange,提问作者Yuejiang_Li
相关产品推荐
相关产品推荐

