当前位置:首页 > 技术 > 正文内容

华为云MRS集群HDFS与MapReduce开发实践

访客 技术 2026年9月28日 12

环境准备

弹性公网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项目:

  1. File → New → Project → Maven Project
  2. 勾选"Create a simple project"选项
  3. 填写工程元数据:
    • 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

相关文章

Linux crontab 详解

1) crontab 是什么cron 是 Linux 的定时任务守护进程;crontab 是用来编辑/查看“按时间周期执行命令”的表(cron table)。常见两类:用户 crontab:每个用户一份(crontab -e 编辑)系统级 crontab / cron.d:可指定执行用户(/etc/crontab、/etc/cron.d/*)2) crontab 时间...

富文本里可以允许的 HTML 属性

一、所有标签默认允许的安全属性(极少)class        (可选)id           (通常建议禁用)title️ 注意:id 容易被滥用做锚点注入,很多系统直接禁用class 允许的话最好只允许固定前缀(如 editor-*)二、a 标签允许属性<a href="" t...

Mac 安装 Node.js 指南

方法一:通过官网安装包(最简单,适合初学者)如果你只是想快速安装并开始使用,这是最直接的方法。访问 Node.js 官网。页面会显示两个版本:LTS (Recommended For Most Users):长期支持版,最稳定。建议选这个。Current:最新特性版,包含最新功能但可能不够稳定。下载 .pkg 安装包并运行。按照安装向导点击“下一步”即可完成。方法二:使用 Homebrew 安装(...

Dom\HTML_NO_DEFAULT_NS 的副作用:自动加闭合标签

在使用Dom\HTMLDocument时,Dom\HTML_NO_DEFAULT_NS 将禁止在解析过程中设置元素的命名空间, 此设置是为了与DOMDocument向后兼容而存在的。当使用它时,已知的一个副作用就是:自动加闭合标签例如 </img> 为什么会这样?当你使用:Dom\HTML_NO_DEFAULT_NS文档会变成 无命名空间模式,此时内部更接近 XML...

Laravel 事件和监听器创建

在 Laravel 中,使用 Artisan 命令创建 Events(事件) 和 Listeners(监听器) 是非常高效的。你可以通过以下几种方式来实现:1. 手动创建单个 Event如果你只想创建一个事件类,可以使用 make:event 命令:Bashphp artisan make:event UserRegistered执行后,文件将生成在 app/Even...

自定义域名解析神器 dnsmasq

什么是 dnsmasq?dnsmasq 是一个轻量级、功能强大的网络服务工具,专为小型和中等规模网络设计。它是一个综合的网络基础设施解决方案[1]。dnsmasq 能做什么?功能说明应用场景DNS 转发与缓存将 DNS 查询转发到上游服务器(ISP、Google DNS 等),并在本地缓存结果加快 DNS 查询速度,减少外部 DNS 流量本地 DNS解析本地网络设备的主机名,无需编辑&n...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。