• linkedu视频
  • 平面设计
  • 电脑入门
  • 操作系统
  • 办公应用
  • 电脑硬件
  • 动画设计
  • 3D设计
  • 网页设计
  • CAD设计
  • 影音处理
  • 数据库
  • 程序设计
  • 认证考试
  • 信息管理
  • 信息安全
菜单
linkedu.com
导航菜单
  • 网页制作
  • 数据库
  • 程序设计
  • 操作系统
  • CMS教程
  • 游戏攻略
  • 脚本语言
  • 平面设计
  • 软件教程
  • 网络安全
  • 电脑知识
  • 服务器
  • 视频教程
  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx
您的位置:首页 > 服务器 >云计算 > MapReduce的两表join操作优化,mapreduce表join

MapReduce的两表join操作优化,mapreduce表join

作者:网友 字体:[增加 减小] 来源:互联网

本文主要包含mapreduce join,mapreduce优化,node mapreduce,mapreduce是什么,mapreduce等服务器相关知识,网友希望可以进行参考

MapReduce的两表join操作优化,mapreduce表join


注:优化前的分析过程详见本博的上篇博文

案例

地址(Address)和人员(Person)的一对多关联

 

原始数据

地址(Address)数据

id AddreName

1 beijing
2 shanghai
3 guangzhou

人员(Person)数据

1 zhangsan 1
2 lisi 2
3 wangwu 1
4 zhaoliu 3
5 maqi 3

 

优化前,我们通过构造一个通用的JavaBean来存储两张表的属性。但是我们发现最后reduce时用List数组来存储和Address地址表区分开来的Person表数组,将造成大量的内存开销,所有我们想到重新构造Map的key类型,在数据进行reduce前重新group分组操作

优化前代码分析

1.自定义JavaBean代码

/*
 * 人员和地址的通用bean
 */
public class Bean implements WritableComparable<Bean> {
	private String userNo = "";
	private String userName = "";
	private String addreNo = "";
	private String addreName = "";
	private int flag;

	public Bean(Bean bean) {
		this.userName = bean.getUserName();
		this.userNo = bean.getUserNo();
		this.addreName = bean.getAddreName();
		this.addreNo = bean.getAddreNo();
		this.flag = bean.getFlag();
	}

	public Bean() {
		super();
		// TODO Auto-generated constructor stub
	}

	public Bean(String userNo, String userName, String addreNo,
			String addreName, int flag) {
		super();
		this.userNo = userNo;
		this.userName = userName;
		this.addreNo = addreNo;
		this.addreName = addreName;
		this.flag = flag;
	}

	public String getUserNo() {
		return userNo;
	}

	public void setUserNo(String userNo) {
		this.userNo = userNo;
	}

	public String getUserName() {
		return userName;
	}

	public void setUserName(String userName) {
		this.userName = userName;
	}

	public String getAddreNo() {
		return addreNo;
	}

	public void setAddreNo(String addreNo) {
		this.addreNo = addreNo;
	}

	public String getAddreName() {
		return addreName;
	}

	public void setAddreName(String addreName) {
		this.addreName = addreName;
	}

	public int getFlag() {
		return flag;
	}

	public void setFlag(int flag) {
		this.flag = flag;
	}

	@Override
	public void write(DataOutput out) throws IOException {
		out.writeUTF(userNo);
		out.writeUTF(userName);
		out.writeUTF(addreNo);
		out.writeUTF(addreName);
		out.writeInt(flag);

	}

	@Override
	public void readFields(DataInput in) throws IOException {
		this.userNo = in.readUTF();
		this.userName = in.readUTF();
		this.addreNo = in.readUTF();
		this.addreName = in.readUTF();
		this.flag = in.readInt();

	}

	@Override
	public int compareTo(Bean arg0) {
		// TODO Auto-generated method stub
		return 0;
	}

	@Override
	public String toString() {
		return "userNo=" + userNo + ", userName=" + userName + ", addreNo="
				+ addreNo + ", addreName=" + addreName;
	}

}

2.Map操作

public class PersonAddrMap extends
		Mapper<LongWritable, Text, IntWritable, Bean> {
	@Override
	protected void map(LongWritable key, Text value,
			Mapper<LongWritable, Text, IntWritable, Bean>.Context context)
			throws IOException, InterruptedException {
		String line = value.toString();
		String str[] = line.split("\t");
		if (str.length == 2) { // 地区信息表
			Bean bean = new Bean();
			bean.setAddreNo(str[0]);
			bean.setAddreName(str[1]);
			bean.setFlag(0); // 0表示地区
			context.write(new IntWritable(Integer.parseInt(str[0])), bean);
		} else {// 人员信息表
			Bean bean = new Bean();
			bean.setUserNo(str[0]);
			bean.setUserName(str[1]);
			bean.setAddreNo(str[2]);
			bean.setFlag(1); // 1表示人员表
			context.write(new IntWritable(Integer.parseInt(str[2])), bean);
		}
	}
}


3.reduce操作

