【大数据】MapReduce 编程--Relation--祖孙辈关系
MapReduce 模型概述
MapReduce 模型可以分为两个主要阶段:
-
Map 阶段:负责将输入数据分解并转换为键值对(
key-value),是数据处理的第一个环节。每一行输入数据都会调用map方法进行处理。 -
Reduce 阶段:负责将 Map 阶段产生的结果进行合并、汇总或过滤,最终输出到结果存储。
hadoop fs -mkdir -p /relation/input
在 HDFS 上创建一个目录/relation/input
-p:如果上级目录不存在则一并创建,防止报错
hadoop fs -put relation.dat /relation/input
将当前目录下的 relation.dat 文件上传到 HDFS 中指定的目录
Map过程
Hadoop MapReduce 中 Mapper 类的一个框架代码
public class MyMap extends Mapper<LongWritable, Text, Text, LongWritable>{
//重写map这个方法
//mapreduce框架每读一行数据就调用一次该方法
protected void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
//具体业务逻辑就写在这个方法体中,而且我们业务要处理的数据已经被框架传递进来,在方法的参数中key-value
//key是这一行数据的起始偏移量,value是这一行的文本内容
}
}
继承自 Mapper<KEYIN, VALUEIN, KEYOUT, VALUEOUT> 泛型类。
导入 Text 类,Hadoop 中的数据类型是 Writable 类型,而不是 JDK 原生的数据类型。Text 是一个 Writable 类型的类,类似于 Java 中的 String 类型
导入 LongWritable 类,它是 Hadoop 中封装的长整型数据类型。Hadoop 不使用 Java 自带的 Long 类型,因为它需要自己定义序列化和反序列化的方式来高效处理大数据
import java.io.IOException;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.mapreduce.Mapper;
public class MyMap extends Mapper<LongWritable, Text, Text, Text> {
//重写map这个方法
//mapreduce框架每读一行数据就调用一次该方法
public void map(LongWritable key, Text value, Context context)
throws IOException, InterruptedException {
//具体业务逻辑就写在这个方法体中,而且我们业务要处理的数据已经被框架传递进来,在方法的参数中key-value
//key是这一行数据的起始偏移量,value是这一行的文本内容
//1:切分名字,用空格隔开,前面的是孩子,后面的是父母
String child = value.toString().split(" ")[0];
String parent = value.toString().split(" ")[1];
//2:产生正序与逆序的key-value同时压入context
context.write(new Text(child), new Text("-" + parent));
context.write(new Text(parent), new Text("+" + child));
}
}
-
LongWritable:这是输入数据中key的类型。在 MapReduce 中,输入的key通常是行的起始偏移量(字节位置)。LongWritable是 Hadoop 封装的一个类,表示长整型数据。
A B
C D
E F每一行的字节位置分别为:
第一行:偏移量
0第二行:偏移量
4(假设每个字符占一个字节,且每行以换行符结束)第三行:偏移量
8
-
Text:这是输入数据中value的类型,表示当前行的数据内容。在这种情况下,每一行文本会被处理为一个Text类型。-----当map方法处理第一行时,value将是"A B",即当前行的文本内容 -
Text:这是输出的key类型,表示 Mapper 输出的键。我们将数据转换成Text类型,因为输出的是字符串形式的数据。 -
Text:这是输出的value类型,表示 Mapper 输出的值。和key一样,它也是字符串类型的数据
map 方法
当前处理行数据的起始偏移量key:这是当前处理行数据的起始偏移量,通常是文件中每一行的字节位置。
当前处理的行文本内容value:这是当前处理的行文本内容,类型是 Text。框架将每行数据的内容传递给 map 方法进行处理。
context:这是一个 Context 对象,它将被用来把处理后的结果输出到框架中,从而进行进一步的处理。你会通过 context.write() 方法将 key-value 对写入到输出中。
value.toString():value 是一个 Text 对象,表示当前行的文本。toString() 将 Text 对象转换成 Java 的 String 类型
split(" "):将当前行的文本按空格进行分割。假设每一行数据都是“孩子 父母”的格式,这一行将被分割成一个包含两个元素的数组,第一个元素是孩子的名字,第二个元素是父母的名字
context.write(new Text(child), new Text("-" + parent));
context.write(new Text(parent), new Text("+" + child));
为了表示亲子关系的方向,在 parent 前面加上一个负号 "-",表示“孩子 → 父母”的关系
在 child 前面加上一个加号 "+",表示“父母 → 孩子”的关系
| 对于输入: Tom Lucy Tom Jack Jone Lucy Jone Jack |
Mapper 输出: Tom -Lucy |
根据每一行的数据,将其拆分为孩子和父母,并产生两对键值对:一个是“孩子 → 父母”的正序关系,另一个是“父母 → 孩子”的逆序关系
context.write(key, value) 会将一个 key-value 对传递给 Hadoop 框架
把
context想象成一个 信箱,你通过context.write()把一封信(key-value)放进去
这些 key-value 对会被暂时存储,通常会传递到 Shuffle 阶段(在 Map 和 Reduce 之间)
在 Shuffle 阶段,框架会按照 key 对这些数据进行分组,然后再传递给 Reducer,使得同一个 key 的所有值都聚集在一起,便于进行进一步的处理
Reduce 阶段接收到的会是类似这样的:
| key | values |
|---|---|
| B | -A, +C |
| C | -B, +D |
| A | +B |
| … | … |
当一个人既是别人的孩子、又是别人的父母时,找出他们的父母(即祖父母)与他们的子女(即孙子孙女)之间的配对
Reduce过程
public class MyReduce extends Reducer<Text, LongWritable, Text, LongWritable>{
//继承Reducer之后重写reduce方法
//第一个参数是key,第二个参数是集合。
//框架在map处理完成之后,将所有key-value对缓存起来,进行分组,然后传递一个组<key,valus{}>,调用一次reduce方法
protected void reduce(Text key, Iterable<LongWritable> values,Context context)
throws IOException, InterruptedException {
}
}
-
key: 是一个中间节点,它既可能是别人的父亲也可能是别人的儿子。 -
values: 是所有相关的 "+child" 或 "-parent" 形式的 value
每个不同的 key会调用一次 reduce 方法
ArrayList<Text> grandparent = new ArrayList<Text>();
ArrayList<Text> grandchild = new ArrayList<Text>();
保存 key 的父母(即 -X 中的 X) ;保存 key 的孩子(即 +Y 中的 Y)
for (Text t : values) {
String s = t.toString();
if (s.startsWith("-")) {
grandparent.add(new Text(s.substring(1)));
} else {
grandchild.add(new Text(s.substring(1)));
}
}
startsWith() 方法 ---------判断一个字符串是否以指定前缀开头str.startsWith("前缀字符串")
substring() 方法---------截取字符串中的一部分str.substring(开始索引, 结束索引);索引从 0 开始
遍历所有的 values----如果是 "-X" 开头,表示 key 是 X 的子女 ⇒ X 是父母,所以放到 grandparent 列表中
substring(1)去掉前缀 - 或 +,只留下人名
key = B,values = [-C, +A, +D]
grandparent = [C]
grandchild = [A, D]
for (int i = 0; i < grandchild.size(); i++) {
for (int j = 0; j < grandparent.size(); j++) {
context.write(new Text(grandchild.get(i) + " "), grandparent.get(j));
}
}
把所有孙子和祖父组合一遍(双重循环组合)
context.write(key, value) 是 Hadoop 提供的方法,用来把处理结果输出
A C
D C
1 收到一个 key 和多个 values(来自 Mapper) 2 遍历 values,将 -号开头的加入父母列表,将+号开头的加入孩子列表3 遍历两个列表组合出所有“孙子 → 祖父”的对应关系 4 使用 context.write()输出每一对结果
import java.io.IOException;
import java.util.ArrayList;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
/***
*
* 1:reduce的四个参数,第一个key-value是map的输出作为reduce的输入,第二个key-value是输出祖父母和孩子的关系,所以
* 是Text,Text的格式;
*/
public class MyReduce extends Reducer<Text, Text, Text, Text> {
//继承Reducer之后重写reduce方法
//第一个参数是key,第二个参数是集合。
//框架在map处理完成之后,将所有key-value对缓存起来,进行分组,然后传递一个组<key,valus{}>,调用一次reduce方法
public void reduce(Text key, Iterable<Text> values, Context context)
throws IOException, InterruptedException {
//1:创建两个List保存祖父母和孩子
ArrayList<Text> grandparent = new ArrayList<Text>();
ArrayList<Text> grandchild = new ArrayList<Text>();
//2:对各个values中的值进行处理,key的父母保存到祖父母列表中,key的子女保存在孩子孩子列表中
for (Text t : values) {
String s = t.toString();
if (s.startsWith("-")) {
grandparent.add(new Text(s.substring(1)));
} else {
grandchild.add(new Text(s.substring(1)));
}
}
//3:再将grandparent与grandchild中的东西,一一对应输出。
for (int i = 0; i < grandchild.size(); i++) {
for (int j = 0; j < grandparent.size(); j++) {
context.write(new Text(grandchild.get(i) + " "), grandparent.get(j));
}
}
}
}
执行MapReduce任务
public class MyRunner {
public static void main(String[] args) throws IOException, ClassNotFoundException, InterruptedException {
//创建配置文件
Configuration conf = new Configuration();
//获取一个作业
Job job = Job.getInstance(conf);
//设置整个job所用的那些类在哪个jar包
job.setJarByClass(MyRunner.class);
//本job使用的mapper和reducer的类
job.setMapperClass(MyMap.class);
job.setReducerClass(MyReduce.class);
//指定reduce的输出数据key-value类型
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(LongWritable.class);
//指定mapper的输出数据key-value类型
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(LongWritable.class);
//指定要处理的输入数据存放路径
FileInputFormat.setInputPaths(job, new Path("hdfs://master:9000/user/cg/input"));
//指定处理结果的输出数据存放路径
FileOutputFormat.setOutputPath(job, new Path("hdfs://master:9000/user/cg/output"));
//将job提交给集群运行
job.waitForCompletion(true);
}
}
-
Configuration:Hadoop 作业的配置对象。 -
Job:代表一个 MapReduce 作业。 -
Path:Hadoop 的路径类。 -
Text:Hadoop 封装的字符串类(用于 key/value)。 -
FileInputFormat和FileOutputFormat:分别指定输入/输出的路径。 -
Scanner:读取控制台用户输入。 -
FileSystem和FSDataInputStream:用于从 HDFS 中读取结果文件。 -
URI:定义 HDFS 地址。
Configuration conf = new Configuration();
Job job = Job.getInstance(conf);
conf 用于加载配置;job 是整个 MapReduce 作业的核心对象
Job就是 Hadoop MapReduce 程序的一次完整 任务。它是 Hadoop 的一个抽象,用来表示你要执行的“计算过程”,包括了以下几个核心要素:
输入数据的路径(你要处理的数据在哪里)
输出数据的路径(你处理完的数据要存放在哪)
Mapper 和 Reducer 的类(这些类告诉 Hadoop 怎么处理数据)
数据类型的定义(你处理的数据是什么样的,输出数据是什么样的)
可以把
Job理解为一个 “任务配置文件”,描述了希望在 Hadoop 上执行的具体操作每个
Job代表一次从头到尾的计算任务Hadoop 集群接收到这个
Job后,它会:
把输入数据按分片(切分)进行分发给多个
Mapper节点;每个
Mapper节点会按照你的逻辑(Mapper类的map方法)进行处理;处理完成后,
Reducer节点会汇总数据并输出结果。执行完后,会得到一个输出文件(例如
part-r-00000),可以查看这个文件中的计算结果
job.setJarByClass(MyRunner.class); // 指定运行类所在的jar包
job.setMapperClass(MyMap.class); // Mapper 类
job.setReducerClass(MyReduce.class); // Reducer 类
job.setJarByClass(MyRunner.class)告诉 Hadoop 这个程序的主类(包含 main() 的类)是哪个,这样它能把这个类所在的 .jar 文件分发到集群中各个节点去运行
设置 Mapper 和 Reducer 的输出 key-value 类型
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(Text.class);
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(Text.class);
inputPath 和 outputPath 是控制台输入,表示 HDFS 中输入文件夹和输出目录
Scanner sc = new Scanner(System.in);
System.out.print("inputPath:");
String inputPath = sc.next();
System.out.print("outputPath:");
String outputPath = sc.next();
拼接完整的 HDFS 路径:
hdfs://master:9000/input 和 hdfs://master:9000/output
提交作业并等待执行完成job.waitForCompletion(true);
读取结果文件内容并输出到控制台
FileSystem fs = FileSystem.get(new URI("hdfs://master:9000"), new Configuration());
Path srcPath = new Path(outputPath+"/part-r-00000");
FSDataInputStream is = fs.open(srcPath);
part-r-00000 是默认的 Reducer 输出文件
is.readLine():按行读取输出文件内容
import java.io.IOException;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
import java.util.Scanner;
import org.apache.hadoop.fs.FSDataInputStream;
import org.apache.hadoop.fs.FileSystem;
import java.net.URI;
public class MyRunner {
public static void main(String[] args) throws IOException, ClassNotFoundException, InterruptedException {
//创建配置文件
Configuration conf = new Configuration();
//获取一个作业
Job job = Job.getInstance(conf);
//设置整个job所用的那些类在哪个jar包
job.setJarByClass(MyRunner.class);
//本job使用的mapper和reducer的类
job.setMapperClass(MyMap.class);
job.setReducerClass(MyReduce.class);
//指定reduce的输出数据key-value类型
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(Text.class);
//指定mapper的输出数据key-value类型
job.setMapOutputKeyClass(Text.class);
job.setMapOutputValueClass(Text.class);
Scanner sc = new Scanner(System.in);
System.out.print("inputPath:");
String inputPath = sc.next();
System.out.print("outputPath:");
String outputPath = sc.next();
//指定要处理的输入数据存放路径
FileInputFormat.setInputPaths(job, new Path("hdfs://master:9000"+inputPath));
//指定处理结果的输出数据存放路径
FileOutputFormat.setOutputPath(job, new Path("hdfs://master:9000"+outputPath));
//将job提交给集群运行
job.waitForCompletion(true);
try {
FileSystem fs = FileSystem.get(new URI("hdfs://master:9000"), new Configuration());
Path srcPath = new Path(outputPath+"/part-r-00000");
FSDataInputStream is = fs.open(srcPath);
System.out.println("Results:");
while(true) {
String line = is.readLine();
if(line == null) {
break;
}
System.out.println(line);
}
is.close();
}catch(Exception e) {
e.printStackTrace();
}
}
}
更多推荐
所有评论(0)