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

hadoop 多数据源连接之DataJoin,hadoopdatajoin

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

本文主要包含hadoop源代码,hadoop源代码下载,hadoop数据分析,hadoop数据分析平台,hadoop数据挖掘等服务器相关知识,网友希望可以进行参考

hadoop 多数据源连接之DataJoin,hadoopdatajoin


一个MapReduce任务很可能访问和处理两个甚至多个数据集,在关系型数据库中,这将是两个或者多个表的连接,但是Hadoop系统没有关系型数据库中那样强大的连接处理功能,因此处理复杂一些。一般来讲,hadoop可以采用这几种数据连接方式:

        1采用DataJoin类库实现Reduce端连接的方法

        2 用全局文件复制实现Map端连接方法

        3 带Map端过滤的Reduce端连接方法

   Hadoop的Mapreduce框架提供了一种较为通用  的多数据源连接方法,该方法用DataJoin类库为程序员提供了完成数据连接所需的编程框架和接口,其处理方法如下:

         为了完成不用数据源的连接操作,我们必须给每个数据源制定一个标签(tag),用来区分数据,就像关系型数据库中表名一样,这里我们需要实现 Text generateInputTag(String inputFile)方法;

         另外,为了进行连接操作,我们必须知道连接的主键是什么,类似于关系型数据库中的key,因此我们需要指定groupKey,这里我们需要实现 Text generateGroupKey(TaggedMapOutput aRecord)

         然后在Map端我们需要把原始数据包装成为一个带标签的数据记录,方便shuffle和Reduce端执行笛卡尔积,所以我们需要实现 TaggedMapOutput generateTaggedMapOutput(Object value);


总结一下Map处理过程:

   Datajoin类库首先提哦功能管理一个抽象基类DataJoinMapperBase,该基类实现了map()方法,帮助程序员对每个数据源下的记录生成一个代标签的数据记录对象。Map端处理过程中,需要指定标签tag和Groupkey,然后包装成为带标签的数据记录对象,在shuffle过程中,这些GroupKey相同的记录被分到同一个Reduce节点上。

Reduce处理过程:

      Reduce节点收到这些带标签的数据记录后,Reduce过程将这些带不同的数据源标签的记录执行笛卡尔积,自动生成所有不同的叉积组合,由程序员实现一个combine()方法,根据应用程序需求将这些具有相同的Groupkey的数据记录进行适当的合并处理,以此完成类似于关系型数据库中不同实体数据记录之间的连接。

     在Reduce阶段我们需要继承DataJoinReduceBase,该基类实现了reduce()方法,我们只是需要实现combine()方法即可,另外我们还是需要继承TaggedMapOutput类,它描述了一个标签化的数据记录,实现了getTag(),setTag()方法,作为Mapper的key_value输出value类型,由于需要I/O,我们需要继承并且实现Writable接口,并且实现getData()方法用以读取记录数据

   下面是数据源:

        user.txt文件:

1,张三,135xxxxxxxx
2,李四,136xxxxxxxx
3,王五,137xxxxxxxx
4,赵六,138xxxxxxxx

order.txt文件:

3,A,13,2013-02-12
1,B,23,2013-02-14
2,C,16,2013-02-17
3,D,25,2013-03-12

这其中需要注意很多小细节,因为没有要求程序员实现Map和reduce方法,所以我们会很容易忽略很多东西,需要注意的东西我在下面一一注释了:

我们必须使用Jobconf 来声明一个job,同时使用JobClient来run job,另外我们在继承TaggedMapOutput的时候默认的无参构造方法中需要初始化data

package joinTest;

import java.io.DataInput;
import java.io.DataOutput;
import java.io.IOException;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.contrib.utils.join.DataJoinMapperBase;
import org.apache.hadoop.contrib.utils.join.DataJoinReducerBase;
import org.apache.hadoop.contrib.utils.join.TaggedMapOutput;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.io.Writable;
import org.apache.hadoop.mapred.FileInputFormat;
import org.apache.hadoop.mapred.FileOutputFormat;
import org.apache.hadoop.mapred.JobClient;
import org.apache.hadoop.mapred.JobConf;
import org.apache.hadoop.mapred.TextInputFormat;
import org.apache.hadoop.mapred.TextOutputFormat;

public class DataJoin{

	public static class MyTaggedMapOutput extends TaggedMapOutput{
		private Writable data;
	    
		public MyTaggedMapOutput(){
			//一定要new一下,不然反序列化时会报空指针异常
			this.data=new Text();
		}
		
		public MyTaggedMapOutput(Writable data){
			this.tag=new Text("");
			this.data=data;
		}

		public void readFields(DataInput in) throws IOException {
			this.tag.readFields(in);
			this.data.readFields(in);
		}

		public void write(DataOutput output) throws IOException {
			this.tag.write(output);
			//this.tag.write(output);     //大问题,粗心写成了this.tag.write(output); 结果一直报错
			this.data.write(output);
		}

		@Override
		public Writable getData() {
			return data;
		}
	}
	
	
	public static class DataJoinMapper extends DataJoinMapperBase{
		
