【发布时间】:2015-06-08 04:37:42
【问题描述】:
我设置了一个简单的测试来从 S3 流式传输文本文件,并在我尝试类似的东西时让它工作
val input = ssc.textFileStream("s3n://mybucket/2015/04/03/")
在存储桶中,我会将日志文件放入其中,一切都会正常工作。
但如果它们是子文件夹,它不会找到放入子文件夹的任何文件(是的,我知道 hdfs 实际上并不使用文件夹结构)
val input = ssc.textFileStream("s3n://mybucket/2015/04/")
所以,我尝试像以前使用标准 spark 应用程序那样简单地使用通配符
val input = ssc.textFileStream("s3n://mybucket/2015/04/*")
但是当我尝试这个时它会抛出一个错误
java.io.FileNotFoundException: File s3n://mybucket/2015/04/* does not exist.
at org.apache.hadoop.fs.s3native.NativeS3FileSystem.listStatus(NativeS3FileSystem.java:506)
at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1483)
at org.apache.hadoop.fs.FileSystem.listStatus(FileSystem.java:1523)
at org.apache.spark.streaming.dstream.FileInputDStream.findNewFiles(FileInputDStream.scala:176)
at org.apache.spark.streaming.dstream.FileInputDStream.compute(FileInputDStream.scala:134)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:300)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:300)
at scala.util.DynamicVariable.withValue(DynamicVariable.scala:57)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:299)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:287)
at scala.Option.orElse(Option.scala:257)
.....
我知道您可以在为标准 Spark 应用程序读取 fileInput 时使用通配符,但似乎在进行流式输入时,它不会这样做,也不会自动处理子文件夹中的文件。我这里有什么遗漏吗??
我最终需要的是一个 24/7 全天候运行的流式作业,它将监控按日期放置日志的 S3 存储桶
类似
s3n://mybucket/<YEAR>/<MONTH>/<DAY>/<LogfileName>
有没有办法把它放在最顶层的文件夹中,它会自动读取出现在任何文件夹中的文件(因为显然日期每天都会增加)?
编辑
因此,在深入研究http://spark.apache.org/docs/latest/streaming-programming-guide.html#basic-sources 的文档时,它指出不支持嵌套目录。
谁能解释一下为什么会这样?
另外,由于我的文件将根据它们的日期嵌套,在我的流应用程序中解决这个问题的好方法是什么?这有点复杂,因为日志需要几分钟才能写入 S3,因此当天写入的最后一个文件可能会写入前一天的文件夹中,即使我们距离新的一天还有几分钟。
【问题讨论】:
-
其实我不确定s3是否支持通配符...
-
确实如此。在过去的 8 个月里,我的工作一直在使用通配符。此外,只是为了进行完整性检查,我刚刚使用通配符输入运行了一项工作,工作正常。我确实注意到,要求您不要执行类似 s3n://mybucket/2015/04* 之类的操作,这有点挑剔,因为线程“main”java.io.IOException 中的异常:不是文件:s3n:/ /mybucket/2015/04/01 这是有道理的,因为它不是一个文件但是如果你这样做 s3n://mybucket/2015/04/* 它会正确解析 days 子文件夹中的所有文件.... 这种of 对我来说就像一个错误。
-
我要投票赞成这个问题。我记得有一个类似的问题,但我不记得我是如何解决的。
-
我很感激。这听起来确实是一种常见的实现方式。
标签: apache-spark hdfs spark-streaming