未分区每日追加数据是否引发PySpark任务maxResultSize报错?
问题解答
结论:是,未分区的追加式存储就是触发报错的根本原因
1. 任务数递增的核心原因
- 你的Hadoop输入数据是每日追加、未分区的,哪怕代码里加了过滤最新一天数据的条件,Spark没有分区元数据的指引,会默认扫描全部历史数据文件
- 随着天数增加,历史文件总量越来越多,Spark为了处理这些文件启动的executor任务数就会逐日递增(每个小文件对应一个任务)
2. 触发maxResultSize报错的逻辑
- 每个executor任务执行完成后,会将任务状态、计算中间结果等数据回传给driver
- 当任务数累积到一定量级时,这些回传数据的总大小会超过
spark.driver.maxResultSize的限制(默认1G),哪怕你实际只需要最新一天的4.5G数据,全量扫描带来的任务回传数据已经撑爆了driver的结果内存阈值 - 调大该参数只是临时缓解,随着历史数据持续增加,迟早会再次触发报错
可行解决方案
- 根本解决(优先):将Hadoop输入数据改为按日期分区存储(例如按
dt=yyyy-MM-dd的分区格式),Spark会通过分区裁剪直接定位到最新一天的分区文件,无需扫描全量历史数据,任务数会稳定在处理当日4.5G数据对应的数量 - 临时过渡方案:
- 如果文件名带有日期标识,可通过
pathGlobFilter直接指定最新一天的文件名,跳过全量扫描 - 调大
spark.driver.maxResultSize的同时,同步提升spark.driver.memory给driver分配更多内存,但这只是短期权宜之计
- 如果文件名带有日期标识,可通过
内容的提问来源于stack exchange,提问作者chris9759
相关产品推荐
相关产品推荐

