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

Dataproc集群PySpark任务Accumulator更新失败问题求助

Dataproc集群PySpark任务Accumulator更新异常问题

现象

在Dataproc集群运行PySpark训练任务时,单任务日志显示执行完成,但日志中出现PythonAccumulatorV2更新失败的Socket异常。本地执行完全正常,集群任务标记为成功,但异常日志持续出现。

任务成功执行日志

23/06/11 15:03:35 INFO LoggerProvider: deleted one hot encode
23/06/11 15:03:35 INFO LoggerProvider: Started Assemble Vector
23/06/11 15:03:36 INFO LoggerProvider: Completed Assemble Vector
23/06/11 15:03:36 INFO LoggerProvider: started training pipeline
23/06/11 15:03:36 INFO LoggerProvider: Started mlflow active run

Accumulator更新异常日志

23/06/11 14:02:53 ERROR DAGScheduler: Failed to update accumulator 0 (org.apache.spark.api.python.PythonAccumulatorV2) for task 54
java.net.SocketException: Connection reset
    at java.net.SocketInputStream.read(SocketInputStream.java:186) ~[?:?]
    at java.net.SocketInputStream.read(SocketInputStream.java:140) ~[?:?]
    at java.net.SocketInputStream.read(SocketInputStream.java:200) ~[?:?]
    at org.apache.spark.api.python.PythonAccumulatorV2.merge(PythonRDD.scala:734) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$updateAccumulators$1(DAGScheduler.scala:1610) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$updateAccumulators$1$adapted(DAGScheduler.scala:1601) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) ~[scala-library-2.12.14.jar:?]
    at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) ~[scala-library-2.12.14.jar:?]
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) ~[scala-library-2.12.14.jar:?]
    at org.apache.spark.scheduler.DAGScheduler.updateAccumulators(DAGScheduler.scala:1601) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGScheduler.handleTaskCompletion(DAGScheduler.scala:1749) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2858) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2803) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2792) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49) ~[spark-core_2.12-3.3.0.jar:3.3.0]
23/06/11 14:02:53 ERROR DAGScheduler: Failed to update accumulator 0 (org.apache.spark.api.python.PythonAccumulatorV2) for task 57
java.net.SocketException: Broken pipe (Write failed)
    at java.net.SocketOutputStream.socketWrite0(Native Method) ~[?:?]
    at java.net.SocketOutputStream.socketWrite(SocketOutputStream.java:110) ~[?:?]
    at java.net.SocketOutputStream.write(SocketOutputStream.java:150) ~[?:?]
    at java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:81) ~[?:?]
    at java.io.BufferedOutputStream.flush(BufferedOutputStream.java:142) ~[?:?]
    at java.io.DataOutputStream.flush(DataOutputStream.java:123) ~[?:?]
    at org.apache.spark.api.python.PythonAccumulatorV2.merge(PythonRDD.scala:732) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$updateAccumulators$1(DAGScheduler.scala:1610) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGScheduler.$anonfun$updateAccumulators$1$adapted(DAGScheduler.scala:1601) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62) ~[scala-library-2.12.14.jar:?]
    at scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55) ~[scala-library-2.12.14.jar:?]
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49) ~[scala-library-2.12.14.jar:?]
    at org.apache.spark.scheduler.DAGScheduler.updateAccumulators(DAGScheduler.scala:1601) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGScheduler.handleTaskCompletion(DAGScheduler.scala:1749) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2858) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2803) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2792) ~[spark-core_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49) ~[spark-core_2.12-3.3.0.jar:3.3.0]

Spark初始化代码

spark = SparkSession.builder.master('yarn')\
    .appName("pyspark-test")\
    .config("spark.jars.packages", "org.mlflow:mlflow-spark:1.30.0,io.findify:s3mock_2.12:0.1.8,org.apache.spark:spark-sql_2.13:3.4.0,org.scala-lang:scala-library:2.13.1")\
    .config("spark.sql.broadcastTimeout", "36000")\
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer")\
    .getOrCreate()
spark.conf.set('temporaryGcsBucket', bucket)
spark.conf.set("spark.sql.legacy.timeParserPolicy","LEGACY")

业务训练代码

logger.info("started training pipeline")
mlflow.pyspark.ml.autolog()
train, test = df.randomSplit([0.7, 0.3], seed = 2018)
with mlflow.start_run() as active_run:
    logger.info("Started mlflow active run")

    rfc = RandomForestClassifier(featuresCol="features", labelCol="label")
    rfc = rfc.fit(train)
    signature = infer_signature( df.select("features"), df.select("label"))
    mlflow.spark.log_model(
            spark_model=rfc, 
            artifact_path=f"rfc-model",
            registered_model_name=f"rfc-model",
            signature=signature
        )
    evaluator = MulticlassClassificationEvaluator(
    labelCol="product_id_in", predictionCol="prediction", metricName="accuracy")
    accuracy = evaluator.evaluate(rfc.transform(test))
    mlflow.log_metric("recall", accuracy)
    logger.info("Completed active run")
mlflow.end_run()
logger.info("completed training pipeline")

问题分析与修复方案

核心原因

  1. 依赖版本冲突:集群自带Spark 3.3.0(Scala 2.12),但初始化配置中引入了Spark 3.4.0(Scala 2.13)和Scala 2.13.1,跨版本依赖导致Python与JVM层通信链路异常,触发Socket连接重置。
  2. MLflow自动日志的隐式Accumulator:mlflow.pyspark.ml.autolog()会自动创建Accumulator收集训练指标,任务完成时Driver可能提前关闭与Executor的连接,导致Accumulator更新请求失败。
  3. Executor资源提前回收:Dataproc动态分配机制可能在任务接近完成时回收空闲Executor,打断Accumulator的更新流程。

修复步骤

  1. 统一依赖版本:
    删除与集群版本冲突的依赖项,修改spark.jars.packages配置:
    .config("spark.jars.packages", "org.mlflow:mlflow-spark:1.30.0,io.findify:s3mock_2.12:0.1.8")
    
  2. 调整MLflow自动日志策略:
    禁用不必要的自动日志功能,减少隐式Accumulator的创建:
    mlflow.pyspark.ml.autolog(log_models=False)  # 根据实际需求调整参数
    
  3. 延长Executor空闲保留时间:
    添加Spark配置,避免Executor被提前回收:
    .config("spark.dynamicAllocation.executorIdleTimeout", "300s")
    .config("spark.dynamicAllocation.shuffleTracking.enabled", "true")
    
  4. 手动控制任务收尾流程:
    在mlflow.end_run()前确保所有Spark作业执行完成,可添加spark.sparkContext.waitForJobs()等待所有任务结束后再关闭MLflow运行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 13:07:00