修改分区后PySpark作业超时/OOM问题咨询
PySpark Delta表分区变更引发OOM问题分析
用户问题
我有一个PySpark作业,原本将数据写入按
year, month, day and hour分区的Delta表,作业耗时2小时,每日运行并导入前一天的完整数据。近期我将分区方式改为单一字段event_date(日期类型),之后作业开始出现内存不足(OOM)错误。我了解Spark任务数与分区数密切相关。改为按event_date分区后,是否意味着单日数据全部由单个任务写入,而之前则是按每日24个小时分区分配到24个任务处理?
解答
- 你的核心猜测是正确的,关键细节补充如下:
- 之前按
year, month, day, hour四级分区时,单日数据会被拆分到24个小时分区中。Spark写入时,会根据数据的分区键值,将对应数据分配到不同任务并行处理——单日数据分散到至少24个任务,每个任务仅负责对应小时分区的数据写入,内存压力被有效分摊。 - 改为
event_date单字段分区后,单日所有数据都会归入同一个分区目录。如果写入前数据的RDD/DataFrame分区数未做调整,Spark会把同属一个event_date的数据集中到少数(甚至单个)任务中处理。若单日数据量较大,单个任务需处理的数据量会远超之前单小时的量级,直接触发OOM错误。
- 之前按
- 额外补充:Spark的任务数不仅和目标表的分区数相关,还取决于写入前数据的分区数(可通过
df.rdd.getNumPartitions()查看)。如果写入前数据分区数大于目标分区数,Spark会将多个数据分区的内容合并到同一个目标分区的任务中处理,这也会加剧单个任务的数据处理压力。
内容的提问来源于stack exchange,提问作者steve
相关产品推荐
相关产品推荐

