Hadoop输入格式深度解析:默认输入格式与自定义实战

Hadoop输入格式深度解析:默认输入格式与自定义实战

    • 引言:输入格式——MapReduce的数据入口
    • 一、Hadoop的默认输入格式:TextInputFormat
      • 1.1 默认输入格式概述
      • 1.2 工作原理
      • 1.3 其他内置输入格式
      • 1.4 TextInputFormat的优缺点
    • 二、为什么需要自定义输入格式?
      • 2.1 典型场景
      • 2.2 内置格式无法满足的需求
    • 三、自定义输入格式的核心接口
      • 3.1 InputFormat类体系
      • 3.2 需要实现的核心方法
        • **1. 继承FileInputFormat,重写createRecordReader**
        • **2. 实现RecordReader,定义解析逻辑**
    • 四、实战案例一:处理非标准行分隔符
      • 4.1 场景描述
      • 4.2 自定义RecordReader实现
      • 4.3 自定义InputFormat
      • 4.4 在作业中使用
    • 五、实战案例二:处理固定宽度文件
      • 5.1 场景描述
      • 5.2 自定义RecordReader实现
    • 六、实战案例三:合并小文件
      • 6.1 场景描述
      • 6.2 自定义CombineFileInputFormat
      • 6.3 ByteArrayRecordReader实现
    • 七、自定义输入格式的最佳实践
      • 7.1 性能考虑
      • 7.2 可配置性设计
      • 7.3 测试与验证
    • 八、总结
      • 8.1 核心要点回顾
      • 8.2 选择指南

🌺The Begin🌺点点关注,收藏不迷路🌺

引言:输入格式——MapReduce的数据入口

在Hadoop MapReduce框架中,输入格式(InputFormat) 是数据处理的起点,它决定了如何读取、分割和解析输入数据。正确选择或设计输入格式,直接影响作业的并行度、数据本地性和处理效率。

本文将深入剖析Hadoop的默认输入格式TextInputFormat的工作原理,并通过完整的实战案例展示如何自定义输入格式来解决特定业务场景中的问题。

输入格式InputFormat

核心功能

验证输入路径

将文件切分为分片

提供RecordReader

记录解析逻辑

默认实现

TextInputFormat

KeyValueTextInputFormat

NLineInputFormat

SequenceFileInputFormat

自定义扩展

继承FileInputFormat

实现RecordReader

定义键值类型

控制分片逻辑

一、Hadoop的默认输入格式:TextInputFormat

1.1 默认输入格式概述

Hadoop的默认输入格式是TextInputFormat 。它是FileInputFormat的子类,专门用于处理纯文本文件。

核心特性

  • 将输入文件按行分割
  • 每行作为一条记录
  • 键(Key):该行在文件中的字节偏移量(LongWritable类型)
  • 值(Value):该行的文本内容(Text类型)
// TextInputFormat的典型输出
// <0, "Hello Hadoop">
// <15, "Welcome to MapReduce">
// <35, "Data processing framework">

1.2 工作原理

记录解析

TextInputFormat处理

输入文件

文件 data.txt
512MB

分片计算
默认等于块大小128MB

分片1
0-128MB

分片2
128-256MB

分片3
256-384MB

分片4
384-512MB

LineRecordReader

LineRecordReader

LineRecordReader

LineRecordReader

逐行读取

<0,'line1'>

<15,'line2'>

关键组件

  1. getSplits():将文件切分为逻辑分片(默认按块大小)
  2. createRecordReader():为每个分片创建LineRecordReader
  3. LineRecordReader:负责实际的行读取,处理跨分片行边界

1.3 其他内置输入格式

输入格式 键类型 值类型 适用场景
TextInputFormat 行偏移量(Long) 行内容(Text) 日志文件、CSV等文本文件
KeyValueTextInputFormat 行中第一个分隔符前的内容 行中剩余内容 键值对格式的数据
NLineInputFormat 行偏移量(Long) 行内容(Text) 每N行作为一个分片
SequenceFileInputFormat 自定义 自定义 Hadoop二进制序列文件
CombineTextInputFormat 行偏移量(Long) 行内容(Text) 合并小文件,减少Map数

