【发布时间】:2020-05-28 19:08:55
【问题描述】:
我在我的项目中使用 Spark-SQL-2.3.1、Kafka、Java 8,并希望将 AWS-S3 用作野蛮存储。
我正在将来自 Kafka 主题的消费数据写入/存储到 S3 存储桶中,如下所示:
ds.writeStream()
.format("parquet")
.option("path", parquetFileName)
.option("mergeSchema", true)
.outputMode("append")
.partitionBy("company_id")
.option("checkpointLocation", checkPtLocation)
.trigger(Trigger.ProcessingTime("25 seconds"))
.start();
但在写作时我收到了FileNotFoundException
Caused by: java.io.FileNotFoundException: No such file or directory: s3a://company_id=216231245/part-00055-f4f87dc9-a620-41bd-9380-de4ba7e70efb.c000.snappy.parquet
at org.apache.hadoop.fs.s3a.S3AFileSystem.s3GetFileStatus(S3AFileSystem.java:1931)
at org.apache.hadoop.fs.s3a.S3AFileSystem.innerGetFileStatus(S3AFileSystem.java:1822)
at org.apache.hadoop.fs.s3a.S3AFileSystem.getFileStatus(S3AFileSystem.java:1763)
我想知道为什么我在写作时收到FileNotFoundException?我不是在读 S3 对吗?
那么这里发生了什么以及如何解决这个问题?
【问题讨论】:
-
发生了什么是检查点和分段上传。 IDK 为什么你会得到它,但是你使用的是 lib 上的每个版本?
-
这不是完整的堆栈跟踪,是吗?另外,如果你不尝试合并 Parquet 模式会发生什么(模式合并引发读取)?
-
这其实是404缓存造成的意外;在创建之前对文件发出的 HEAD 请求可以缓存在负载均衡器中,然后很快就会中断重命名。尚未发布的 Hadoop 版本(稍微)更好,但要安全地提交工作,您必须使用 jegan 讨论的那种 s3a 提交器
标签: apache-spark amazon-s3 apache-spark-sql spark-structured-streaming