寻找一种方法来连续处理写入hdfs的文件

Art*_*Art 1 hadoop bigdata hdfs apache-spark apache-flink

我正在寻找一种可以:

  1. 监视新文件的hdfs dir并在它们出现时处理它们.
  2. 它还应该在作业/应用程序开始工作之前处理目录中的文件.
  3. 它应该有检查点,以便在重新启动时从它离开的地方继续.

我看了一下apache spark:它可以读取新添加的文件,并且可以处理重启以从它离开的地方继续.我找不到办法使它也处理同一作业范围内的旧文件(所以只有1和3).

我看了一下apache flink:它确实处理了新旧文件.但是,一旦重新启动作业,它将再次开始处理所有作业(1和2).

这是一个非常常见的用例.我错过了火花/叮当的东西,这使得它成为可能吗?还有其他工具可以在这里使用吗?

小智 5

使用Flink流式传输,您可以完全按照建议处理目录中的文件,当您重新启动时,它将从停止的位置开始处理.它被称为连续文件处理.

您唯一需要做的就是1)为您的工作启用检查点,2)启动您的计划:

    Time period = Time.minutes(10)
    env.readFile(inputFormat, "hdfs:// … /logs",
                 PROCESS_CONTINUOUSLY, 
                 period.toMilliseconds, 
                 FilePathFilter.createDefaultFilter())
Run Code Online (Sandbox Code Playgroud)

该功能相当新,在dev邮件列表中有一个关于如何进一步改进其功能的积极讨论.

希望这可以帮助!