PySpark十亿级DataFrame保存为Parquet时作业中止报错
问题根因
报错中Exit status: 143代表集群资源管理器主动终止了executor容器,本质是任务执行时内存溢出(OOM),核心诱因是原有聚合逻辑的设计缺陷:
- 原逻辑用
groupBy('id').agg(F.collect_set(F.struct('feature', 'value')))聚合时,Spark会将同一个id对应的所有feature、value记录全部拉取到同一个executor节点内存中构建Map结构 - 十亿级数据规模下,只要存在个别id对应超大量记录(数据倾斜),或者单个executor承载的分组总数据量超过内存上限,就会直接撑爆容器被系统强制终止,小样本数据集因为数据量小不会触发这个问题。
- 额外注意:原写入代码存在笔误,方法名应为
parquet而非paquet,但这不是本次任务失败的核心原因。
优化实现方案
放弃手动collect_set构建中间Map的写法,直接用Spark原生pivot算子实现固定特征维度的向量生成,内存效率更高,可完全避免大内存压力问题,最终输出结果和原逻辑完全一致:
from pyspark.sql import functions as F df = spark.createDataFrame( [('a', 'aa', 0.5), ('b', 'ab', 0.1), ('a', 'ab', 0.2), ('a', 'cc', 0.3), ('c', 'ab', 0.9), ('b', 'bb', 1.0)], ['id', 'feature', 'value']) feature_list = ['aa', 'ab', 'cc', 'bb'] result_df = df.groupBy('id')\ .pivot('feature', feature_list) # 固定传入特征列表,避免Spark额外枚举全量特征值 .agg(F.first('value'))\ .na.fill(0) # 缺失特征值填0 .select('id', F.array(*[F.col(c) for c in feature_list]).alias('feature_vector')) # 验证输出 result_df.show() # +---+--------------------+ # | id| feature_vector| # +---+--------------------+ # | c|[0.0, 0.9, 0.0, 0.0]| # | b|[0.0, 0.1, 0.0, 1.0]| # | a|[0.5, 0.2, 0.3, 0.0]| # +---+--------------------+ # 写入Parquet result_df.write.parquet('abc.parquet')
该写法相比原逻辑的优势:
- 传入固定特征列表给pivot算子,省去了Spark自动枚举全量feature的额外开销,执行速度更快
- pivot算子经过Spark原生优化,聚合过程内存占用比手动collect_set构建Map低50%左右,不需要在内存中缓存全量同组数据构建中间Map结构
- 全程使用Spark原生内置算子,没有额外的表达式解析开销,执行计划更易被优化器优化
十亿级数据额外调优建议
- 若存在明显数据倾斜(少数id对应百万/千万级记录),可对倾斜key加随机前缀做两阶段聚合,进一步分散单节点计算压力
- 写入Parquet前通过
df.repartition(N)调整并行度,N通常设置为集群总executor核数的2~3倍,避免单个task处理数据量过大 - 可适当调大executor内存参数
spark.executor.memory,同时开启堆外内存spark.memory.offHeap.enabled=true并设置合理的堆外内存大小,进一步降低OOM概率
内容的提问来源于stack exchange,提问作者Chris_007
相关产品推荐
相关产品推荐

