【问题标题】:How do I delete files in hdfs directory after reading it using scala?使用scala读取hdfs目录后如何删除文件?
【发布时间】:2017-12-19 15:10:56
【问题描述】:

我使用 fileStream 从 Spark(流上下文)中读取 hdfs 目录中的文件。如果我的 Spark 在一段时间后关闭并启动,我想读取目录中的新文件。我不想读取目录中已被 Spark 读取和处理的旧文件。我在这里尽量避免重复。

val lines = ssc.fileStream[LongWritable, Text, TextInputFormat]("/home/File")

任何代码 sn-ps 可以帮助?

【问题讨论】:

    标签: scala hadoop apache-spark spark-streaming


    【解决方案1】:

    您可以使用FileSystem API:

    import org.apache.hadoop.fs.{FileSystem, Path}
    
    val fs = FileSystem.get(sc.hadoopConfiguration)
    
    val outPutPath = new Path("/abc")
    
    if (fs.exists(outPutPath))
      fs.delete(outPutPath, true)
    

    【讨论】:

      【解决方案2】:

      fileStream 已经为您处理了 - 来自它的 Scaladoc

      创建一个输入流,用于监控与 Hadoop 兼容的文件系统中的新文件,并使用给定的键值类型和输入格式读取它们。

      这意味着fileStream 只会加载新文件(在流式上下文启动后创建),在您启动流式应用程序之前文件夹中已经存在的任何文件都将被忽略。

      【讨论】:

      • 我想在这里解释一下容错。假设我在 hdfs 中有 1 到 10 个文件,并且 spark 流已经读取了这个文件。现在我的火花流不幸停止了。我在 hdfs 中有文件说 1 到 20,其中 1 到 10 个文件已经被 spark 流解析,并且新添加了 11 到 20 个文件。现在我开始 Spark Streaming,我可以看到文件 1-30。由于我在 hdfs 中的第 21 个文件时启动了 spark,所以我的 spark styreaming 会丢失文件 11-20。如何丢失文件。
      猜你喜欢
      • 1970-01-01
      • 2021-06-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多