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

Spark Structured Streaming写入Delta表时GCS连接重置报错排查

Spark Structured Streaming写入Delta表(GCS存储)随机失败问题

我有一个从Kafka Topic读取数据并写入Delta表的Spark Structured Streaming查询,该查询正常运行数月无异常,但近一周随机在运行15-20小时后失败,报错信息为java.io.IOException: Failed to get result: java.io.IOException: Error listing gs://bucket_name/table_name/_delta_log/。数据存储在Google Cloud Storage(GCS)Bucket中,无法确定该错误源于Delta表格式问题还是GCS读取问题。

完整错误日志

java.io.IOException: Failed to get result: java.io.IOException: Error listing gs://bucket_name/table_name/_delta_log/
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageFileSystem.getFromFuture(GoogleCloudStorageFileSystem.java:898)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageFileSystem.listFileInfo(GoogleCloudStorageFileSystem.java:1037)
    at com.google.cloud.hadoop.fs.gcs.GoogleHadoopFileSystemBase.listStatus(GoogleHadoopFileSystemBase.java:856)
    at org.apache.spark.sql.delta.storage.HadoopFileSystemLogStore.listFrom(HadoopFileSystemLogStore.scala:69)
    at org.apache.spark.sql.delta.storage.DelegatingLogStore.listFrom(DelegatingLogStore.scala:96)
    at org.apache.spark.sql.delta.SnapshotManagement.listFrom(SnapshotManagement.scala:64)
    at org.apache.spark.sql.delta.SnapshotManagement.listFrom$(SnapshotManagement.scala:63)
    at org.apache.spark.sql.delta.DeltaLog.listFrom(DeltaLog.scala:59)
    at org.apache.spark.sql.delta.SnapshotManagement.getLogSegmentForVersion(SnapshotManagement.scala:88)
    at org.apache.spark.sql.delta.SnapshotManagement.getLogSegmentForVersion$(SnapshotManagement.scala:82)
    at org.apache.spark.sql.delta.DeltaLog.getLogSegmentForVersion(DeltaLog.scala:59)
    at org.apache.spark.sql.delta.SnapshotManagement.$anonfun$updateInternal$1(SnapshotManagement.scala:289)
    at com.databricks.spark.util.DatabricksLogging.recordOperation(DatabricksLogging.scala:77)
    at com.databricks.spark.util.DatabricksLogging.recordOperation$(DatabricksLogging.scala:67)
    at org.apache.spark.sql.delta.DeltaLog.recordOperation(DeltaLog.scala:59)
    at org.apache.spark.sql.delta.metering.DeltaLogging.recordDeltaOperation(DeltaLogging.scala:106)
    at org.apache.spark.sql.delta.metering.DeltaLogging.recordDeltaOperation$(DeltaLogging.scala:91)
    at org.apache.spark.sql.delta.DeltaLog.recordDeltaOperation(DeltaLog.scala:59)
    at org.apache.spark.sql.delta.SnapshotManagement.updateInternal(SnapshotManagement.scala:287)
    at org.apache.spark.sql.delta.SnapshotManagement.updateInternal$(SnapshotManagement.scala:286)
    at org.apache.spark.sql.delta.DeltaLog.updateInternal(DeltaLog.scala:59)
    at org.apache.spark.sql.delta.SnapshotManagement.$anonfun$update$1(SnapshotManagement.scala:248)
    at org.apache.spark.sql.delta.DeltaLog.lockInterruptibly(DeltaLog.scala:152)
    at org.apache.spark.sql.delta.SnapshotManagement.update(SnapshotManagement.scala:248)
    at org.apache.spark.sql.delta.SnapshotManagement.update$(SnapshotManagement.scala:244)
    at org.apache.spark.sql.delta.DeltaLog.update(DeltaLog.scala:59)
    at org.apache.spark.sql.delta.DeltaLog.startTransaction(DeltaLog.scala:171)
    at org.apache.spark.sql.delta.DeltaLog.withNewTransaction(DeltaLog.scala:185)
    at org.apache.spark.sql.delta.sources.DeltaSink.addBatch(DeltaSink.scala:54)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$16(MicroBatchExecution.scala:586)
    at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:103)
    at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:163)
    at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:90)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775)
    at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runBatch$15(MicroBatchExecution.scala:584)
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:357)
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:355)
    at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runBatch(MicroBatchExecution.scala:584)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$2(MicroBatchExecution.scala:226)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken(ProgressReporter.scala:357)
    at org.apache.spark.sql.execution.streaming.ProgressReporter.reportTimeTaken$(ProgressReporter.scala:355)
    at org.apache.spark.sql.execution.streaming.StreamExecution.reportTimeTaken(StreamExecution.scala:68)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.$anonfun$runActivatedStream$1(MicroBatchExecution.scala:194)
    at org.apache.spark.sql.execution.streaming.ProcessingTimeExecutor.execute(TriggerExecutor.scala:57)
    at org.apache.spark.sql.execution.streaming.MicroBatchExecution.runActivatedStream(MicroBatchExecution.scala:188)
    at org.apache.spark.sql.execution.streaming.StreamExecution.$anonfun$runStream$1(StreamExecution.scala:334)
    at scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23)
    at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775)
    at org.apache.spark.sql.execution.streaming.StreamExecution.org$apache$spark$sql$execution$streaming$StreamExecution$$runStream(StreamExecution.scala:317)
    at org.apache.spark.sql.execution.streaming.StreamExecution$$anon$1.run(StreamExecution.scala:244)
