【问题标题】:Check if file exists on remote HDFS from local spark-submit从本地 spark-submit 检查远程 HDFS 上是否存在文件
【发布时间】:2020-07-27 01:30:32
【问题描述】:

我正在开发一个 Java 程序,专门用于在 HDFS 文件系统(位于 HDFS_IP)上使用 Spark。 我的目标之一是检查路径hdfs://HDFS_IP:HDFS_PORT/path/to/file.json 的HDFS 上是否存在文件。在本地调试我的程序时,我发现我无法使用以下代码访问这个远程文件

private boolean existsOnHDFS(String path) {
     Configuration conf = new Configuration();
     FileSystem fs;
     Boolean fileDoesExist = false ;
     try {
         fs = FileSystem.get(conf);
         fileDoesExist = fs.exists(new Path(path)) ;
     } catch (IOException e) {
            e.printStackTrace();
     }
     return fileDoesExist ;
 }

实际上,fs.exists 试图在我的本地 FS 中而不是 HDFS 中查找文件 hdfs://HDFS_IP:HDFS_PORT/path/to/file.json。顺便说一句,让hdfs://HDFS_IP:HDFS_PORT 前缀使fs.existscrash 并抑制它回答false,因为/path/to/file.json 在本地不存在。

fs 的适当配置是什么才能在本地和从 Hadoop 集群执行 Java 程序时正常工作?

编辑:我终于放弃并将错误修复传递给我团队中的其他人。感谢那些试图帮助我的人!

【问题讨论】:

    标签: java apache-spark hadoop hdfs


    【解决方案1】:

    问题是您向 FileSystem 传递了一个空的 conf 文件。

    你应该这样创建你的文件系统:

    FileSystem.get(spark.sparkContext().hadoopConfiguration());
    

    当 spark 是 SparkSession 对象时。

    正如你在文件系统的代码中看到的:

     /**
       * Returns the configured filesystem implementation.
       * @param conf the configuration to use
       */
      public static FileSystem get(Configuration conf) throws IOException {
        return get(getDefaultUri(conf), conf);
      }
    
      /** Get the default filesystem URI from a configuration.
       * @param conf the configuration to use
       * @return the uri of the default filesystem
       */
      public static URI getDefaultUri(Configuration conf) {
        return URI.create(fixName(conf.get(FS_DEFAULT_NAME_KEY, DEFAULT_FS)));
      }
    

    它根据作为参数传递的配置创建 URI,它在 DEFAULT_FS 为时查找键 FS_DEFAULT_NAME_KEY(fs.defaultFS):

      public static final String  FS_DEFAULT_NAME_DEFAULT = "file:///";
    

    【讨论】:

    • 感谢 ShemTov。我使用我的spark.sparkContext().hadoopConfiguration() 创建了我的FileSystem 对象,但fs.exists 仍然查看我的本地存储而不是HDFS。你有别的想法吗?
    • 是的,我猜这是因为您的本地环境在资源中没有 hdfs-site.xml / core-site.xml。您可以通过两种方式解决它: 1. 将 hdfs-site.xml 和 core-site.xml 添加到资源目录(如果您使用 intellij,请右键单击该目录 -> 将目录标记为资源)。 2. 有点作弊但你可以做 spark.sparkContext().hadoopConfiguration().set("fs.defaultFS","hdfs://server:port/") 然后将它传递给 FileSystem.get()。我更喜欢选项一。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-04-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多