【问题标题】:Spark Streaming textFileStream not supporting wildcardsSpark Streaming textFileStream 不支持通配符
【发布时间】: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


【解决方案1】:

通过扩展 FileInputDStream 可以创建一些“丑陋但有效的解决方案”。 写sc.textFileStream(d)相当于

new FileInputDStream[LongWritable, Text, TextInputFormat](streamingContext, d).map(_._2.toString)

您可以创建将扩展 FileInputDStream 的 CustomFileInputDStream。自定义类将从 FileInputDStream 类复制计算方法,并根据需要调整 findNewFiles 方法。

从以下位置更改 findNewFiles 方法:

 private def findNewFiles(currentTime: Long): Array[String] = {
    try {
      lastNewFileFindingTime = clock.getTimeMillis()

  // Calculate ignore threshold
  val modTimeIgnoreThreshold = math.max(
    initialModTimeIgnoreThreshold,   // initial threshold based on newFilesOnly setting
    currentTime - durationToRemember.milliseconds  // trailing end of the remember window
  )
  logDebug(s"Getting new files for time $currentTime, " +
    s"ignoring files older than $modTimeIgnoreThreshold")
  val filter = new PathFilter {
    def accept(path: Path): Boolean = isNewFile(path, currentTime, modTimeIgnoreThreshold)
  }
  val newFiles = fs.listStatus(directoryPath, filter).map(_.getPath.toString)
  val timeTaken = clock.getTimeMillis() - lastNewFileFindingTime
  logInfo("Finding new files took " + timeTaken + " ms")
  logDebug("# cached file times = " + fileToModTime.size)
  if (timeTaken > slideDuration.milliseconds) {
    logWarning(
      "Time taken to find new files exceeds the batch size. " +
        "Consider increasing the batch size or reducing the number of " +
        "files in the monitored directory."
    )
  }
  newFiles
} catch {
  case e: Exception =>
    logWarning("Error finding new files", e)
    reset()
    Array.empty
}

}

到:

  private def findNewFiles(currentTime: Long): Array[String] = {
    try {
      lastNewFileFindingTime = clock.getTimeMillis()

      // Calculate ignore threshold
      val modTimeIgnoreThreshold = math.max(
        initialModTimeIgnoreThreshold,   // initial threshold based on newFilesOnly setting
        currentTime - durationToRemember.milliseconds  // trailing end of the remember window
      )
      logDebug(s"Getting new files for time $currentTime, " +
        s"ignoring files older than $modTimeIgnoreThreshold")
      val filter = new PathFilter {
        def accept(path: Path): Boolean = isNewFile(path, currentTime, modTimeIgnoreThreshold)
      }
      val directories = fs.listStatus(directoryPath).filter(_.isDirectory)
      val newFiles = ArrayBuffer[FileStatus]()

      directories.foreach(directory => newFiles.append(fs.listStatus(directory.getPath, filter) : _*))

      val timeTaken = clock.getTimeMillis() - lastNewFileFindingTime
      logInfo("Finding new files took " + timeTaken + " ms")
      logDebug("# cached file times = " + fileToModTime.size)
      if (timeTaken > slideDuration.milliseconds) {
        logWarning(
          "Time taken to find new files exceeds the batch size. " +
            "Consider increasing the batch size or reducing the number of " +
            "files in the monitored directory."
        )
      }
      newFiles.map(_.getPath.toString).toArray
    } catch {
      case e: Exception =>
        logWarning("Error finding new files", e)
        reset()
        Array.empty
    }
  }

将检查所有一级子文件夹中的文件,您可以将其调整为使用批处理时间戳来访问相关的“子目录”。

如前所述,我创建了 CustomFileInputDStream 并通过调用激活它:

new CustomFileInputDStream[LongWritable, Text, TextInputFormat](streamingContext, d).map(_._2.toString)

这似乎符合我们的预期。

当我写出这样的解决方案时,我必须添加一些考虑因素:

  • 您打破了 Spark 封装并创建了一个自定义类,随着时间的推移您将不得不单独支持该类。

  • 我相信这样的解决方案是最后的手段。如果您的用例可以通过不同的方式实现,通常最好避免这样的解决方案。

  • 如果您将在 S3 上有很多“子目录”并且要检查每个子目录,那么您将付出代价。

  • 如果 Databricks 不支持嵌套文件只是因为可能的性能损失,这将是非常有趣的,也许还有更深层次的原因我没有想到。

【讨论】:

  • 我有一个类似的用例,如果找不到替代方案,我正在考虑走这条路。我有使用 YYYY-MM-DD-HH 格式的日期分区子文件夹。每小时都会创建一个新文件夹并将文件上传到其中。所以我不一定要扫描所有子文件夹(只有最后三个),也不会遇到性能问题。我更担心此类代码的可维护性和重新启动的状态管理(上次扫描​​哪个小时文件夹+文件等)。看看您是否可以分享您对这个甚至适用于您的自定义 FileDstream 的代码的想法。
  • 如果您在流中使用检查点目录,那么当您重新启动应用程序时,您将首先重新安排应在应用程序停机期间执行的所有批处理。例如,如果您的流式传输间隔为 1 分钟,并且您的应用程序在 10:00 关闭并在 10:30 备份,那么当它启动时,应用程序将尝试执行 10:01、10:02 等批处理。现在,如果您以您扫描的文件夹派生自 currentTime 的方式实现 findNewFiles(currentTime),您将能够在重新启动后扫描“正确”的文件。
  • 请注意,currentTime实际上并不是CURRENT时间,而是批处理的时间。在这种方法中我能想到的唯一问题是,如果您的文件不是不可变的。例如,您在 10:10 将一些数据写入文件 A 并在 10:20 覆盖此数据,那么如果您的应用程序在 10:10-10:20 之间关闭,您将失去对 A 的第一次写入。这确实是一个问题,但我不熟悉在这种情况下使用可变文件的许多组织。
  • 谢谢你们!我们最终使用 AWS lambda 作为在 S3 中处理文件的第一步。
【解决方案2】:

我们有同样的问题。我们用逗号连接了子文件夹名称。

List<String> paths = new ArrayList<>();
SimpleDateFormat sdf = new SimpleDateFormat("yyyy/MM/dd");

try {        
    Date start = sdf.parse("2015/02/01");
    Date end = sdf.parse("2015/04/01");

    Calendar calendar = Calendar.getInstance();
    calendar.setTime(start);        

    while (calendar.getTime().before(end)) {
        paths.add("s3n://mybucket/" + sdf.format(calendar.getTime()));
        calendar.add(Calendar.DATE, 1);
    }                
} catch (ParseException e) {
    e.printStackTrace();
}

String joinedPaths = StringUtils.join(",", paths.toArray(new String[paths.size()]));
val input = ssc.textFileStream(joinedPaths);

希望通过这种方式解决你的问题。

【讨论】:

  • 酷。你如何处理更大的结束日期?通过编译和重新启动程序?还是我错过了什么?
猜你喜欢
  • 2016-10-08
  • 1970-01-01
  • 1970-01-01
  • 2015-11-29
  • 1970-01-01
  • 1970-01-01
  • 2018-01-13
  • 2018-04-21
  • 2011-01-26
相关资源
最近更新 更多