Spark使用mapPartitions作业数多于分区数及日志异常问题咨询
Spark mapPartitions重复执行、日志丢失问题解决方案
一、作业数远多于分区数问题
核心原因
- Spark RDD惰性计算机制:如果
yearMonthQueryRDD后续被多个行动算子(比如count、collect、save等)调用,或者没有做持久化,每次行动都会重新触发整个计算链路,mapPartitions里的提交作业逻辑会被重复执行。 - Spark任务重试机制:默认
spark.task.maxFailures配置为4,也就是单个分区任务最多重试3次,每次重试都会重新执行mapPartitions的所有逻辑,导致同个分区重复提交作业。 - 分区逻辑不符合预期:如果按日期列分区时用了哈希分区,极端情况下同个日期值可能被分配到多个分区,或者你实际的分区数和预期不一致。
- 代码隐藏bug:你的代码中
partitionSize变量没有累加逻辑,打印的总大小永远为0,另外每个分区内的代码是单线程执行的,不需要使用AtomicLong、AtomicReference这类线程安全类,直接用普通var即可。
修复方案
- 对处理后的RDD立即做持久化,避免重算:
import org.apache.spark.storage.StorageLevel val yearMonthQueryRDD = yearMonthQueryDF.rdd.mapPartitions( // 你的原有逻辑 ).persist(StorageLevel.MEMORY_AND_DISK_SER)
- 仅调用一次行动算子触发计算,比如只调一次
yearMonthQueryRDD.count(),不要重复调用行动算子。 - 把作业提交逻辑改成幂等:提交SQS或者调用lambda时,用
TaskContext.getPartitionId()+ 分区日期值作为唯一幂等键,就算重复提交也不会生成重复作业。 - 提前确认分区数符合预期:执行
println(yearMonthQueryDF.rdd.getNumPartitions)打印实际分区数,和你预期的日期分区数量对比。 - 修复代码中的
partitionSize累加逻辑,在遍历记录时添加partitionSize += recordSize。
二、mapPartitions日志无法找到问题
核心原因
Spark Executor默认的日志配置中,slf4j等日志框架的INFO级别日志默认输出到stderr,stdout仅收集println打印的内容,很多集群默认只会保留stderr的错误日志,或者你配置的日志级别过低导致INFO日志被过滤。
修复方案
- 临时调试时把日志级别改成ERROR,ERROR级别的日志一定会输出到stderr,可以快速验证逻辑是否执行:
logger.error("partitionIndex = {}",TaskContext.getPartitionId()); logger.error("No.of Elements in Partition [{}] are =[{}] Total Size=[{}]", paritionColumnName.get(), partitionCounter.get(), partitionSize.get())
- 调整Executor的log4j配置,把你业务代码对应的包日志级别设置为INFO,指定日志输出路径,确保日志被收集。
- 临时调试可以用
println打印关键信息,内容会输出到stdout,方便快速定位。 - 如果你用的是YARN/K8s集群,可以配置Spark参数保留Executor日志:
spark.executor.logs.rolling.maxRetainedFiles=10,避免日志被滚动删除。
内容的提问来源于stack exchange,提问作者Sateesh K
相关产品推荐
相关产品推荐

