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
相关产品推荐
相关产品推荐

