HDFSをJava APIで操作する

1. 環境の準備

Hadoopのインストールパッケージ内のshareディレクトリには、Hadoopに関連するすべてのコンテンツが含まれています。インストールパッケージのshareフォルダに入ると、dochadoopの2つのフォルダが見つかります。docにはHadoopの全ドキュメントが、hadoopフォルダにはHadoop開発に必要なすべてのJARファイルと、それらの依存関係が対応するlibフォルダに配置されています。

今回使用するのは、hadoopフォルダ内のcommonおよびhdfsフォルダの内容です。プロジェクトを作成する際には、これら2つのフォルダ内のJARファイルをプロジェクトに追加し、さらに各フォルダのlib内のすべてのJARファイルを依存関係としてプロジェクトに追加する必要があります(commonhdfslibにあるJARファイルに重複がある可能性があるため、重複を避けるように注意してください)。また、sourceフォルダには各モジュールのソースコードが含まれており、IDEでソースコードを参照したい場合は、これらの内容もプロジェクトに追加する必要があります。

2. HDFS操作の流れ

  1. HDFSの設定を取得し、HDFSオペレーティングシステム全体を取得します。
  2. HDFSオペレーティングシステムを開き、操作を実行します。
    1. フォルダの操作:作成、削除、変更、検索
    2. ファイルのアップロードとダウンロード
    3. ファイルのIO操作——HDFS間でのコピー

3. 具体的な操作例

1. HDFSファイルシステムオブジェクトの取得(3つの方法)


/**
 * HDFSファイルシステムオブジェクトを設定から取得します。
 * 3つの方法があります:
 *  1. 設定ファイルを直接使用する方法
 *  2. URIパスを指定し、設定ファイルからファイルシステムを作成する方法
 *  3. リソースファイルを直接追加する方法
 * @return FileSystem
 */
public FileSystem getHadoopFileSystem() {
    FileSystem hdfsClient = null;
    Configuration hadoopConfig = new Configuration();

    // 方法1:ローカルに設定ファイルがある場合、直接設定ファイルからHDFSオブジェクトを作成します。
    // この場合、hdfsのアクセスパスを指定する必要があります。
    hadoopConfig.set("fs.defaultFS", "hdfs://your-namenode:9000");

    try {
        hdfsClient = FileSystem.get(hadoopConfig);
    } catch (IOException e) {
        e.printStackTrace();
    }

    // 方法2:ローカルにHadoopシステムがないが、リモートアクセスが可能な場合。
    // 指定されたURIとユーザー名を使用して、リモートの設定情報にアクセスします。
    /*hadoopConfig = new Configuration();
    String hdfsUserName = "your-username";

    URI hdfsUri = null;
    try {
        hdfsUri = new URI("hdfs://your-namenode:9000");
    } catch (URISyntaxException e) {
        e.printStackTrace();
    }

    try {
        hdfsClient = FileSystem.get(hdfsUri, hadoopConfig, hdfsUserName);
    } catch (IOException | InterruptedException e) {
        e.printStackTrace();
    }*/

    // 方法3:リソースファイルを直接追加する方法
    /*hadoopConfig.addResource(new Path("/path/to/your/core-site.xml"));
    hadoopConfig.addResource(new Path("/path/to/your/hdfs-site.xml"));

    try {
        hdfsClient = FileSystem.get(hadoopConfig);
    } catch (IOException e) {
        e.printStackTrace();
    }*/

    return hdfsClient;
}

2. フォルダの作成


/**
 * 指定されたパスにディレクトリを作成します。これはシェルコマンドのmkdir -pと同様で、親ディレクトリが存在しない場合でも作成できます。
 * JavaのIO操作と同様に、ここでもPathオブジェクトに対して操作を行いますが、これはHDFSのPathオブジェクトです。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param pathStr 作成するディレクトリのパス
 * @return 作成が成功した場合はtrue、それ以外はfalse
 */