1.4 TextInputFormat的优缺点

优点

  • 简单通用,人类可读
  • 支持所有文本格式
  • 易于调试(可直接用hadoop fs -cat查看)

缺点

  • 存储效率低(比二进制格式多占用30-50%空间)
  • 解析开销大(字符串转换、类型判断)
  • 不保留数据类型信息
  • 不支持复杂数据结构嵌套

二、为什么需要自定义输入格式?

2.1 典型场景

自定义输入格式场景

特殊文件格式

非标准分隔符

固定宽度记录

二进制格式

自定义序列化

特殊解析逻辑

跳过文件头

按业务规则分组

数据过滤

格式转换

性能优化

合并小文件

预聚合

索引读取

列式裁剪

2.2 内置格式无法满足的需求

场景 问题 解决方案
非\n行分隔符 某些文件使用特殊符号(如!@!)作为行分隔符 自定义RecordReader,识别特殊分隔符
二进制格式 图片、音频、视频文件 自定义InputFormat,直接读取二进制数据
固定宽度文件 每行按固定列宽存储 自定义RecordReader,按宽度切分字段
跳过文件头 CSV文件包含列名行 自定义RecordReader,跳过前N行
多行记录 XML/JSON跨越多行 自定义RecordReader,按逻辑记录切分

三、自定义输入格式的核心接口

3.1 InputFormat类体系

InputFormat抽象类

FileInputFormat
基于文件的实现

DBInputFormat
数据库输入

TextInputFormat
默认

KeyValueTextInputFormat
键值对

SequenceFileInputFormat
序列文件

自定义InputFormat
继承实现

3.2 需要实现的核心方法

自定义输入格式通常需要继承FileInputFormat并实现两个核心组件:

1. 继承FileInputFormat,重写createRecordReader
public class CustomInputFormat extends FileInputFormat<KeyType, ValueType> {
    @Override
    public RecordReader<KeyType, ValueType> createRecordReader(
            InputSplit split, TaskAttemptContext context) 
            throws IOException, InterruptedException {
        // 返回自定义的RecordReader实例
        return new CustomRecordReader();
    }
    // 可选:重写isSplitable控制文件是否可切分
    @Override
    protected boolean isSplitable(JobContext context, Path filename) {
        // 对于不可切分的格式(如gzip压缩),返回false
        return super.isSplitable(context, filename);
    }
    // 可选:重写getSplits控制分片逻辑
    @Override
    public List<InputSplit> getSplits(JobContext job) throws IOException {
        // 自定义分片策略,如合并小文件
        return super.getSplits(job);
    }
}
2. 实现RecordReader,定义解析逻辑
public class CustomRecordReader extends RecordReader<KeyType, ValueType> {
    // 初始化方法,在任务开始时调用
    @Override
    public void initialize(InputSplit split, TaskAttemptContext context) 
            throws IOException, InterruptedException {
        // 打开文件流,定位到分片起始位置
    }
    // 读取下一条记录,返回true表示还有记录
    @Override
    public boolean nextKeyValue() throws IOException, InterruptedException {
        // 读取下一条记录,填充key和value
    }
    // 获取当前键
    @Override
    public KeyType getCurrentKey() throws IOException, InterruptedException {
        return key;
    }
    // 获取当前值
    @Override
    public ValueType getCurrentValue() throws IOException, InterruptedException {
        return value;
    }
    // 获取处理进度
    @Override
    public float getProgress() throws IOException, InterruptedException {
        // 返回0.0-1.0之间的进度值
    }
    // 关闭资源
    @Override
    public void close() throws IOException {
        // 关闭文件流等资源
    }
}

四、实战案例一:处理非标准行分隔符

4.1 场景描述

某些数据文件不是以换行符\n分隔行,而是使用特殊字符组合如!@!作为行分隔符。默认的TextInputFormat无法正确解析这种文件。