		protected Text generateInputTag(String inputFile){
				//String datasource=inputFile.substring(inputFile.lastIndexOf("/")+1).split("\\.")[0];
			    String datasource = inputFile.split("-")[0];
				System.out.println("datasource:"+datasource);
				return new Text(datasource);
		}
		
		protected TaggedMapOutput generateTaggedMapOutput(Object value){
			TaggedMapOutput tm=new MyTaggedMapOutput((Text)value);
			tm.setTag(this.inputTag);
			return tm;
		}

		@Override
		protected Text generateGroupKey(TaggedMapOutput aRecord) {
			String line=aRecord.getData().toString();
			//String groupkey=line.split("\\s")[0];
			String groupkey=line.split(",")[0];
			return new Text(groupkey);
		}
	}
	
	public static class DataJoinReducer extends DataJoinReducerBase{

		@Override
		protected TaggedMapOutput combine(Object[] tags, Object[] values) {
			if(tags.length<2){
				return null;
			}
			String output = "";
		   /* for(int i=0;i<tags.length;i++){
		    	TaggedMapOutput tat=(MyTaggedMapOutput)values[i];System.out.println("tags:"+tags[i]+"    values:"+tat.getData().toString());
		    	if(i==0){
		    	     output=tat.getData().toString();//System.out.println("i==0  output:"+output);
		    	}else{
		    		output+="\t";
		    		String [] s=tat.getData().toString().split("\\s",2);
		    		System.out.println("s.length:"+s.length);
		    		output+=s[0];
		    	}*/
			for(int j=0;j<tags.length;j++){
				//TaggedMapOutput taOutput=(TaggedMapOutput)tags[j];
				TaggedMapOutput taggedMapOutput=(TaggedMapOutput)values[j];
				System.out.println("tag:"+taggedMapOutput.getTag()+"  value:"+taggedMapOutput.getData().toString());
			}
			     for(int i=0;i<values.length;i++){
			    	 TaggedMapOutput tat=(MyTaggedMapOutput)values[i];
			    	 String  recordLine=((Text)tat.getData()).toString();
			    	 String [] tokens=recordLine.split(",",2);System.out.println("data:"+recordLine);
			    	 if(i>0)
			    		 output+=",";
			    	output+=tokens[1];
		    }
		    TaggedMapOutput tag=new MyTaggedMapOutput(new Text(output));
		    tag.setTag((Text)tags[0]);
		    return tag;
		}
	}
	
	/*
	 * 这里一定要注意,FileInputFormat和FileOutputFormat一定要是org.apache.hadoop.mapred下面的包
	 */
	public static int run(String args[]) throws IOException{
		       Configuration conf=new Configuration();
				//Configuration conf=getConf();
				JobConf job=new JobConf(conf, DataJoin.class);
				//Job job=new Job(conf,"DataJoin");
				job.setJobName("DataJoin");
			//	job.setJarByClass(DataJoin.class);
				
				job.setMapperClass(DataJoinMapper.class);
				job.setReducerClass(DataJoinReducer.class);
				job.setInputFormat(TextInputFormat.class);
				job.setOutputFormat(TextOutputFormat.class);
				job.setMapOutputKeyClass(Text.class);
				job.setMapOutputValueClass(MyTaggedMapOutput.class);
				job.setOutputKeyClass(Text.class);
				job.setOutputValueClass(Text.class);
				
				//job.set("mapred.textoutputformat.separator", "\t");
				job.set(&q
  


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

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

  • Hadoop源代码分析(三七),hadoop源代码
  • Hadoop 源代码分析(六)RPC-Client,hadooprpc-client
  • Hadoop 源代码分析(三)对象序列化,hadoop序列化
  • hadoop 多数据源连接之DataJoin,hadoopdatajoin

相关文章

  • 机器学习-统计学习方法概论,学习方法概论
  • 通过 JMX 获取Hadoop/HBase监控数据,jmxhadoop
  • mahout简介及安装配置
  • Spark的日志配置,Spark日志配置
  • hive之实现列转行,hive列转行
  • openstack中Nova组件images的所有python API 汇总,openstacknova
  • RegionServer功能职责,regionserver职责
  • Hadoop之——HDFS命令,hadoophdfs
  • inet_ntoa issue,inet_ntoa
  • Elasticsearch之Nested Object Mapping,elasticsearch

文章分类

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

最近更新的内容

    • Install Docker Mac OS X,installdocker
    • Linux、hive、sqoop常用脚本,hivesqoop
    • 2015 OpenCloud峰会总结,2015opencloud峰会
    • S3C2416 按键驱动,s3c2416按键驱动
    • HBase常用操作之namespace,hbasenamespace
    • Mahout的BreimanExample例子分析,mahoutexample
    • Hadoop Installation on Linux,hadoopinstallation
    • 配置KVM虚拟机的网络,Bridge和Nat方式,kvmnat
    • hadoop做HA后,hbase修改,hadoophahbase修改
    • myeclipse2013 建的web项目没有web.xml,myeclipse中web.xml

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

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