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临时文件的线程被提前中断。具体诱因包括:
- Spark在
trigger=once模式下,任务核心数据处理完成后会主动中断后台线程,但此时GCS连接器的异步写入逻辑还未完成checkpoint临时文件的写入,导致抛出异常。 - 当前使用的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
相关产品推荐
相关产品推荐