4.2 自定义RecordReader实现

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 java.io.IOException;
public class CustomDelimiterRecordReader extends RecordReader<LongWritable, Text> {
    private LongWritable key;
    private Text value;
    private FSDataInputStream inputStream;
    private long start;
    private long end;
    private long pos;
    private byte[] delimiter;
    private int delimiterLength;
    @Override
    public void initialize(InputSplit split, TaskAttemptContext context) 
            throws IOException, InterruptedException {
        // 获取分片信息
        FileSplit fileSplit = (FileSplit) split;
        Path path = fileSplit.getPath();
        Configuration conf = context.getConfiguration();
        // 获取自定义分隔符(可从配置中读取)
        String delimStr = conf.get("custom.record.delimiter", "!@!");
        delimiter = delimStr.getBytes("UTF-8");
        delimiterLength = delimiter.length;
        // 打开文件流
        FileSystem fs = path.getFileSystem(conf);
        inputStream = fs.open(path);
        // 定位到分片起始位置
        start = fileSplit.getStart();
        end = start + fileSplit.getLength();
        inputStream.seek(start);
        pos = start;
        key = new LongWritable();
        value = new Text();
    }
    @Override
    public boolean nextKeyValue() throws IOException, InterruptedException {
        if (pos >= end) {
            return false;
        }
        // 读取数据直到遇到分隔符或到达分片末尾
        ByteArrayOutputStream buffer = new ByteArrayOutputStream();
        // 处理可能的分片边界情况
        boolean boundaryFound = false;
        while (pos < end && !boundaryFound) {
            int b = inputStream.read();
            if (b == -1) {
                break;
            }
            pos++;
            buffer.write(b);
            // 检查是否匹配分隔符
            if (b == delimiter[delimiterLength - 1]) {
                byte[] data = buffer.toByteArray();
                if (endsWithDelimiter(data)) {
                    // 找到完整记录,去掉末尾的分隔符
                    byte[] record = new byte[data.length - delimiterLength];
                    System.arraycopy(data, 0, record, 0, data.length - delimiterLength);
                    value.set(record);
                    boundaryFound = true;
                }
            }
        }
        // 设置键为记录的起始偏移量
        key.set(pos - buffer.size());
        return true;
    }
    private boolean endsWithDelimiter(byte[] data) {
        if (data.length < delimiterLength) {
            return false;
        }
        for (int i = 0; i < delimiterLength; i++) {
            if (data[data.length - delimiterLength + i] != delimiter[i]) {
                return false;
            }
        }
        return true;
    }
    @Override
    public LongWritable getCurrentKey() {
        return key;
    }
    @Override
    public Text getCurrentValue() {
        return value;
    }
    @Override
    public float getProgress() {
        if (start == end) {
            return 0.0f;
        }
        return Math.min(1.0f, (pos - start) / (float) (end - start));
    }
    @Override
    public void close() throws IOException {
        if (inputStream != null) {
            inputStream.close();
        }
    }
}

4.3 自定义InputFormat

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.FileInputFormat;
public class CustomDelimiterInputFormat extends FileInputFormat<LongWritable, Text> {
    @Override
    public RecordReader<LongWritable, Text> createRecordReader(
            InputSplit split, TaskAttemptContext context) {
        // 返回自定义的RecordReader
        return new CustomDelimiterRecordReader();
    }
    // 确保文件可切分
    @Override
    protected boolean isSplitable(JobContext context, Path filename) {
        // 根据实际情况返回,如果是压缩文件可能返回false
        return true;
    }
}

4.4 在作业中使用

public class CustomDelimiterJob {
    public static void main(String[] args) throws Exception {
        Configuration conf = new Configuration();
        // 设置自定义分隔符
        conf.set("custom.record.delimiter", "!@!");
        Job job = Job.getInstance(conf, "Custom Delimiter Job");
        // 设置自定义输入格式
        job.setInputFormatClass(CustomDelimiterInputFormat.class);
        // 其他配置...
        job.setJarByClass(CustomDelimiterJob.class);
        job.setMapperClass(MyMapper.class);
        FileInputFormat.addInputPath(job, new Path(args[0]));
        FileOutputFormat.setOutputPath(job, new Path(args[1]));
        System.exit(job.waitForCompletion(true) ? 0 : 1);
    }
}

五、实战案例二:处理固定宽度文件

5.1 场景描述

传统系统导出的固定宽度文件,每条记录长度固定(如100字节),需要按列宽解析。

