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

Dataproc上PySpark任务抛出IOException但仍成功完成问题排查

问题

在Google Cloud Dataproc上运行PySpark结构化流任务,trigger设为once,从原始层GCS Bucket读取Parquet数据,经业务规则处理后以Delta格式写入可信层GCS Bucket。环境配置:

  • Dataproc 2.1-debian11镜像
  • Spark 3.3.0
  • Delta 2.3.0

任务运行时抛出java.io.IOException错误,但未终止且最终成功完成,错误日志如下:

ERROR MicroBatchExecution: Query [id = a7859181-5d6e-4cb4-a4bc-93bbba0de6bf, runId = bd43ed9f-7300-4df5-b806-85badb86db5f] terminated with error
java.io.IOException: Failed to write 941418 bytes in 'gs://promo-bucket-data/bk-promo-co-grupoexito-datalake-dev/NAPSE/PROMOCIONES/PROMO_ENCABEZADO/_checkpoint/sources/0/.0.0eaeaaaa-014a-4d64-b7d8-1609c46c0c75.tmp'
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.util.BaseAbstractGoogleAsyncWriteChannel.write(BaseAbstractGoogleAsyncWriteChannel.java:136) ~[gcs-connector-hadoop3-2.2.16.jar:?]
    at java.nio.channels.Channels.writeFullyImpl(Channels.java:74) ~[?:?]
    at java.nio.channels.Channels.writeFully(Channels.java:97) ~[?:?]
    at java.nio.channels.Channels$1.write(Channels.java:172) ~[?:?]
    at java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:81) ~[?:?]
    at java.io.BufferedOutputStream.flush(BufferedOutputStream.java:142) ~[?:?]
    at java.io.FilterOutputStream.close(FilterOutputStream.java:182) ~[?:?]
    at com.google.cloud.hadoop.fs.gcs.GoogleHadoopOutputStream.lambda$close$2(GoogleHadoopOutputStream.java:144) ~[gcs-connector-hadoop3-2.2.16.jar:?]
    at com.google.cloud.hadoop.fs.gcs.GhfsStorageStatistics.trackDuration(GhfsStorageStatistics.java:77) ~[gcs-connector-hadoop3-2.2.16.jar:?]
    at com.google.cloud.hadoop.fs.gcs.GoogleHadoopOutputStream.close(GoogleHadoopOutputStream.java:136) ~[gcs-connector-hadoop3-2.2.16.jar:?]
    at org.apache.hadoop.fs.FSDataOutputStream$PositionCache.close(FSDataOutputStream.java:77) ~[hadoop-client-api-3.3.3.jar:?]
    at org.apache.hadoop.fs.FSDataOutputStream.close(FSDataOutputStream.java:106) ~[hadoop-client-api-3.3.3.jar:?]
    at org.apache.spark.sql.execution.streaming.CheckpointFileManager$RenameBasedFSDataOutputStream.close(CheckpointFileManager.scala:152) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.HDFSMetadataLog.$anonfun$addNewBatchByStream$2(HDFSMetadataLog.scala:176) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at scala.runtime.java8.JFunction0$mcZ$sp.apply(JFunction0$mcZ$sp.java:23) ~[scala-library-2.12.14.jar:?]
    at scala.Option.getOrElse(Option.scala:189) ~[scala-library-2.12.14.jar:?]
    at org.apache.spark.sql.execution.streaming.HDFSMetadataLog.addNewBatchByStream(HDFSMetadataLog.scala:171) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.HDFSMetadataLog.add(HDFSMetadataLog.scala:116) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.CompactibleFileStreamLog.add(CompactibleFileStreamLog.scala:168) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.FileStreamSourceLog.add(FileStreamSourceLog.scala:66) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.FileStreamSource.fetchMaxOffset(FileStreamSource.scala:198) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.FileStreamSource.latestOffset(FileStreamSource.scala:340) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$constructNextBatch$4(MicroBatchExecution.scala:448) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$constructNextBatch$2(MicroBatchExecution.scala:447) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:286) ~[scala-library-2.12.14.jar:?]
    at scala.collection.Iterator.foreach(Iterator.scala:943) ~[scala-library-2.12.14.jar:?]
    at scala.collection.Iterator.foreach$(Iterator.scala:943) ~[scala-library-2.12.14.jar:?]
    at scala.collection.AbstractIterator.foreach(Iterator.scala:1431) ~[scala-library-2.12.14.jar:?]
    at scala.collection.IterableLike.foreach(IterableLike.scala:74) ~[scala-library-2.12.14.jar:?]
    at scala.collection.IterableLike.foreach$(IterableLike.scala:73) ~[scala-library-2.12.14.jar:?]
    at scala.collection.AbstractIterable.foreach(Iterable.scala:56) ~[scala-library-2.12.14.jar:?]
    at scala.collection.TraversableLike.map(TraversableLike.scala:286) ~[scala-library-2.12.14.jar:?]
    at scala.collection.TraversableLike.map$(TraversableLike.scala:279) ~[scala-library-2.12.14.jar:?]
    at scala.collection.AbstractTraversable.map(Traversable.scala:108) ~[scala-library-2.12.14.jar:?]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$constructNextBatch$1(MicroBatchExecution.scala:436) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at scala.runtime.java8.JFunction0$mcZ$sp.apply(JFunction0$mcZ$sp.java:23) ~[scala-library-2.12.14.jar:?]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.withProgressLocked(MicroBatchExecution.scala:687) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.constructNextBatch(MicroBatchExecution.scala:432) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:237) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) ~[scala-library-2.12.14.jar:?]
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:375) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:373) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$1(MicroBatchExecution.scala:218) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.SingleBatchExecutor.execute(TriggerExecutor.scala:39) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:212) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:307) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) ~[scala-library-2.12.14.jar:?]
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:779) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:285) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
    at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:208) ~[spark-sql_2.12-3.3.0.jar:3.3.0]
