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

Spark中稀疏向量聚合求和优化及Parquet写入问题排查

你的PySpark代码优化指南:解决Executor崩溃与低效计算问题

刚接触PySpark就碰到大规模数据的聚合和写入问题,确实容易踩坑。我帮你拆解下当前代码里的几个核心低效点,以及对应的优化方案,应该能解决你的ExecutorLostFailure和Shuffle过大的问题:

1. Python UDF是性能杀手,用原生函数替换

你当前用自定义UDF把稀疏向量转成数组,这是最大的性能瓶颈之一。Python UDF会触发JVM和Python进程之间的频繁序列化/反序列化,对于50M行的数据来说,这个开销会被放大到无法接受的程度。

优化方案:Spark 3.0+ 提供了原生的vector_to_array函数,可以直接将稀疏/稠密向量转换为数组,完全不需要UDF:

from pyspark.sql.functions import vector_to_array

# 替换原来的UDF逻辑
df_bayes = pipeline.fit(df_bayes).transform(df_bayes)
# 直接用vector_to_array转换稀疏向量
df_bayes = df_bayes.select('id', vector_to_array(col('KPvec')).alias('KPvec'))

如果你的Spark版本低于3.0,也可以利用Spark对向量类型的原生支持,直接跳过转数组的步骤(后面会讲)。

2. 数组逐个元素求和的方式完全错误

看你的执行计划里有5万多个sum操作,这是因为你用了[F.sum(F.col('KPvec')[i]) for i in range(len(kids))]——相当于把数组的每个元素都当成单独的列来求和,这会让Spark生成巨量的计算任务,直接导致Shuffle数据暴增到6GB,Executor负载过高崩溃。

优化方案:直接对稀疏向量进行聚合,Spark原生支持向量类型的sum操作,而且稀疏向量的求和只会计算非零元素,效率极高:

# 去掉转数组的步骤,直接对稀疏向量求和
df_bayes = pipeline.fit(df_bayes).transform(df_bayes)
# 按id分组,直接sum KPvec(稀疏向量类型)
df_bayes = df_bayes.groupby('id').agg(F.sum(col('KPvec')).alias('KPvec')).cache()

如果业务上必须要数组类型的结果,聚合后再用vector_to_array转一次即可,这比先转数组再求和高效10倍以上。

3. 分区策略不合理,导致多次无意义Shuffle

你先把数据repartition到15000个分区,之后写入又重新repartition,这会导致多次Shuffle操作。对于50M行的数据来说,15000个分区太多了——每个分区只有几万行,Shuffle的网络开销远大于计算开销。

优化方案:把分区数设置为集群CPU核心数的2-4倍(比如集群有100核,设置200-400个分区):

# 根据集群情况设置合理的分区数
num_partitions = 300
# 聚合前repartition到合适的分区数
df_bayes = df_bayes.repartition(num_partitions, col('id'))
# 聚合后缓存
df_bayes = df_bayes.groupby('id').agg(F.sum(col('KPvec')).alias('KPvec')).cache()

4. 写入方式错误,生成大量小文件导致崩溃

你用partitionBy("id")写入Parquet,如果聚合后有50M个不同的id,这会生成50M个小文件——IO开销直接拉满,Executor根本扛不住。

优化方案:放弃partitionBy("id"),用合理的分区数控制文件大小,同时可以设置参数限制每个文件的记录数:

df_bayes.repartition(num_partitions).write.format("parquet") \
    .option("path", "s3://...output.parquet") \
    .option("spark.sql.files.maxRecordsPerFile", 1000000)  # 每个文件最多100万条记录
    .option("mergeSchema", "true") \
    .saveAsTable("output")

如果一定要按id做某种分区,可以考虑用bucketBy,但要注意bucketBy是针对表的,且需要提前指定桶数,适合后续查询优化,而不是用来解决当前写入崩溃的问题。

优化后的整体代码示例

from pyspark.ml.feature import StringIndexer, OneHotEncoderEstimator, Pipeline
from pyspark.sql.functions import col, sum

indexer = StringIndexer(inputCol="kpID", outputCol="KPindex")
encoder = OneHotEncoderEstimator(inputCols=[indexer.getOutputCol()], outputCols=["KPvec"])
pipeline = Pipeline(stages=[indexer, encoder])

# 数据转换+聚合(直接操作稀疏向量)
df_bayes = pipeline.fit(df_bayes).transform(df_bayes)
num_partitions = 300
df_bayes = df_bayes.repartition(num_partitions, col('id'))
# 直接对稀疏向量求和,效率拉满
df_bayes = df_bayes.groupby('id').agg(sum(col('KPvec')).alias('KPvec')).cache()

# 写入Parquet
df_bayes.repartition(num_partitions).write.format("parquet") \
    .option("path", "s3://...output.parquet") \
    .option("spark.sql.files.maxRecordsPerFile", 1000000) \
    .saveAsTable("output")

预期优化效果

  • 聚合时间会从188s大幅降低(预计能降到几十秒内)
  • Shuffle数据量会从6GB降到几百MB级别
  • 写入Parquet时不会再出现ExecutorLostFailure错误

内容的提问来源于stack exchange,提问作者user1860389

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 19:32:43