public class PersonAddreRedu extends
		Reducer<IntWritable, Bean, NullWritable, Text> {
	@Override
	protected void reduce(IntWritable key, Iterable<Bean> values,
			Reducer<IntWritable, Bean, NullWritable, Text>.Context context)
			throws IOException, InterruptedException {
		Bean Addre = null;
		List<Bean> peoples = new ArrayList<Bean>();
		/*
		 * 如果values的第一个元素信息就是地址Addre的信息的话,
		 * 我们就不再需要一个List来缓存person信息了,values后面的全是人员信息
		 * 将减少巨大的内存空间
		 */
		/*
		 * partitioner和shuffer的过程:
		 * partitioner的主要功能是根据reduce的数量将map输出的结果进行分块,将数据送入到相应的reducer.
		 * 所有的partitioner都必须实现partitioner接口并实现getPartition方法,该方法的返回值为int类型,并且取值范围在0~(numOfReducer-1),
		 * 从而能将map的输出输入到对应的reducer中,对于某个mapreduce过程,hadoop框架定义了默认的partitioner为HashPartioner,
		 * 该partitioner使用key的hashCode来决定将该key输送到哪个reducer;
		 * shuffle将每个partitioner输出的结果根据key进行group以及排序,将具有相同key的value构成一个values的迭代器,并根据key进行排序分别调用
		 * 开发者定义的reduce方法进行排序,因此mapreducer的所以key必须实现comparable接口的compareto()方法从而能实现两个key对象的比较
		 */
		/*
		 * 我们需要自定义key的数据结构(shuffle按照key进行分组)来满足共同addreNo的情况下地址表的更小需求
		 * 
		 */
		for (Bean bean : values) {
			if (bean.getFlag() == 0) { // 表示地区表
				Addre = new Bean(bean);

			} else {
				peoples.add(new Bean(bean)); // 添加到peoplelist中
			}
		}
		for (Bean peo : peoples) { // 给peoplelist添加地区名字
			peo.setAddreName(Addre.getAddreName());
			context.write(NullWritable.get(), new Text(peo.toString()));
		}
	}
}

4.job操作

public class PersonAddreMain {
	public static void main(String[] args) throws Exception {
		Configuration conf = new Configuration();
		Job job = new Job(conf);
		job.setJarByClass(PersonAddreMain.class);

		job.setMapperClass(PersonAddrMap.class);
		job.setMapOutputKeyClass(IntWritable.class);
		job.setMapOutputValueClass(Bean.class);

		job.setReducerClass(PersonAddreRedu.class);
		job.setOutputKeyClass(NullWritable.class);
		job.setOutputValueClass(Text.class);

		FileInputFormat.addInputPath(job, new Path(args[0]));
		FileOutputFormat.setOutputPath(job, new Path(args[1]));
		job.waitForCompletion(true);
	}
}


 

具体优化分析

在Reduce类的reduce()方法中如果values的第一个元素信息就是地址Addre的信息的话,我们就不再需要一个List来缓存person信息了,values后面的全是人员信息将减少巨大的内存空间。

 

partitioner和shuffer的过程
partitioner的主要功能是根据reduce的数量将map输出的结果进行分块,将数据送入到相应的reducer.
所有的partitioner都必须实现partitioner接口并实现getPartition方法,该方法的返回值为int类型,并且取值范围在0~(numOfReducer-1),
从而能将map的输出输入到对应的reducer中,对于某个mapreduce过程,hadoop框架定义了默认的partitioner为HashPartioner,
该partitioner使用key的hashCode来决定将该key输送到哪个reducer;

 shuffle将每个partitioner输出的结果根据key进行group以及排序,将具有相同key的value构成一个values的迭代器,并根据key进行排序分别调用
 开发者定义的reduce方法进行排序,因此mapreducer的所以key必须实现comparable接口的compareto()方法从而能实现两个key对象的比较

我们需要自定义key的数据结构(shuffle按照key进行分组)来满足共同addreNo的情况下地址表的更小需求

 

优化后

1.JavaBean操作

/*
 * 人员和地址的通用bean
 * 用作map输出的value
 */
public class Bean implements WritableComparable<Bean> {
	private String userNo = " ";
	private String userName = " ";
	private String addreNo = " ";
	private String addreName = " ";

	public Bean(Bean bean)
  


 
分享到:QQ空间新浪微博腾讯微博微信百度贴吧QQ好友复制网址打印

您可能想查找下面的文章:

  • MapReduce的两表join操作优化,mapreduce表join
  • MapReduce的两表join一般操作,mapreduce表join

相关文章

  • Cloud Foundry service broker开发部署实例解析(上),foundrybroker
  • Hive 外部表 分区表,hive外部表分区表
  • Hadoop之——HDFS随笔,hadoophdfs
  • Linux、hive、sqoop常用脚本,hivesqoop
  • Spark SQL and DataFrame Guide(1.4.1)——之DataFrames,sparkdataframe
  • 可穿戴KEY带来的身份认证的革命,key身份认证
  • Dynamics CRM 2015 Update 1 系列(2): Upsert API,dynamicsupsert
  • Spark Streaming和Flume集成指南V1.4.1,flumev1.4.1
  • CRM市场及Zoho模式,CRM市场Zoho模式
  • Libvirt中windows虚拟机的动态内存管理,libvirt虚拟机

文章分类

  • windows
  • 服务器硬件
  • 服务器运维
  • 云计算
  • 虚拟化
  • IIS教程
  • Linux
  • Apache
  • Ftp
  • DNS
  • Nginx

最近更新的内容

    • 云服务器之间实时文件同步和文件备份的最简单高效的免费方案,实时文件备份
    • 【Spark1.3官方翻译】 Spark Submit提交应用程序,spark1.3spark
    • 学习札记:CISCO云计算,札记cisco云计算
    • 【Spark1.3官方翻译】Spark集群模式概览,spark1.3spark
    • Spark中广播变量知识点
    • HBase Shell的基本用法,HBaseShell用法
    • apache hadoop-2.6.0-CDH5.4.1 安装笔记,hadoop2.6.0安装
    • ElasticSearch的分布式安装
    • 从事Cloud行业需要掌握的基本技能清单,cloud基本技能
    • hive grouping sets 和 cube 用法,hivegrouping

关于我们 - 联系我们 - 免责声明 - 网站地图

©2020-2025 All Rights Reserved. linkedu.com 版权所有