Caused by: java.nio.channels.ClosedByInterruptException
    at java.nio.channels.spi.AbstractInterruptibleChannel.end(AbstractInterruptibleChannel.java:199) ~[?:?]
    at java.nio.channels.Channels$WritableByteChannelImpl.write(Channels.java:466) ~[?:?]
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.util.BaseAbstractGoogleAsyncWriteChannel.write(BaseAbstractGoogleAsyncWriteChannel.java:133) ~[gcs-connector-hadoop3-2.2.16.jar:?]
    ... 53 more

原因分析

从错误栈的Caused by: java.nio.channels.ClosedByInterruptException可以定位核心问题:写入GCS checkpoint临时文件的线程被提前中断。具体诱因包括:

  1. Spark在trigger=once模式下,任务核心数据处理完成后会主动中断后台线程,但此时GCS连接器的异步写入逻辑还未完成checkpoint临时文件的写入,导致抛出异常。
  2. 当前使用的GCS连接器(2.2.16版本)的异步写入逻辑与Spark的线程中断机制存在时序冲突,临时文件写入未完成就被终止,但任务核心流程已经完成,因此最终任务仍能成功结束。

解决办法

1. 升级GCS连接器版本

旧版本GCS连接器存在线程中断处理的缺陷,升级到3.2.0及以上版本(适配Hadoop 3.x),新版本优化了异步写入的中断处理逻辑,可避免此类异常。
在Dataproc集群初始化时指定连接器版本:

gcloud dataproc clusters create <cluster-name> \
    --image-version=2.1-debian11 \
    --properties spark:spark.hadoop.fs.gs.impl=com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystem,spark:spark.hadoop.fs.AbstractFileSystem.gs.impl=com.google.cloud.hadoop.fs.gcs.GoogleHadoopFS \
    --metadata gcs-connector-version=3.2.0

2. 调整Spark checkpoint相关配置

增加GCS写入的缓冲和同步参数,给临时文件写入足够的完成时间:

spark.conf.set("spark.hadoop.fs.gs.outputstream.sync.interval", "1000")
spark.conf.set("spark.hadoop.fs.gs.outputstream.buffer.size", "4194304")

3. 优化trigger=once任务的等待逻辑

显式等待查询完成,确保所有后台写入操作结束后再结束任务:

query = df.writeStream \
    .format("delta") \
    .option("checkpointLocation", "gs://your-checkpoint-path") \
    .trigger(triggerOnce=True) \
    .start("gs://your-target-path")

# 等待查询完全终止,避免线程提前中断
query.awaitTermination()

4. 屏蔽非致命异常日志

如果上述方法无法完全消除异常,且任务最终能成功完成,可调整日志级别屏蔽该类非致命错误输出:

# 全局降低日志级别至WARN
spark.sparkContext.setLogLevel("WARN")

或在log4j配置中定向屏蔽:

log4j.logger.org.apache.spark.sql.execution.streaming.MicroBatchExecution=WARN

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 23:49:51