【问题标题】:Looking for a way to continuously process files written to hdfs寻找一种方法来持续处理写入 hdfs 的文件
【发布时间】:2017-04-16 05:34:57
【问题描述】:

我正在寻找一种可以:

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

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

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

这是一个应该很常见的用例。我是否在 spark/flink 中遗漏了一些使它成为可能的东西?还有其他工具可以在这里使用吗?

【问题讨论】:

  • 你考虑过 Apache NiFi 吗?啊,也许你更喜欢从头开始手工编码......

标签: hadoop apache-spark hdfs apache-flink bigdata


【解决方案1】:

使用 Flink 流,您可以完全按照您的建议处理目录中的文件,并且当您重新启动时,它将从中断处开始处理。它被称为连续文件处理。

您唯一需要做的就是 1) 为您的工作启用检查点和 2) 启动您的程序:

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

该功能相当新,开发者邮件列表中正在积极讨论如何进一步改进其功能。

希望这会有所帮助!

【讨论】:

  • 顺便说一句,您必须启用检查点才能从您离开的地方开始。如果不是,那么一切都将被重新处理。
  • 这也不起作用.. 仅当某些文件处于处理过程中时才会存储检查点信息。完成后 - 没有关于已处理内容的信息。检查点目录为空,因此在作业重新启动后,所有内容都会再次处理。
【解决方案2】:

我建议您稍微修改文件摄取并合并 Kafka,这样每次您在 HDFS 中放入新文件时,都会在 Kafka 队列中放入一条消息。然后使用 Spark 流从队列中读取文件名,然后从 hdfs 和进程中读取文件。

检查点是一个真正的痛苦,也不能保证你想要什么。带有 spark 的 Kafka 将能够保证只有一次语义。

Flume 有一个 SpoolDirSource ,你也可以看看。

【讨论】:

    【解决方案3】:

    最好的方法是维护一个状态机。维护一个包含所有已处理文件的表或文件。

    应用程序在启动时读取文件列表并在 set/map 中保持相同。任何已处理的新文件/旧文件都可以查看和验证。

    摄取文件夹也需要维护一些文件状态。已处理的类似文件用一些 ext 重命名。失败的文件被移动到失败的文件夹,被拒绝到被拒绝的文件夹。等等

    您可以使用 spark/flink 完成所有这些工作。技术不是这里的瓶颈

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-12-18
      • 1970-01-01
      • 1970-01-01
      • 2011-10-23
      • 2011-10-02
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多