Caused by: java.util.concurrent.ExecutionException: java.io.IOException: Error listing gs://bucket_name/table_name/_delta_log/
    at java.util.concurrent.FutureTask.report(FutureTask.java:122)
    at java.util.concurrent.FutureTask.get(FutureTask.java:192)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageFileSystem.getFromFuture(GoogleCloudStorageFileSystem.java:892)
    ... 52 more
Caused by: java.io.IOException: Error listing gs://bucket_name/table_name/_delta_log/
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl.listStorageObjectsAndPrefixesPage(GoogleCloudStorageImpl.java:1476)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl.listStorageObjectsAndPrefixes(GoogleCloudStorageImpl.java:1446)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl.listObjectInfo(GoogleCloudStorageImpl.java:1599)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageFileSystem.lambda$listFileInfo$7(GoogleCloudStorageFileSystem.java:1017)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)
Caused by: java.net.SocketException: Connection reset
    at java.net.SocketInputStream.read(SocketInputStream.java:210)
    at java.net.SocketInputStream.read(SocketInputStream.java:141)
    at org.conscrypt.ConscryptEngineSocket$SSLInputStream.readFromSocket(ConscryptEngineSocket.java:920)
    at org.conscrypt.ConscryptEngineSocket$SSLInputStream.processDataFromSocket(ConscryptEngineSocket.java:884)
    at org.conscrypt.ConscryptEngineSocket$SSLInputStream.readUntilDataAvailable(ConscryptEngineSocket.java:799)
    at org.conscrypt.ConscryptEngineSocket$SSLInputStream.read(ConscryptEngineSocket.java:772)
    at java.io.BufferedInputStream.fill(BufferedInputStream.java:246)
    at java.io.BufferedInputStream.read1(BufferedInputStream.java:286)
    at java.io.BufferedInputStream.read(BufferedInputStream.java:345)
    at sun.net.www.MeteredStream.read(MeteredStream.java:134)
    at java.io.FilterInputStream.read(FilterInputStream.java:133)
    at sun.net.www.protocol.http.HttpURLConnection$HttpInputStream.read(HttpURLConnection.java:3454)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.http.javanet.NetHttpResponse$SizeValidatingInputStream.read(NetHttpResponse.java:164)
    at com.google.cloud.hadoop.repackaged.gcs.com.fasterxml.jackson.core.json.UTF8StreamJsonParser._loadMore(UTF8StreamJsonParser.java:225)
    at com.google.cloud.hadoop.repackaged.gcs.com.fasterxml.jackson.core.json.UTF8StreamJsonParser.parseEscapedName(UTF8StreamJsonParser.java:2019)
    at com.google.cloud.hadoop.repackaged.gcs.com.fasterxml.jackson.core.json.UTF8StreamJsonParser.slowParseName(UTF8StreamJsonParser.java:1925)
    at com.google.cloud.hadoop.repackaged.gcs.com.fasterxml.jackson.core.json.UTF8StreamJsonParser._parseName(UTF8StreamJsonParser.java:1709)
    at com.google.cloud.hadoop.repackaged.gcs.com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:766)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.jackson2.JacksonParser.nextToken(JacksonParser.java:53)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parse(JsonParser.java:466)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parseValue(JsonParser.java:787)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parseArray(JsonParser.java:641)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parseValue(JsonParser.java:744)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parse(JsonParser.java:451)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parseValue(JsonParser.java:787)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parse(JsonParser.java:360)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonParser.parse(JsonParser.java:335)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonObjectParser.parseAndClose(JsonObjectParser.java:79)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.json.JsonObjectParser.parseAndClose(JsonObjectParser.java:73)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.http.HttpResponse.parseAs(HttpResponse.java:456)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.api.client.googleapis.services.AbstractGoogleClientRequest.execute(AbstractGoogleClientRequest.java:565)
    at com.google.cloud.hadoop.repackaged.gcs.com.google.cloud.hadoop.gcsio.GoogleCloudStorageImpl.listStorageObjectsAndPrefixesPage(GoogleCloudStorageImpl.java:1468)
    ... 7 more

简化后的代码

val df = spark
    .readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", kafkaUrl)
    .option("subscribe", topic)
    .option("startingOffsets", "earliest")
    .option("failOnDataLoss", true)
    .option("maxOffsetsPerTrigger", 3000000)
    .load()
    .selectExpr(s"""deserialize('$topic',value) AS value""")

df.writeStream
      .queryName("myquery")
      .format("delta")
      .partitionBy("date_time", "bin")
      .outputMode("append")
      .option("checkpointLocation", checkpoint)
      .option("dataChange", "false")
      .toTable(tableName)

Spark Submit配置

spark-submit  --driver-memory 6G --executor-memory 8G --executor-cores 4 --num-executors 3  --master yarn --deploy-mode cluster --conf spark.yarn.maxAppAttempts=2 --conf spark.yarn.am.attemptFailuresValidityInterval=1h

我曾尝试搜索类似问题但未找到有效解决方案。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 17:55:54