千家信息网

FileInputFormat如何导读getSplits

发表于:2024-11-19 作者:千家信息网编辑
千家信息网最后更新 2024年11月19日,这篇文章主要为大家展示了"FileInputFormat如何导读getSplits",内容简而易懂,条理清晰,希望能够帮助大家解决疑惑,下面让小编带领大家一起研究并学习一下"FileInputForm
千家信息网最后更新 2024年11月19日FileInputFormat如何导读getSplits

这篇文章主要为大家展示了"FileInputFormat如何导读getSplits",内容简而易懂,条理清晰,希望能够帮助大家解决疑惑,下面让小编带领大家一起研究并学习一下"FileInputFormat如何导读getSplits"这篇文章吧。

/**
* Generate the list of files and make them into FileSplits.
* @param job the job context
* @throws IOException
*/
public List getSplits(JobContext job) throws IOException {
Stopwatch sw = new Stopwatch().start();
//获得一个InputSplit能够包含的最小值
long minSize = Math.max(getFormatMinSplitSize(), getMinSplitSize(job));
//获得一个InputSplit能够包含的最大值
long maxSize = getMaxSplitSize(job);

// generate splits
List splits = new ArrayList();
List files = listStatus(job);
/*
* 由此可知,如果有一百万个小文件,就会循环一百万次,并且至少生成一百万个InputSplit,就至少含有一百万个map任务
* 如果一个InputSplit的默认大小是一个block大小,即64M
* 一个20M的文件会产生一个InputSplit,一个Map任务
* 一个80M的文件会产生两个InputSplit,两个Map任务
* 两个分别为20M的文件总共产生两个InputSplit,两个Map任务
* 一个20M、一个70M的文件总共会产生三个InputSplit,三个Map任务
*/
for (FileStatus file: files) {
Path path = file.getPath();
long length = file.getLen();
if (length != 0) {
BlockLocation[] blkLocations;
if (file instanceof LocatedFileStatus) {
blkLocations = ((LocatedFileStatus) file).getBlockLocations();
} else {
FileSystem fs = path.getFileSystem(job.getConfiguration());
blkLocations = fs.getFileBlockLocations(file, 0, length);
} if (isSplitable(job, path)) {
//拿到hdfs默认的block块大小
long blockSize = file.getBlockSize();
//计算一个InputSplit的大小
long splitSize = computeSplitSize(blockSize, minSize, maxSize);

long bytesRemaining = length;
while (((double) bytesRemaining)/splitSize > SPLIT_SLOP) {
int blkIndex = getBlockIndex(blkLocations, length-bytesRemaining);
splits.add(makeSplit(path, length-bytesRemaining, splitSize,
blkLocations[blkIndex].getHosts(),
blkLocations[blkIndex].getCachedHosts()));
bytesRemaining -= splitSize;
}

if (bytesRemaining != 0) {
int blkIndex = getBlockIndex(blkLocations, length-bytesRemaining);
splits.add(makeSplit(path, length-bytesRemaining, bytesRemaining,
blkLocations[blkIndex].getHosts(),
blkLocations[blkIndex].getCachedHosts()));
}
} else { // not splitable
splits.add(makeSplit(path, 0, length, blkLocations[0].getHosts(),
blkLocations[0].getCachedHosts()));
}
} else {
//Create empty hosts array for zero length files
splits.add(makeSplit(path, 0, length, new String[0]));
}
}
// Save the number of input files for metrics/loadgen
job.getConfiguration().setLong(NUM_INPUT_FILES, files.size());
sw.stop();
if (LOG.isDebugEnabled()) {
LOG.debug("Total # of splits generated by getSplits: " + splits.size() + ", TimeTaken: " + sw.elapsedMillis());
}
return splits;
}

以上是"FileInputFormat如何导读getSplits"这篇文章的所有内容,感谢各位的阅读!相信大家都有了一定的了解,希望分享的内容对大家有所帮助,如果还想学习更多知识,欢迎关注行业资讯频道!

0