【问题标题】:Merging hdfs files合并 hdfs 文件
【发布时间】:2019-02-23 02:27:49
【问题描述】:

我在 HDFS 中有 1000 多个文件可用,命名约定为 1_fileName.txtN_fileName.txt。每个文件的大小为 1024 MB。 我需要将这些文件合并为一个(HDFS)并保持文件的顺序。说5_FileName.txt 应该只附加在4_fileName.txt 之后

执行此操作的最佳和最快方法是什么。

有没有什么方法可以在不复制数据节点之间的实际数据的情况下执行这种合并?例如:获取此文件的块位置并在名称节点中使用这些块位置创建一个新条目(文件名)?

【问题讨论】:

标签: hadoop hdfs


【解决方案1】:

没有有效的方法来做到这一点,您需要将所有数据移动到一个节点,然后再回到 HDFS。

执行此操作的命令行脚本如下:

hadoop fs -text *_fileName.txt | hadoop fs -put - targetFilename.txt

这会将匹配 glob 的所有文件分类到标准输出,然后您将通过管道将该流传输到 put 命令并将流输出到名为 targetFilename.txt 的 HDFS 文件

您遇到的唯一问题是您所采用的文件名结构 - 如果您有固定宽度,将数字部分补零会更容易,但在当前状态下,您会得到一个意想不到的字典顺序(1、10、 100、1000、11、110 等)而不是数字顺序(1、2、3、4 等)。您可以通过将 scriptlet 修改为:

hadoop fs -text [0-9]_fileName.txt [0-9][0-9]_fileName.txt \
    [0-9][0-9[0-9]_fileName.txt | hadoop fs -put - targetFilename.txt

【讨论】:

  • 一个小修复:要分类多个文件,我们应该使用 'hadoop fs -cat' 而不是 'hadoop fs -text'。顺便说一句,我喜欢从非零填充文件名中获取正确顺序的技巧。
  • text 和 cat 是一样的,但是 text 也适用于压缩和序列文件
【解决方案2】:

有一个API方法org.apache.hadoop.fs.FileUtil.copyMerge可以执行这个操作:

public static boolean copyMerge(
                    FileSystem srcFS,
                    Path srcDir,
                    FileSystem dstFS,
                    Path dstFile,
                    boolean deleteSource,
                    Configuration conf,
                    String addString)

它按字母顺序读取srcDir 中的所有文件,并将其内容附加到dstFile。

【讨论】:

【解决方案3】:

如果你可以使用火花。可以这样做

sc.textFile("hdfs://...../part*).coalesce(1).saveAsTextFile("hdfs://...../filename)

希望这可行,因为 spark 以分布式方式工作,您不必将文件复制到一个节点中。虽然只是一个警告,但如果文件非常大,在 spark 中合并文件可能会很慢。

【讨论】:

  • 这能保证特定的订单吗?这是 OP 要求的一部分。
  • 不是!即使 Spark shuffle 非常高效(可以压缩数据),这也不会尊重顺序。
【解决方案4】:

由于文件顺序很重要,而字典顺序不能满足目的,因此为这个任务编写一个映射器程序似乎是一个不错的选择,它可能会定期运行。 当然没有reducer,把它写成一个HDFS map任务是高效的,因为它可以将这些文件合并到一个输出文件中,而无需跨数据节点移动太多数据。由于源文件位于 HDFS 中,并且由于 mapper 任务会尝试数据亲和性,因此它可以合并文件而无需跨不同数据节点移动文件。

映射程序将需要一个自定义 InputSplit(在输入目录中获取文件名并根据需要对其进行排序)和一个自定义 InputFormat。

映射器可以使用 hdfs append 或原始输出流,它可以写入 byte[]。

我正在考虑的 Mapper 程序的粗略草图是这样的:

public class MergeOrderedFileMapper extends MapReduceBase implements Mapper<ArrayWritable, Text, ??, ??> 
{
    FileSystem fs;

    public void map(ArrayWritable sourceFiles, Text destFile, OutputCollector<??, ??> output, Reporter reporter) throws IOException 
    {

        //Convert the destFile to Path.
        ...
        //make sure the parent directory of destFile is created first.
        FSDataOutputStream destOS = fs.append(destFilePath);
        //Convert the sourceFiles to Paths.
        List<Path> srcPaths;
        ....
        ....
            for(Path p: sourcePaths) {

                FSDataInputStream srcIS = fs.open(p);
                byte[] fileContent
                srcIS.read(fileContent);
                destOS.write(fileContent);
                srcIS.close();
                reporter.progress();  // Important, else mapper taks may timeout.
            }
            destOS.close();


        // Delete source files.

        for(Path p: sourcePaths) {
            fs.delete(p, false);
            reporter.progress();
        }

    }
}

【讨论】:

    【解决方案5】:

    我为 PySpark 编写了一个实现,因为我们经常使用它。

    以 Hadoop 的 copyMerge() 为模型,并使用相同的低级别 Hadoop API 来实现这一目标。

    https://github.com/Tagar/abalon/blob/v2.3.3/abalon/spark/sparkutils.py#L335

    它保持文件名的字母顺序。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-11-08
      • 2021-12-16
      • 2017-07-20
      • 2013-10-30
      • 1970-01-01
      • 1970-01-01
      • 2021-01-31
      • 2018-03-31
      相关资源
      最近更新 更多