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")
问题分析与修复方案
核心原因
- 依赖版本冲突:集群自带Spark 3.3.0(Scala 2.12),但初始化配置中引入了Spark 3.4.0(Scala 2.13)和Scala 2.13.1,跨版本依赖导致Python与JVM层通信链路异常,触发Socket连接重置。
- MLflow自动日志的隐式Accumulator:
mlflow.pyspark.ml.autolog()会自动创建Accumulator收集训练指标,任务完成时Driver可能提前关闭与Executor的连接,导致Accumulator更新请求失败。 - Executor资源提前回收:Dataproc动态分配机制可能在任务接近完成时回收空闲Executor,打断Accumulator的更新流程。
修复步骤
- 统一依赖版本:
删除与集群版本冲突的依赖项,修改spark.jars.packages配置:.config("spark.jars.packages", "org.mlflow:mlflow-spark:1.30.0,io.findify:s3mock_2.12:0.1.8") - 调整MLflow自动日志策略:
禁用不必要的自动日志功能,减少隐式Accumulator的创建:mlflow.pyspark.ml.autolog(log_models=False) # 根据实际需求调整参数 - 延长Executor空闲保留时间:
添加Spark配置,避免Executor被提前回收:.config("spark.dynamicAllocation.executorIdleTimeout", "300s") .config("spark.dynamicAllocation.shuffleTracking.enabled", "true") - 手动控制任务收尾流程:
在mlflow.end_run()前确保所有Spark作业执行完成,可添加spark.sparkContext.waitForJobs()等待所有任务结束后再关闭MLflow运行。
内容的提问来源于stack exchange,提问作者Rohan Aswani
相关产品推荐
相关产品推荐