public boolean createDirectory(FileSystem hdfsClient, String pathStr) {
    boolean success = false;
    Path hdfsPath = new Path(pathStr);

    try {
        // パスが既に存在していても、このメソッドは正常に動作します。
        success = hdfsClient.mkdirs(hdfsPath);
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        try {
            hdfsClient.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
    return success;
}

3. フォルダの削除


/**
 * 指定されたパスのファイルまたはディレクトリを削除します。
 * Javaと同様に、hadoop.fsパッケージのPathオブジェクトが必要です。
 * delete(Path p)は非推奨となっており、delete(Path p, boolean recursive)が推奨されています。
 * 2番目のboolean引数は、ファイルの削除方法を制御し、rm -rコマンドと同等の再帰的削除を意味します。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param pathStr 削除するパス
 * @return 削除が成功した場合はtrue、それ以外はfalse
 */
public boolean deletePath(FileSystem hdfsClient, String pathStr) {
    boolean success = false;
    Path hdfsPath = new Path(pathStr);

    try {
        success = hdfsClient.delete(hdfsPath, true);
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        try {
            hdfsClient.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
    return success;
}

4. ファイルのリネーム


/**
 * ファイルまたはディレクトリの名前を変更します。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param oldPathStr 古いパス
 * @param newPathStr 新しいパス
 * @return リネームが成功した場合はtrue、それ以外はfalse
 */
public boolean renamePath(FileSystem hdfsClient, String oldPathStr, String newPathStr) {
    boolean success = false;
    Path oldPath = new Path(oldPathStr);
    Path newPath = new Path(newPathStr);

    try {
        success = hdfsClient.rename(oldPath, newPath);
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        try {
            hdfsClient.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
    return success;
}

5. フォルダの再帰的なリスト表示


/**
 * 指定されたディレクトリ内のファイルとサブディレクトリを再帰的にリスト表示します。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param rootPathStr ルートディレクトリのパス
 * @return ファイルとディレクトリのパスのセット
 */
public Set<String> listFilesRecursively(FileSystem hdfsClient, String rootPathStr) {
    Set<String> pathSet = new HashSet<>();
    Path rootPath = new Path(rootPathStr);

    try {
        FileStatus[] fileStatuses = hdfsClient.listStatus(rootPath);
        if (fileStatuses == null || fileStatuses.length == 0) {
            // ディレクトリが空の場合、そのパス自体を追加
            pathSet.add(rootPath.toUri().getPath());
        } else {
            for (FileStatus status : fileStatuses) {
                if (status.isFile()) {
                    // ファイルの場合、そのパスを追加
                    pathSet.add(status.getPath().toUri().getPath());
                } else {
                    // ディレクトリの場合、再帰的に処理
                    pathSet.addAll(listFilesRecursively(hdfsClient, status.getPath().toString()));
                }
            }
        }
    } catch (IOException e) {
        e.printStackTrace();
    }
    return pathSet;
}

6. ファイルの存在確認と種類の判定


/**
 * ファイルの存在、種類(ディレクトリかファイルか)を確認します。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param pathStr 確認するパス
 */
public void checkFileStatus(FileSystem hdfsClient, String pathStr) {
    boolean isExists = false;
    boolean isDirectory = false;
    boolean isFile = false;

    Path hdfsPath = new Path(pathStr);

    try {
        isExists = hdfsClient.exists(hdfsPath);
        isDirectory = hdfsClient.isDirectory(hdfsPath);
        isFile = hdfsClient.isFile(hdfsPath);
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        try {
            hdfsClient.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }

    if (!isExists) {
        System.out.println("指定されたパスは存在しません。");
    } else {
        System.out.println("パスが存在します。");
        if (isDirectory) {
            System.out.println("これはディレクトリです。");
        } else if (isFile) {
            System.out.println("これはファイルです。");
        }
    }
}

7. 設定情報の表示


/**
 * Hadoop設定のすべてのキーと値のペアを表示します。
 */
public void displayConfiguration() {
    Configuration hadoopConfig = new Configuration();
    hadoopConfig.set("fs.defaultFS", "hdfs://your-namenode:9000");

    Iterator> iterator = hadoopConfig.iterator();
    while (iterator.hasNext()) {
        Map.Entry entry = iterator.next();
        System.out.println(entry.getKey() + " = " + entry.getValue());
    }
}

8. ファイルのダウンロード


/**
 * HDFSからローカルファイルシステムにファイルをダウンロードします。
 * ダウンロード先のパスの最後の部分は、ダウンロードされるファイル名となります。
 * copyToLocalFile(Path src, Path dst)メソッドは、ファイルをコピーします。
 * 引数にboolean値を追加すると、moveToLocalFile()メソッドとして動作します。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param hdfsFilePath HDFS上のファイルパス
 * @param localDestinationPath ローカルの保存先ディレクトリパス
 */
public void downloadFileToLocal(FileSystem hdfsClient, String hdfsFilePath, String localDestinationPath) {
    Path srcPath = new Path(hdfsFilePath);
    Path dstPath = new Path(localDestinationPath);

    try {
        hdfsClient.copyToLocalFile(srcPath, dstPath);
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        try {
            hdfsClient.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

9. ファイルのアップロード


/**
 * ローカルファイルシステムからHDFSにファイルをアップロードします。
 * アップロード先のパスが存在しない場合は自動的に作成されます。
 * 同名のファイルが存在する場合は上書きされます。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param localFilePath アップロードするローカルファイルのパス
 * @param hdfsDestinationPath HDFS上の保存先パス
 */
public void uploadFileToHDFS(FileSystem hdfsClient, String localFilePath, String hdfsDestinationPath) {
    Path localPath = new Path(localFilePath);
    Path hdfsPath = new Path(hdfsDestinationPath);

    try {
        hdfsClient.copyFromLocalFile(localPath, hdfsPath);
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        try {
            hdfsClient.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

10. ファイルのコピー(HDFS間)


/**
 * HDFS内のファイルを別のパスにコピーします。
 * FSDataInputStreamを使用してファイルを開き、FSDataOutputStreamを使用して書き込み先を開きます。
 * IOUtils.copyBytes()を使用して、バッファリングされた読み書きを実行します。
 * @param hdfsClient HDFSファイルシステムクライアント
 * @param sourcePath コピー元のHDFSパス
 * @param destinationPath コピー先のHDFSパス
 */
public void copyFileWithinHDFS(FileSystem hdfsClient, String sourcePath, String destinationPath) {
    Path srcPath = new Path(sourcePath);
    Path dstPath = new Path(destinationPath);

    FSDataInputStream hdfsInputStream = null;
    FSDataOutputStream hdfsOutputStream = null;

    try {
        hdfsInputStream = hdfsClient.open(srcPath);
        hdfsOutputStream = hdfsClient.create(dstPath);

        // バッファサイズを指定してコピーを実行します。最後のboolean引数はストリームをクローズするかどうかを指定します。
        IOUtils.copyBytes(hdfsInputStream, hdfsOutputStream, 1024 * 1024 * 64, false);
    } catch (IOException e) {
        e.printStackTrace();
    } finally {
        try {
            if (hdfsOutputStream != null) hdfsOutputStream.close();
            if (hdfsInputStream != null) hdfsInputStream.close();
        } catch (IOException e) {
            e.printStackTrace();
        }
    }
}

タグ: Hadoop HDFS Java FileSystem API

8月2日 22:50 投稿