【问题标题】:Hadoop Map Reduce CustomSplit/CustomRecordReaderHadoop Mapreduce 自定义拆分/自定义记录读取器
【发布时间】:2012-11-13 03:57:20
【问题描述】:

我有一个巨大的文本文件,我想拆分文件,以便每个块有 5 行。我实现了自己的 GWASInputFormat 和 GWASRecordReader 类。但是我的问题是,在以下代码(我从http://bigdatacircus.com/2012/08/01/wordcount-with-custom-record-reader-of-textinputformat/ 复制)中,在 initialize() 方法中,我有以下几行

FileSplit split = (FileSplit) genericSplit;
final Path file = split.getPath();
Configuration conf = context.getConfiguration();

我的问题是,在我的 GWASRecordReader 类中调用 initialize() 方法时,文件是否已经拆分?我以为我是在 GWASRecordReader 类中做的(拆分)。让我知道我的思考过程是否就在这里。

package com.test;

import java.io.IOException;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataInputStream;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.InputSplit;
import org.apache.hadoop.mapreduce.RecordReader;
import org.apache.hadoop.mapreduce.TaskAttemptContext;
import org.apache.hadoop.mapreduce.lib.input.FileSplit;
import org.apache.hadoop.util.LineReader;

public class GWASRecordReader extends RecordReader<LongWritable, Text> {

private final int NLINESTOPROCESS = 5;
private LineReader in;
private LongWritable key;
private Text value = new Text();
private long start = 0;
private long pos = 0;
private long end = 0;
private int maxLineLength;

public void close() throws IOException {
    if(in != null) {
        in.close();
    }
}

public LongWritable getCurrentKey() throws IOException, InterruptedException {
    return key;
}

public Text getCurrentValue() throws IOException, InterruptedException {
    return value;
}

public float getProgress() throws IOException, InterruptedException {
    if(start == end) {
        return 0.0f;
    }
    else {
        return Math.min(1.0f, (pos - start)/(float) (end - start));
    }
}

public void initialize(InputSplit genericSplit, TaskAttemptContext context) throws IOException {
    FileSplit split = (FileSplit) genericSplit;
    final Path file = split.getPath();
    Configuration conf = context.getConfiguration();
    this.maxLineLength = conf.getInt("mapred.linerecordreader.maxlength",Integer.MAX_VALUE);
    FileSystem fs = file.getFileSystem(conf);
    start = split.getStart();
    end = start + split.getLength();
    System.out.println("---------------SPLIT LENGTH---------------------" + split.getLength());
    boolean skipFirstLine = false;
    FSDataInputStream filein = fs.open(split.getPath());

    if(start != 0) {
        skipFirstLine = true;
        --start;
        filein.seek(start);
    }

    in = new LineReader(filein, conf);
    if(skipFirstLine) {
        start += in.readLine(new Text(),0,(int)Math.min((long)Integer.MAX_VALUE, end - start));
    }
    this.pos = start;
}

public boolean nextKeyValue() throws IOException, InterruptedException {
    if (key == null) {
        key = new LongWritable();
    }

    key.set(pos);

    if (value == null) {
        value = new Text();
    }
    value.clear();
    final Text endline = new Text("\n");
    int newSize = 0;
    for(int i=0; i<NLINESTOPROCESS;i++) {
        Text v = new Text();
        while( pos < end) {
            newSize = in.readLine(v ,maxLineLength, Math.max((int)Math.min(Integer.MAX_VALUE, end - pos), maxLineLength));
            value.append(v.getBytes(), 0, v.getLength());
            value.append(endline.getBytes(),0,endline.getLength());
            if(newSize == 0) {
                break;
            }
            pos += newSize;
            if(newSize < maxLineLength) {
                break;
            }
        }
    }

    if(newSize == 0) {
        key = null;
        value = null;
        return false;
    } else {
        return true;
    }
}
}

【问题讨论】:

    标签: java hadoop


    【解决方案1】:

    是的,输入文件已经被分割。基本上是这样的:

    your input file(s) -&gt; InputSplit -&gt; RecordReader -&gt; Mapper...

    基本上,InputSplit 将输入分解为块,RecordReader 将这些块分解为键/值对。请注意,InputSplit 和 RecordReader 将由您使用的 InputFormat 确定。例如,TextInputFormat 使用FileSplit 分解输入,然后LineRecordReader 使用位置作为键处理每一行,并将行本身作为值。 因此,在您的GWASInputFormat 中,您需要查看您使用哪种FileSplit 来查看它传递给GWASRecordReader 的内容。

    我建议查看NLineInputFormat,它“将 N 行输入拆分为一个拆分”。它也许可以完全按照您自己的意愿去做。

    如果您尝试一次获取 5 行作为值,并将第一行的行号作为键,我会说您可以使用自定义 NLineInputFormat 和自定义 LineRecordReader 来做到这一点。我认为您不必担心输入拆分,因为输入格式可以将其拆分为那 5 行块。您的RecordReader 将与LineRecordReader 非常相似,但不是获取块开始的字节位置,而是获取行号。所以除了那个小改动之外,代码几乎是相同的。因此,您基本上可以复制和粘贴NLineInputFormat 和LineRecordReader,然后让输入格式使用获取行号的记录阅读器。代码将非常相似。

    【讨论】:

    • 非常感谢。这清除了一些东西。我想跟踪输入文件的行号,并将行号与输入记录一起作为映射器的值输入。所以看起来我必须使用我自己的拆分,因为我现在所做的方式已经拆分了文件。你能告诉我我有什么选择吗(我想我需要重写 computeSplitSize() 方法)。我搜索了整个网络,但找不到具体的答案,如果我们能做到这一点
    • @user1707141 我更新了我的答案来解决这个问题。让我知道这是否有意义或者我需要更好地解释。
    • 这个答案的第一部分救了我的命,我遇到了类似的问题,我知道我关心的是 RecordReader。我想找到由字符串分隔的记录,所以我找到了这篇很棒的文章。我知道这是一个迟到的答案,但它可能对任何需要这个的人有用:hadoopi.wordpress.com/2013/05/31/…
    猜你喜欢
    • 2015-11-13
    • 1970-01-01
    • 2014-04-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-11-04
    • 2023-03-12
    相关资源
    最近更新 更多