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
Lucy  +Tom
Tom   -Jack
Jack  +Tom
Jone  -Lucy
Lucy  +Jone
Jone  -Jack
Jack  +Jone

根据每一行的数据,将其拆分为孩子和父母,并产生两对键值对:一个是“孩子 → 父母”的正序关系,另一个是“父母 → 孩子”的逆序关系

context.write(key, value) 会将一个 key-value 对传递给 Hadoop 框架

把 context 想象成一个 信箱,你通过 context.write() 把一封信(key-value)放进去

这些 key-value 对会被暂时存储,通常会传递到 Shuffle 阶段(在 Map 和 Reduce 之间)

在 Shuffle 阶段,框架会按照 key 对这些数据进行分组,然后再传递给 Reducer,使得同一个 key 的所有值都聚集在一起,便于进行进一步的处理


Reduce 阶段接收到的会是类似这样的:

keyvalues
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();
        }
    } 
}

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