【问题标题】:Apache Flink: read file from HDFSApache Flink:从 HDFS 读取文件
【发布时间】:2018-05-18 20:06:12
【问题描述】:

所以我必须检索存储在 HDFS 中的文件的内容并对其进行某些分析。

问题是,我什至无法读取文件并将其内容写入本地文件系统中的另一个文本文件。 (我是 Flink 新手,这只是一个测试,以确保我正确读取文件)

HDFS 中的文件是纯文本文件。这是我的代码:

public class readFromHdfs {

    public static void main(String[] args) throws Exception {

        final ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();

        DataSet<String> lines = env.readTextFile("hdfs://localhost:9000//test/testfile0.txt");

        lines.writeAsText("/tmp/hdfs_file.txt"); 

        env.execute("File read from HDFS");
    }
}

/tmp 运行后没有输出。

这是一个非常简单的代码,我不确定它是否有问题,或者我只是做错了其他事情。正如我所说,我对 Flink 完全陌生

此外,该作业在 Web 仪表板中显示为失败。 flink 日志内容如下:https://pastebin.com/rvkXPGHU

提前致谢

编辑:我通过增加任务槽的数量解决了这个问题。网络仪表板显示了一个可用的任务槽,它根本没有抱怨没有足够的槽,所以我认为不可能是这样。

无论如何,writeAsText 并没有像我预期的那样工作。我从 testfile0.txt 读取内容没问题,但它没有将它们写入 hdfs_file.txt。相反,它使用该名称创建一个目录,其中包含 8 个文本文件,其中 6 个完全为空。另外两个包含testfile0.txt(大部分在1.txt,最后一个chunk在2.txt)。

虽然这并不重要,因为文件的内容已正确存储在 DataSet 中,所以我可以继续分析数据。

【问题讨论】:

    标签: java hdfs apache-flink


    【解决方案1】:

    它按预期工作 - 您已将完整作业的并行度(以及输出格式)设置为 8,因此每个插槽都会创建自己的文件(您可能知道并发写入单个文件是不安全的)。如果您只需要 1 个输出文件,您应该 writeAsText(...).setParalellis(1) 覆盖全局并行属性。

    如果你想在本地文件系统而不是 HDFS 中获取输出,你应该在路径中显式设置“file://”协议,因为对于 Hadoop,flink 默认看起来是“hdfs://”。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-06-14
      • 2017-05-06
      • 1970-01-01
      • 1970-01-01
      • 2021-06-18
      • 2015-05-11
      相关资源
      最近更新 更多