华为云MRS集群HDFS与MapReduce开发实践
环境准备
弹性公网IP配置
在华为云控制台中,进入网络服务模块,选择弹性公网IP EIP服务。配置参数如下:
- 计费方式:按需付费
- 地域:华北-北京四
- 线路类型:全动态BGP
- 带宽计费:按实际流量
- 带宽峰值:100Mbps
- IPv6支持:关闭
- 资源命名:eip-bigdata1
- 购买数量:1个
MRS集群部署
进入大数据服务模块,选择MapReduce服务,采用自定义购买模式。核心配置项:
- 地域与计费:华北-北京四,按需付费
- 集群标识:mrs-bigdata
- 版本选择:MRS 3.1.0 WXL 普通版
- 组件勾选:Hadoop、HBase、Hive、Zookeeper、Ranger
- 网络环境:vpc-bigdata虚拟私有云,subnet-bigdata子网,sg-bigdata安全组
- 公网访问:绑定已创建的弹性公网IP
节点规格配置:
- 实例类型:通用计算增强型 c6.2xlarge.4(8核32GB)
- Master节点:3台,系统盘480GB高IO,数据盘600GB高IO
- Core节点:2台,相同磁盘配置
- 拓扑优化:在Master节点启用DN、NM、RS角色
安全认证设置:
- Kerberos认证:禁用
- 管理账号:admin(自定义强密码)
- 系统账号:root(自定义强密码)
JDK环境部署
集群默认仅含JRE运行时,需手动安装JDK开发套件。执行以下命令获取JDK 1.8:
wget https://sandbox-expriment-files.obs.cn-north-1.myhuaweicloud.com/hccdp/HCCDP/jdk-8u341-linux-x64.tar.gz
解压至目标目录:
tar -zxvf jdk-8u341-linux-x64.tar.gz -C /home/user/
HDFS文件系统操作实战
HDFS作为Hadoop生态的分布式存储底座,为上层计算引擎提供高可靠的数据存取能力。以下通过Java API演示核心文件操作。
Maven工程搭建
启动Eclipse IDE,创建Maven项目:
- File → New → Project → Maven Project
- 勾选"Create a simple project"选项
- 填写工程元数据:
- GroupId: com.huawei
- ArtifactId: HDFSClient
- Version: 1.0-SNAPSHOT
- Packaging: jar
开发环境配置
JDK版本切换:进入Window → Preferences → Java → Compiler,将合规级别调整为1.8。随后在Installed JREs中添加本地JDK路径/home/user/jdk1.8.0_341,并设为默认运行时。
构建路径更新:右键项目 → Build Path → Configure Build Path,移除旧版JRE系统库,添加新配置的JDK 1.8环境。
依赖管理配置
编辑pom.xml引入Hadoop客户端库:
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.huawei</groupId>
<artifactId>HDFSClient</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<maven.compiler.source>1.8</maven.compiler.source>
<maven.compiler.target>1.8</maven.compiler.target>
</properties>
<repositories>
<repository>
<id>huaweicloud</id>
<url>https://mirrors.huaweicloud.com/repository/maven/</url>
</repository>
</repositories>
<dependencies>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.1.1</version>
</dependency>
</dependencies>
<build>
<plugins>
<plugin>
<artifactId>maven-assembly-plugin</artifactId>
<configuration>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
</configuration>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>assembly</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
核心功能实现
路径存在性检测
创建com.huawei.hdfs包,新建PathChecker类:
package com.huawei.hdfs;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import java.io.IOException;
public class PathChecker {
public static void main(String[] args) throws IOException {
Configuration cfg = new Configuration();
FileSystem hdfs = FileSystem.get(cfg);
Path dirPath = new Path("/user/demo/hdfs");
Path filePath = new Path("/user/demo/hdfs/sample.txt");
System.out.println(dirPath + (hdfs.exists(dirPath) ? " 目录已存在" : " 目录不存在"));
System.out.println(filePath + (hdfs.exists(filePath) ? " 文件已存在" : " 文件不存在"));
hdfs.close();
}
}
空文件创建
新建EmptyFileCreator类:
package com.huawei.hdfs;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import java.io.IOException;
public class EmptyFileCreator {
public static void main(String[] args) throws IOException {
Configuration cfg = new Configuration();
FileSystem hdfs = FileSystem.get(cfg);
Path target = new Path("/user/demo/hdfs/sample.txt");
boolean created = hdfs.createNewFile(target);
System.out.println(created ? "文件创建成功" : "创建失败,文件已存在");
hdfs.close();
}
}
带内容文件写入
新建ContentFileWriter类:
package com.huawei.hdfs;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataOutputStream;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import java.io.BufferedWriter;
import java.io.IOException;
import java.io.OutputStreamWriter;
public class ContentFileWriter {
public static void main(String[] args) throws IOException {
Configuration cfg = new Configuration();
FileSystem hdfs = FileSystem.get(cfg);
Path target = new Path("/user/demo/hdfs/data.txt");
FSDataOutputStream outStream = hdfs.create(target);
BufferedWriter writer = new BufferedWriter(new OutputStreamWriter(outStream));
String[] lines = {"hadoop", "spark", "flink"};
for (String line : lines) {
writer.write(line);
writer.newLine();
}
writer.close();
outStream.close();
hdfs.close();
System.out.println("数据写入完成: " + target);
}
}
文件内容读取
新建FileContentReader类:
package com.huawei.hdfs;
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 java.io.BufferedReader;
import java.io.IOException;
import java.io.InputStreamReader;
public class FileContentReader {
public static void main(String[] args) throws IOException {
Configuration cfg = new Configuration();
FileSystem hdfs = FileSystem.get(cfg);
Path source = new Path(args[0]);
FSDataInputStream inStream = hdfs.open(source);
BufferedReader reader = new BufferedReader(new InputStreamReader(inStream));
String line;
while ((line = reader.readLine()) != null) {
System.out.println(line);
}
reader.close();
inStream.close();
hdfs.close();
}
}
文件删除操作
新建ResourceRemover类:
package com.huawei.hdfs;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import java.io.IOException;
public class ResourceRemover {
public static void main(String[] args) throws IOException {
Configuration cfg = new Configuration();
FileSystem hdfs = FileSystem.get(cfg);
Path target = new Path(args[0]);
if (hdfs.exists(target)) {
hdfs.delete(target, true);
System.out.println("删除成功");
} else {
System.out.println("目标不存在");
}
hdfs.close();
}
}
部署与验证
执行Maven打包:
cd ~/eclipse-workspace/HDFSClient/target/
上传至集群节点:
scp HDFSClient-jar-with-dependencies.jar root@<EIP>:/root/
远程登录并执行验证:
ssh root@<EIP>
yarn jar HDFSClient-jar-with-dependencies.jar com.huawei.hdfs.PathChecker
MapReduce计算框架开发
MapReduce采用分治思想处理海量数据,将计算任务拆分为Map阶段的数据转换与Reduce阶段的结果聚合。
词频统计实现
新建Maven项目MRJob,pom.xml关键依赖:
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-client</artifactId>
<version>3.1.1</version>
<scope>provided</scope>
</dependency>
创建com.huawei.mr包,实现WordFrequency类:
package com.huawei.mr;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import java.io.IOException;
public class WordFrequency {
public static class TokenizerMapper extends Mapper<LongWritable, Text, Text, LongWritable> {
private final static LongWritable one = new LongWritable(1);
private Text word = new Text();
@Override
protected void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
String[] tokens = value.toString().split("\\s+");
for (String token : tokens) {
word.set(token);
context.write(word, one);
}
}
}
public static class SumReducer extends Reducer<Text, LongWritable, Text, LongWritable> {
@Override
protected void reduce(Text key, Iterable<LongWritable> values, Context context)
throws IOException, InterruptedException {
long sum = 0;
for (LongWritable val : values) {
sum += val.get();
}
context.write(key, new LongWritable(sum));
}
}
public static void main(String[] args) throws Exception {
Configuration conf = new Configuration();
Job job = Job.getInstance(conf, "word frequency");
job.setJarByClass(WordFrequency.class);
job.setMapperClass(TokenizerMapper.class);
job.setCombinerClass(SumReducer.class);
job.setReducerClass(SumReducer.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(LongWritable.class);
FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}
二次排序机制
实现自定义排序需构建复合键类。创建CompositeSortKey:
package com.huawei.mr.sort;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.io.WritableComparable;
import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;
public class CompositeSortKey implements WritableComparable<CompositeSortKey> {
private Text primary = new Text();
private IntWritable secondary = new IntWritable();
public CompositeSortKey() {}
public CompositeSortKey(String first, int second) {
this.primary.set(first);
this.secondary.set(second);
}
@Override
public void write(DataOutput out) throws IOException {
primary.write(out);
secondary.write(out);
}
@Override
public void readFields(DataInput in) throws IOException {
primary.readFields(in);
secondary.readFields(in);
}
@Override
public int compareTo(CompositeSortKey other) {
int cmp = this.primary.compareTo(other.primary);
if (cmp != 0) return cmp;
return this.secondary.compareTo(other.secondary);
}
public Text getPrimary() { return primary; }
public IntWritable getSecondary() { return secondary; }
}
自定义分组比较器PrimaryKeyGroup:
package com.huawei.mr.sort;
import org.apache.hadoop.io.WritableComparable;
import org.apache.hadoop.io.WritableComparator;
public class PrimaryKeyGroup extends WritableComparator {
protected PrimaryKeyGroup() {
super(CompositeSortKey.class, true);
}
@Override
public int compare(WritableComparable a, WritableComparable b) {
CompositeSortKey k1 = (CompositeSortKey) a;
CompositeSortKey k2 = (CompositeSortKey) b;
return k1.getPrimary().compareTo(k2.getPrimary());
}
}
自定义分区器CategoryPartitioner:
package com.huawei.mr.sort;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.mapreduce.Partitioner;
public class CategoryPartitioner extends Partitioner<CompositeSortKey, IntWritable> {
@Override
public int getPartition(CompositeSortKey key, IntWritable value, int numPartitions) {
String category = key.getPrimary().toString();
if (category.startsWith("A")) return 0;
if (category.startsWith("B")) return 1;
return 2;
}
}
主作业类SecondarySortJob:
package com.huawei.mr.sort;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.conf.Configured;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.input.KeyValueTextInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import org.apache.hadoop.util.Tool;
import org.apache.hadoop.util.ToolRunner;
import java.io.IOException;
public class SecondarySortJob extends Configured implements Tool {
public static class SortMapper extends Mapper<Text, Text, CompositeSortKey, IntWritable> {
private CompositeSortKey outKey = new CompositeSortKey();
private IntWritable outValue = new IntWritable();
@Override
protected void map(Text key, Text value, Context context)
throws IOException, InterruptedException {
outKey = new CompositeSortKey(key.toString(), Integer.parseInt(value.toString()));
outValue.set(Integer.parseInt(value.toString()));
context.write(outKey, outValue);
}
}
public static class SortReducer extends Reducer<CompositeSortKey, IntWritable, Text, Text> {
@Override
protected void reduce(CompositeSortKey key, Iterable<IntWritable> values, Context context)
throws IOException, InterruptedException {
for (IntWritable val : values) {
context.write(key.getPrimary(), new Text(val.toString()));
}
}
}
@Override
public int run(String[] args) throws Exception {
Configuration conf = getConf();
conf.set("mapreduce.input.keyvaluelinerecordreader.key.value.separator", "\t");
Job job = Job.getInstance(conf, "secondary sort");
job.setJarByClass(SecondarySortJob.class);
job.setMapperClass(SortMapper.class);
job.setMapOutputKeyClass(CompositeSortKey.class);
job.setMapOutputValueClass(IntWritable.class);
job.setPartitionerClass(CategoryPartitioner.class);
job.setGroupingComparatorClass(PrimaryKeyGroup.class);
job.setNumReduceTasks(3);
job.setReducerClass(SortReducer.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(Text.class);
job.setInputFormatClass(KeyValueTextInputFormat.class);
FileInputFormat.addInputPath(job, new Path(args[0]));
FileSystem fs = FileSystem.get(conf);
Path output = new Path(args[1]);
if (fs.exists(output)) fs.delete(output, true);
FileOutputFormat.setOutputPath(job, output);
return job.waitForCompletion(true) ? 0 : 1;
}
public static void main(String[] args) throws Exception {
ToolRunner.run(new SecondarySortJob(), args);
}
}
作业提交与验证
打包上传后,准备测试数据:
vim word_input.txt
# 输入内容(Tab分隔)
hello world
hello hadoop
hadoop spark
hdfs dfs -mkdir -p /user/demo/input
hdfs dfs -put word_input.txt /user/demo/input/
执行词频统计:
yarn jar MRJob-jar-with-dependencies.jar com.huawei.mr.WordFrequency /user/demo/input /user/demo/output
查看结果:
hdfs dfs -cat /user/demo/output/part-r-00000
执行二次排序作业:
yarn jar MRJob-jar-with-dependencies.jar com.huawei.mr.sort.SecondarySortJob /user/demo/sort_input /user/demo/sort_output