`
qindongliang1922
  • 浏览: 2189057 次
  • 性别: Icon_minigender_1
  • 来自: 北京
博客专栏
7265517b-f87e-3137-b62c-5c6e30e26109
证道Lucene4
浏览量:117677
097be4a0-491e-39c0-89ff-3456fadf8262
证道Hadoop
浏览量:126079
41c37529-f6d8-32e4-8563-3b42b2712a50
证道shell编程
浏览量:60034
43832365-bc15-3f5d-b3cd-c9161722a70c
ELK修真
浏览量:71405
社区版块
存档分类
最新评论

如何使用Hadoop的JobControl

阅读更多
使用Hadoop里面的MapReduce来处理海量数据是非常简单方便的,但有时候我们的应用程序,往往需要多个MR作业,来计算结果,比如说一个最简单的使用MR提取海量搜索日志的TopN的问题,注意,这里面,其实涉及了两个MR作业,第一个是词频统计,第两个是排序求TopN,这显然是需要两个MapReduce作业来完成的。其他的还有,比如一些数据挖掘类的作业,常常需要迭代组合好几个作业才能完成,这类作业类似于DAG类的任务,各个作业之间是具有先后,或相互依赖的关系,比如说,这一个作业的输入,依赖上一个作业的输出等等。


在Hadoop里实际上提供了,JobControl类,来组合一个具有依赖关系的作业,在新版的API里,又新增了ControlledJob类,细化了任务的分配,通过这两个类,我们就可以轻松的完成类似DAG作业的模式,这样我们就可以通过一个提交来完成原来需要提交2次的任务,大大简化了任务的繁琐度。具有依赖式的作业提交后,hadoop会根据依赖的关系,先后执行的job任务,每个任务的运行都是独立的。

下面来看下散仙的例子,组合一个词频统计+排序的作业,测试数据如下:


秦东亮;72
秦东亮;34
秦东亮;100
三劫;899
三劫;32
三劫;1
a;45
b;567
b;12

代码如下:
package com.qin.test.hadoop.jobctrol;

import java.io.IOException;

import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.io.WritableComparator;
import org.apache.hadoop.mapred.JobConf;
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.jobcontrol.ControlledJob;
import org.apache.hadoop.mapreduce.lib.jobcontrol.JobControl;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;




/**
 * Hadoop的版本是1.2的
 * JDK环境1.6
 * 使用ControlledJob+JobControl新版API
 * 完成组合式任务 
 * 第一个任务是统计词频
 * 第二个任务是降序排序
 * 
 * 如果使用MapReduce作业来完成的话,则需要跑2个MR任务
 * 但是如果我们使用了JobControl+ControlledJob就可以在
 * 一个类里面完成类型的DAG依赖式的作业
 * 
 * 
 * @author qindongliang
 * 
 * 
 * 
 * ***/
public class MyHadoopControl {
	
	
	
	/***
	 *
	 *MapReduce作业1的Mapper
	 *
	 *LongWritable 1  代表输入的key值,默认是文本的位置偏移量
	 *Text 2          每行的具体内容
	 *Text 3          输出的Key类型
	 *Text 4          输出的Value类型
	 * 
	 * */
	private static class SumMapper extends Mapper<LongWritable, Text, Text, IntWritable>{
		
		private Text t=new Text();
        private IntWritable one=new IntWritable(1);
		
        /**
         * 
         * map阶段输出词频
         * 
         * 
         * **/
		@Override
		protected void map(LongWritable key, Text value,Context context)
				throws IOException, InterruptedException {
			String data=value.toString();
			String words[]=data.split(";");
		  if(words[0].trim()!=null){
			  t.set(" "+words[0]);//赋值K
			  one.set(Integer.parseInt(words[1]));
			  context.write(t, one);
		  }	
		}
		
	}
	
	/**
	 * MapReduce作业1的Reducer
	 * 负责词频累加,并输出
	 * 
	 * **/
	private static class SumReduce extends Reducer<Text, IntWritable, IntWritable, Text>{
		
		//存储词频对象
		private IntWritable iw=new IntWritable();
		
		@Override
		protected void reduce(Text key, Iterable<IntWritable> value,Context context)
				throws IOException, InterruptedException {
		 
			
			int sum=0;
			for(IntWritable count:value){
				sum+=count.get();//累加词频
			}
			iw.set(sum);//设置词频
			context.write(iw, key);//输出数据
			
			
			
			
			
		}
		
	}
	
	
	/**
	 * MapReduce作业2排序的Mapper
	 * 
	 * **/
	private static class SortMapper  extends Mapper<LongWritable, Text, IntWritable, Text>{
		
		
		IntWritable iw=new IntWritable();//存储词频
		private Text t=new Text();//存储文本
		@Override
		protected void map(LongWritable key, Text value,Context context)throws IOException, InterruptedException {
			
			String words[]=value.toString().split(" ");
		   System.out.println("数组的长度: "+words.length);
			System.out.println("Map读入的文本: "+value.toString());
			System.out.println("=====>  "+words[0]+"  =====>"+words[1]);
			 if(words[0]!=null){
				 iw.set(Integer.parseInt(words[0].trim()));
				 t.set(words[1].trim());
				 context.write(iw, t);//map阶段输出,默认按key排序
			 }
			
			 
			
		}
		
		
		
	}
	
	
	/**
	 * MapReduce作业2排序的Reducer
	 * 
	 * **/
	private static class SortReduce extends Reducer<IntWritable, Text, Text, IntWritable>{
		
		
		
		/**
		 * 
		 * 输出排序内容
		 * 
		 * **/
		@Override
		protected void reduce(IntWritable key, Iterable<Text> value,Context context)
				throws IOException, InterruptedException {
		 
			 for(Text t:value){
				 context.write(t, key);//输出排好序后的K,V
			 }
			
		}
		
	}
	
	
	
	
	/***
	 * 排序组件,在排序作业中,需要使用
	 * 按key的降序排序
	 * 
	 * **/
		public static class DescSort extends  WritableComparator{

			 public DescSort() {
				 super(IntWritable.class,true);//注册排序组件
			}
			 @Override
			public int compare(byte[] arg0, int arg1, int arg2, byte[] arg3,
					int arg4, int arg5) {
				return -super.compare(arg0, arg1, arg2, arg3, arg4, arg5);//注意使用负号来完成降序
			}
			 
			 @Override
			public int compare(Object a, Object b) {
		 
				return   -super.compare(a, b);//注意使用负号来完成降序
			}
			
		}
	
	
	
	
	
	
    /**
     * 驱动类
     * 
     * **/
 	public static void main(String[] args)throws Exception {
	
		
		   JobConf conf=new JobConf(MyHadoopControl.class); 
		   conf.set("mapred.job.tracker","192.168.75.130:9001");
		   conf.setJar("tt.jar");
		 
		 System.out.println("模式:  "+conf.get("mapred.job.tracker"));;
		
		
		/**
		 * 
		 *作业1的配置
		 *统计词频
		 * 
		 * **/
		Job job1=new Job(conf,"Join1");
	    job1.setJarByClass(MyHadoopControl.class);
		
	    job1.setMapperClass(SumMapper.class);
		job1.setReducerClass(SumReduce.class);
		
		job1.setMapOutputKeyClass(Text.class);//map阶段的输出的key
		job1.setMapOutputValueClass(IntWritable.class);//map阶段的输出的value
		
		job1.setOutputKeyClass(IntWritable.class);//reduce阶段的输出的key
		job1.setOutputValueClass(Text.class);//reduce阶段的输出的value
		
	
		
		//加入控制容器
		ControlledJob ctrljob1=new  ControlledJob(conf);
		ctrljob1.setJob(job1);
		
		
		FileInputFormat.addInputPath(job1, new Path("hdfs://192.168.75.130:9000/root/input"));
        FileSystem fs=FileSystem.get(conf);
		 
		 Path op=new Path("hdfs://192.168.75.130:9000/root/op");
		 
		 if(fs.exists(op)){
			 fs.delete(op, true);
			 System.out.println("存在此输出路径,已删除!!!");
		 }
		FileOutputFormat.setOutputPath(job1, op);
		
	/**========================================================================*/
		
		/**
		 * 
		 *作业2的配置
		 *排序
		 * 
		 * **/
		Job job2=new Job(conf,"Join2");
	    job2.setJarByClass(MyHadoopControl.class);
		
	    //job2.setInputFormatClass(TextInputFormat.class);
	    
	    
	    job2.setMapperClass(SortMapper.class);
		job2.setReducerClass(SortReduce.class);
		
		job2.setSortComparatorClass(DescSort.class);//按key降序排序
		
		job2.setMapOutputKeyClass(IntWritable.class);//map阶段的输出的key
		job2.setMapOutputValueClass(Text.class);//map阶段的输出的value
		
		job2.setOutputKeyClass(Text.class);//reduce阶段的输出的key
		job2.setOutputValueClass(IntWritable.class);//reduce阶段的输出的value
		
		

		//作业2加入控制容器
		ControlledJob ctrljob2=new ControlledJob(conf);
		ctrljob2.setJob(job2);
		
		/***
		 * 
		 * 设置多个作业直接的依赖关系
		 * 如下所写:
		 * 意思为job2的启动,依赖于job1作业的完成
		 * 
		 * **/
		ctrljob2.addDependingJob(ctrljob1);
		
		
		
		//输入路径是上一个作业的输出路径
		FileInputFormat.addInputPath(job2, new Path("hdfs://192.168.75.130:9000/root/op/part*"));
        FileSystem fs2=FileSystem.get(conf);
		 
		 Path op2=new Path("hdfs://192.168.75.130:9000/root/op2");
		 if(fs2.exists(op2)){
			 fs2.delete(op2, true);
			 System.out.println("存在此输出路径,已删除!!!");
		 }
		FileOutputFormat.setOutputPath(job2, op2);
		
		// System.exit(job2.waitForCompletion(true) ? 0 : 1);
		
		
		
		/**====================================================================***/
		
		
		
		
		
		
		/**
		 * 
		 * 主的控制容器,控制上面的总的两个子作业
		 * 
		 * **/
		JobControl jobCtrl=new JobControl("myctrl");
		//ctrljob1.addDependingJob(ctrljob2);// job2在job1完成后,才可以启动
		//添加到总的JobControl里,进行控制
		
		jobCtrl.addJob(ctrljob1); 
		jobCtrl.addJob(ctrljob2);
		
 
		//在线程启动
		Thread  t=new Thread(jobCtrl);
		t.start();
		
		while(true){
			
			if(jobCtrl.allFinished()){//如果作业成功完成,就打印成功作业的信息
				System.out.println(jobCtrl.getSuccessfulJobList());
				
				jobCtrl.stop();
				break;
			}
			
			if(jobCtrl.getFailedJobList().size()>0){//如果作业失败,就打印失败作业的信息
				System.out.println(jobCtrl.getFailedJobList());
				
				jobCtrl.stop();
				break;
			}
			
		}
		
		
		
		
	 
		
		
		
		
		
		
		
		
		
		
		
		
		
		
 
		
		
	}
	
	
	
	

}

运行日志如下:
模式:  192.168.75.130:9001
存在此输出路径,已删除!!!
存在此输出路径,已删除!!!
WARN - JobClient.copyAndConfigureFiles(746) | Use GenericOptionsParser for parsing the arguments. Applications should implement Tool for the same.
INFO - FileInputFormat.listStatus(237) | Total input paths to process : 1
WARN - NativeCodeLoader.<clinit>(52) | Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
WARN - LoadSnappy.<clinit>(46) | Snappy native library not loaded
WARN - JobClient.copyAndConfigureFiles(746) | Use GenericOptionsParser for parsing the arguments. Applications should implement Tool for the same.
INFO - FileInputFormat.listStatus(237) | Total input paths to process : 1
[job name:	Join1
job id:	myctrl0
job state:	SUCCESS
job mapred id:	job_201405092039_0001
job message:	just initialized
job has no depending job:	
, job name:	Join2
job id:	myctrl1
job state:	SUCCESS
job mapred id:	job_201405092039_0002
job message:	just initialized
job has 1 dependeng jobs:
	 depending job 0:	Join1
]

处理的结果如下:
三劫	932
b	579
秦东亮	206
a	45

可以看出,结果是正确的。程序运行成功,上面只是散仙测的2个MapReduce作业的组合,更多的组合其实和上面的一样。
总结:在配置多个作业时,Job的配置尽量分离单独写,不要轻易拷贝修改,这样很容易出错的,散仙在配置的时候,就是拷贝了一个,结果因为少修改了一个地方,在运行时候一直报错,最后才发现,原来少改了某个地方,这一点需要注意一下。
分享到:
评论
2 楼 qindongliang1922 2016-03-11  
vs_diy 写道
写得很不错。thx

1 楼 vs_diy 2016-03-10  
写得很不错。thx

相关推荐

    hadoop0.23.9离线api

    org.apache.hadoop.mapreduce.lib.jobcontrol org.apache.hadoop.mapreduce.lib.join org.apache.hadoop.mapreduce.lib.map org.apache.hadoop.mapreduce.lib.output org.apache.hadoop.mapreduce.lib.partition...

    hadoop-3.1.3-src.tar.gz

    - **核心类库**:如`org.apache.hadoop.fs.FileSystem`、`org.apache.hadoop.mapreduce.Job`等,提供了与HDFS交互和MapReduce编程的基本接口。 4. **开发与调试** - **Hadoop API**:学习如何使用Hadoop API开发...

    Hadoop 2.9.0 已废除属性集

    - mapreduce.jobcontrol.createdir.ifnotexist:已废除,替代为mapreduce.job.dir.createdir.ifnotexist。 - mapreduce.jobhistory.location:已废除,替代为mapreduce.jobhistory.location。 - mapreduce.job...

    基于hadoop的商品推荐引擎.zip

    JobControl是Hadoop中的任务调度管理工具,用于协调和控制MapReduce作业的执行流程。在本项目中,我们需要执行一系列MapReduce任务,JobControl可以确保这些任务按照预设顺序正确地启动和停止,同时提供失败重试机制...

    Hadoop平台搭建方案_hadoop_

    在IT行业中,Hadoop是一个广泛使用的开源框架,用于处理和存储大规模数据。它基于分布式计算模型,能够高效地处理PB级别的数据。本篇文章将详细解析如何构建一个Hadoop平台,以便于你理解并实践Hadoop的搭建过程。 ...

    Hadoop技术内幕:深入解析MapReduce架构设计与实现原理

    MapReduce编程模型3.1 MapReduce编程模型概述3.1.1 MapReduce编程接口体系结构3.1.2 新旧...Hadoop工作流3.5.1 JobControl的实现原理3.5.2 ChainMapperChainReducer的实现原理3.5.3 Hadoop工作流引擎3.6 小结

    hadoop教程

    JobControl 是 Map/Reduce 框架中的一个概念,负责对作业的执行进行控制和监控,以便提高作业的执行效率和可靠性。 数据压缩是 Map/Reduce 框架中的一个概念,负责对数据进行压缩,以便提高数据的传输和存储效率。 ...

    Hadoop_MapReduce教程

    - **JobControl**:JobControl 提供了一组用于管理多个作业的工具。 - **数据压缩**:为了节省存储空间和网络带宽,Hadoop 支持多种压缩算法,可以在作业的不同阶段启用压缩功能。 #### 例子:WordCount v1.0 - **...

    Hadoop 权限管理.pdf

    本文将详细介绍Hadoop中的权限管理机制,包括服务级别授权(Service Level Authorization)、作业队列访问控制(Access Control on Job Queues)以及文件系统权限控制(DFS Permission)。 #### 二、服务级别授权...

    《Hadoop开发者》第四期

    - **ACL**:ZooKeeper使用访问控制列表(Access Control List, ACL)来管理权限。 - **权限类型**:包括CREATE、DELETE、READ、WRITE、ADMIN五种基本权限。 - **认证机制**:支持简单的认证机制,可以通过用户名和密码...

    Hadoop开发者第四期

    2. **ACLs(Access Control Lists)**:使用ACLs来定义节点上的访问权限,包括读、写、创建、删除等操作。 3. **权限模型**:ZooKeeper支持两种权限模型:digest和super。 #### 六、ZooKeeper服务器工作原理和流程 ...

    Hadoop开发第四期

    - **解包**:使用`hadoop fs -untar`命令将tar包解压到HDFS上。 #### 十三、Hadoop技术论坛运营数据分享 - **数据分析**: - **用户行为分析**:通过对用户的行为数据进行分析,了解用户的兴趣偏好。 - **内容...

    spring-hadoop-getting-started:Spring for Apache Hadoop 入门示例

    2. Bean定义:Spring Hadoop将Hadoop的类,如FileSystem、Configuration、Job等,作为Spring的Bean进行管理,使得开发者可以方便地注入和使用。 3. MapReduce支持:Spring Hadoop提供了MapReduce作业的抽象,包括...

    Hadoop架构下的大数据安全存储技术研究.docx

    ResourceManager负责资源调度,NodeManager监控每个节点的资源使用情况,JobHistory Server则记录和存储作业执行的历史信息,便于后续分析和优化。 然而,随着大数据的广泛应用,数据安全问题日益凸显。大数据存储...

    hadoop_the_definitive_guide_3nd_edition

    JobControl 182 Apache Oozie 182 6. How MapReduce Works . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 187 Anatomy of a MapReduce Job Run 187 Classic ...

    20、MapReduce 工作流介绍

    `JobControl` 允许用户创建并管理多个`ControlledJob` 实例,`ControlledJob` 是对标准Hadoop `Job` 的封装,它可以设置作业之间的依赖性。当依赖的前序作业完成后,`JobControl` 会自动提交并启动后续作业,确保...

    MapReduce-Code:Hadoop平台下的MapReduce源码分析

    MapReduce源码分析(主要四大模块,其他表示父目录下的.java文件的总称):1.org.apache.hadoop.mapred(旧版MapReduceAPI):( 1).jobcontrol(job作业直接控制类)(2 ).join :(作业作业中用于模仿数据连接处理...

    mahout0.9测试详细傻瓜说明

    hadoop jar mahout-examples-0.9-job.jar org.apache.mahout.clustering.syntheticcontrol.kmeans.Job ``` 算法运行完成后,你可以通过以下命令检查输出结果: ```bash # 查看输出目录 hadoop fs –ls /user/root/...

    oozie介绍及使用详解

    - **Control Node**:控制流程的走向,如Start Node、End Node、Decision Node、Fork Node和Join Node。 理解Oozie节点的性质和作用对于有效构建和管理工作流至关重要。 总结来说,Oozie是大数据环境中的核心调度...

Global site tag (gtag.js) - Google Analytics