5.2 自定义RecordReader实现

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 java.io.IOException;
public class FixedWidthRecordReader extends RecordReader<LongWritable, Text> {
    private LongWritable key;
    private Text value;
    private FSDataInputStream inputStream;
    private long start;
    private long end;
    private long pos;
    private int recordLength;
    private byte[] recordBuffer;
    @Override
    public void initialize(InputSplit split, TaskAttemptContext context) 
            throws IOException, InterruptedException {
        FileSplit fileSplit = (FileSplit) split;
        Path path = fileSplit.getPath();
        Configuration conf = context.getConfiguration();
        // 获取记录长度配置(字节数)
        recordLength = conf.getInt("fixed.record.length", 100);
        recordBuffer = new byte[recordLength];
        // 打开文件流
        FileSystem fs = path.getFileSystem(conf);
        inputStream = fs.open(path);
        // 定位到分片起始位置
        start = fileSplit.getStart();
        end = start + fileSplit.getLength();
        // 对齐到记录边界
        if (start % recordLength != 0) {
            // 如果不是记录边界,向前移动到下一条记录起始
            start = start - (start % recordLength) + recordLength;
        }
        inputStream.seek(start);
        pos = start;
        key = new LongWritable();
        value = new Text();
    }
    @Override
    public boolean nextKeyValue() throws IOException, InterruptedException {
        if (pos >= end) {
            return false;
        }
        // 读取固定长度的记录
        int bytesRead = inputStream.read(recordBuffer, 0, recordLength);
        if (bytesRead < recordLength) {
            return false; // 不完整的记录
        }
        // 设置键为记录起始偏移量
        key.set(pos);
        // 设置值为记录内容
        value.set(recordBuffer, 0, recordLength);
        pos += recordLength;
        return true;
    }
    // 其他方法实现...
}

六、实战案例三:合并小文件

6.1 场景描述

HDFS中大量小文件会导致NameNode内存压力和过多的Map任务。CombineFileInputFormat可以将多个小文件合并到一个分片处理。

6.2 自定义CombineFileInputFormat

import org.apache.hadoop.io.BytesWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.lib.input.CombineFileInputFormat;
import org.apache.hadoop.mapreduce.lib.input.CombineFileSplit;
import org.apache.hadoop.mapreduce.InputSplit;
import org.apache.hadoop.mapreduce.RecordReader;
import org.apache.hadoop.mapreduce.TaskAttemptContext;
import java.io.IOException;
public class CustomCombineInputFormat 
        extends CombineFileInputFormat<Text, BytesWritable> {
    public CustomCombineInputFormat() {
        // 设置最大分片大小(例如256MB)
        super.setMaxSplitSize(268435456);
    }
    @Override
    public RecordReader<Text, BytesWritable> createRecordReader(
            InputSplit split, TaskAttemptContext context) throws IOException {
        // 使用CombineFileRecordReader处理合并后的分片
        return new CombineFileRecordReader<>(
            (CombineFileSplit) split,
            context,
            ByteArrayRecordReader.class
        );
    }
}

6.3 ByteArrayRecordReader实现

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.BytesWritable;
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.CombineFileSplit;
import java.io.IOException;
public class ByteArrayRecordReader extends RecordReader<Text, BytesWritable> {
    private CombineFileSplit split;
    private int currentIndex;
    private float currentProgress;
    private Text currentKey;
    private BytesWritable currentValue;
    @Override
    public void initialize(InputSplit genericSplit, TaskAttemptContext context) 
            throws IOException {
        this.split = (CombineFileSplit) genericSplit;
        this.currentIndex = 0;
        this.currentProgress = 0;
    }
    @Override
    public boolean nextKeyValue() throws IOException {
        if (currentIndex >= split.getNumPaths()) {
            return false;
        }
        // 获取当前文件的路径和偏移
        Path filePath = split.getPath(currentIndex);
        long offset = split.getOffset(currentIndex);
        long length = split.getLength(currentIndex);
        // 读取文件内容
        Configuration conf = new Configuration();
        FileSystem fs = filePath.getFileSystem(conf);
        FSDataInputStream in = fs.open(filePath);
        byte[] content = new byte[(int) length];
        in.readFully(offset, content);
        in.close();
        // 设置键为文件名
        currentKey = new Text(filePath.getName());
        // 设置值为文件内容
        currentValue = new BytesWritable(content);
        currentIndex++;
        currentProgress = (float) currentIndex / split.getNumPaths();
        return true;
    }
    @Override
    public Text getCurrentKey() {
        return currentKey;
    }
    @Override
    public BytesWritable getCurrentValue() {
        return currentValue;
    }
    @Override
    public float getProgress() {
        return currentProgress;
    }
    @Override
    public void close() {
        // 不需要额外关闭
    }
}

