You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 07:39:20