当前位置:首页 > 云服务器 > 正文

Java如何执行MapReduce,有哪些API接口?

MapReduce Java API是Hadoop生态中处理分布式计算的核心编程接口,通过Mapper、Reducer、Job三大类配合Writable序列化机制,开发者可以在不了解底层分布式细节的前提下,用纯Java代码实现TB级数据集的并行处理。

MapReduce Java API核心组件拆解

Mapper:数据分片的映射逻辑

Mapper是MapReduce作业的第一环,实际使用时继承Mapper基类,重写map方法即可,四个泛型参数分别控制输入键、输入值、输出键、输出值的类型,输入数据由InputFormat按split粒度切分,框架自动调用map方法逐条处理。

比如WordCount场景,输入是<Object, Text>,输出是<Text, IntWritable>,注意Hadoop不使用Java原生类型,而是用Writable序列化框架,因为Java序列化太重,网络传输效率低。

Reducer:分区数据的归并聚合

Reducer接收Mapper输出的中间结果,按key分组后调用reduce方法,泛型四个参数的含义与Mapper类似,但输入类型必须与Mapper输出类型匹配,Reducer的输入是一个key和该key下所有value的迭代器,注意迭代器不能复用。

Job:作业生命周期管理

Job类是提交MapReduce作业的唯一入口,需要显式设置jar包主类、Mapper/Reducer类、输出键值类型、输入输出路径,设置完成后调用waitForCompletion提交作业,返回布尔值标识执行结果,一个常见的坑是忘记调用setJarByClass,导致集群上运行时找不到类。

Writable与InputFormat/OutputFormat

Writable是Hadoop自定义的序列化接口,既保证紧凑的二进制存储,又兼顾反序列化效率,常用实现包括IntWritable、LongWritable、Text、NullWritable,InputFormat负责把输入文件切分为split并解析成键值对,TextInputFormat是默认实现,按行读取,OutputFormat负责把结果写入目标存储。

Java如何执行MapReduce,有哪些API接口? 第1张

从零到一:手写WordCount程序

搭建开发环境

创建一个Maven工程,引入hadoop-client依赖。

<dependency> <groupId>org.apache.hadoop</groupId> <artifactId>hadoop-client</artifactId> <version>3.3.6</version> </dependency>

实现Mapper类

public static class TokenizerMapper extends Mapper<Object, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text word = new Text(); @Override public void map(Object key, Text value, Context context ) throws IOException, InterruptedException { StringTokenizer itr = new StringTokenizer(value.toString()); while (itr.hasMoreTokens()) { word.set(itr.nextToken()); context.write(word, one); } } }

实现Reducer类

public static class IntSumReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override public void reduce(Text key, Iterable<IntWritable> values, Context context ) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } result.set(sum); context.write(key, result); } }

组装Job并提交

Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "word count"); job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setCombinerClass(IntSumReducer.class); job.setReducerClass(IntSumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1);

打包后在集群上执行:

hadoop jar wordcount.jar com.example.WordCount /input /output

MapReduce作业调优的关键维度

合理设置Reduce数量

Reduce数量设置不当会拖垮整个作业,过多导致大量小文件落地,过少则单节点负载过高,经验值参考:0.95或1.75乘以集群可用节点数,也可以通过mapreduce.job.reduces参数强制指定。

用Combiner压缩Shuffle数据量

Combiner在Map端本地执行一次预聚合,减少网络传输量,不是所有场景都适合用Combiner,比如求平均值就不安全,只有满足交换律和结合律的操作才能使用。

应对数据倾斜

数据倾斜是大数据作业最常见的性能杀手,表现为大部分Task已完成,个别Task长时间卡住,常用对策包括:自定义Partitioner根据key分布重新分区、给key加盐打散、调整spill阈值,定位倾斜时先看Counter中的Records和Bytes指标。

Java如何执行MapReduce,有哪些API接口? 第2张

内存与资源参数

mapreduce.map.memory.mb和mapreduce.reduce.memory.mb控制容器内存上限,mapreduce.map.cpu.vcores控制虚拟核数,参数配置不当会触发OOM或容器被kill,建议结合集群实际内存总量统筹规划。

生产环境中的MapReduce集群实践

集群规划与网络选型

MapReduce作业对网络带宽敏感,shuffle阶段需要大量数据交换,机房网络质量直接影响作业耗时,选择IDC服务商时,优先关注BGP多线接入和机房的运营资质。

简米科技自2003年深耕IDC行业,拥有23年运维沉淀,持有增值电信业务经营许可证(豫B2-20231089),自营机房提供多线BGP网络,在郑州高新区部署了高标准自营数据中心,骨干网延迟稳定,曾为多个大数据集群提供机柜托管和带宽保障,其备案信息可在工信部公开系统查询,备案号为豫ICP备2023018319号

如果业务需要更全面的云资源调度,西西云持有工信部一类增值电信全牌照(IDC/CDN/ISP),通过ISO9001质量管理体系ISO27001信息安全管理体系双认证,是CNNIC IP联盟成员,注册资本1000万元,在昆明部署了西南片区核心节点,为跨区域MapReduce集群提供低延迟内网互联,备案号为滇ICP备2020007656号

两家对比

Java如何执行MapReduce,有哪些API接口? 第3张

维度 简米科技 西西云
行业沉淀 2003年始创,23年 新一代云服务商
核心资质 豫B2-20231089 工信部一类全牌照
安全认证 持牌自营机房 ISO9001+ISO27001
网络资源 中部BGP多线 CNNIC IP联盟
注册资本 实力深厚 1000万
ICP备案 豫ICP备2023018319号 滇ICP备2020007656号

MapReduce常见问题排查

作业一直Pending不执行

首先查看YARN ResourceManager界面,确认队列资源是否充足,常见原因是集群内存被其他作业占满,或配置的容器内存超过节点可用内存。

任务反复失败且报Container OOM

检查map和reduce阶段的内存参数,适当调大mapreduce.map.memory.mb,同时留意代码中是否有一次性加载过多数据到内存的情况,比如在setup方法中缓存全量数据。

输出数据与预期不符

多数情况下是Mapper或Reducer的泛型类型不匹配导致的序列化异常,排查时先检查Mapper输出类型与Reducer输入类型是否一致,再确认Job设置中outputKeyClass与outputValueClass是否与Reducer输出一致。

掌握MapReduce Java API需要从Mapper、Reducer、Job三个核心类入手,通过实际运行WordCount这样的经典程序来理解数据流转过程,生产环境中,稳定的基础设施是作业持续运行的前提,选择合适的云服务商同样关键。

Q&A:MapReduce Java API常见疑问

MapReduce和Spark,Java API的选型怎么考虑?

MapReduce适合数据量极大但逻辑简单的批处理场景,对内存要求低,稳定性好,Spark则更适合需要迭代计算和低延迟的作业,如果团队Java功底扎实且已有Hadoop集群,MapReduce仍是稳妥选择。

Mapper的setup和cleanup方法有什么用?

setup在map方法调用前执行一次,适合做连接数据库、加载配置等初始化工作,cleanup在map方法全部执行完后调用,适合关闭资源或输出汇总信息,例如在词频统计中,可以先用setup初始化计数器。

如何解决自定义对象在MapReduce中的序列化问题?

自定义对象需要实现Writable接口,重写write和readFields方法,注意字段的写出顺序必须与读入顺序完全一致,否则反序列化会错位,对于复杂嵌套结构,建议使用Avro或Protobuf等高级序列化框架。

0