七、自定义输入格式的最佳实践

7.1 性能考虑

优化点 建议 说明
缓冲区大小 使用8KB-64KB缓冲区 减少磁盘IO调用次数
字节操作 避免频繁的字节数组拷贝 使用ByteArrayOutputStream时注意容量
对象重用 重用Key和Value对象 减少GC压力
分片边界处理 正确处理跨分片的记录 避免数据丢失或重复
压缩文件 不可切分的格式返回isSplitable=false 确保数据完整性

7.2 可配置性设计

public class ConfigurableRecordReader extends RecordReader {
    private String encoding;
    private String delimiter;
    private boolean skipHeader;
    @Override
    public void initialize(InputSplit split, TaskAttemptContext context) {
        Configuration conf = context.getConfiguration();
        // 从配置读取参数,提供默认值
        encoding = conf.get("custom.encoding", "UTF-8");
        delimiter = conf.get("custom.delimiter", ",");
        skipHeader = conf.getBoolean("custom.skip.header", false);
        // 使用这些参数初始化
    }
}

7.3 测试与验证

public class CustomInputFormatTest {
    @Test
    public void testRecordReader() throws Exception {
        // 创建测试文件
        Path testFile = new Path("test.txt");
        FileSystem fs = FileSystem.getLocal(new Configuration());
        FSDataOutputStream out = fs.create(testFile);
        out.writeBytes("record1!@!record2!@!record3");
        out.close();
        // 创建分片
        FileSplit split = new FileSplit(testFile, 0, 
            fs.getFileStatus(testFile).getLen(), new String[0]);
        // 创建RecordReader
        CustomDelimiterRecordReader reader = new CustomDelimiterRecordReader();
        reader.initialize(split, new TaskAttemptContext());
        // 验证读取结果
        assertTrue(reader.nextKeyValue());
        assertEquals(0L, reader.getCurrentKey().get());
        assertEquals("record1", reader.getCurrentValue().toString());
        // 清理
        reader.close();
        fs.delete(testFile, true);
    }
}

八、总结

8.1 核心要点回顾

输入格式总结

默认格式

TextInputFormat

行偏移量为键

行内容为值

适用于文本文件

内置格式

KeyValueTextInputFormat

NLineInputFormat

SequenceFileInputFormat

CombineTextInputFormat

自定义场景

特殊分隔符

固定宽度

二进制格式

合并小文件

实现要点

继承FileInputFormat

实现RecordReader

处理分片边界

配置参数化

8.2 选择指南

场景 推荐方案 理由
普通文本文件 TextInputFormat 简单够用,无需自定义
键值对格式 KeyValueTextInputFormat 内置支持,配置简单
每N行一个Map NLineInputFormat 精准控制并行度
大量小文件 CombineTextInputFormat 合并处理,减少Map数
二进制格式 SequenceFileInputFormat Hadoop原生,高效紧凑
特殊格式 自定义InputFormat 完全控制解析逻辑

核心启示:输入格式是MapReduce作业的"第一公里"。选择正确的输入格式能事半功倍,而自定义输入格式则是解决复杂数据格式问题的终极武器。理解其设计原理,能让你的Hadoop作业更加高效、健壮。


互动问题:你在实际项目中遇到过哪些需要自定义输入格式的场景?是处理XML/JSON多行记录,还是解析老旧系统的固定宽度文件?欢迎在评论区分享你的经验和解决方案!

在这里插入图片描述

🌺The End🌺点点关注,收藏不迷路🌺
© 版权声明

相关文章