Hadoop学习总结之四:Map-Reduce的过程解析_第1页
Hadoop学习总结之四:Map-Reduce的过程解析_第2页
Hadoop学习总结之四:Map-Reduce的过程解析_第3页
Hadoop学习总结之四:Map-Reduce的过程解析_第4页
Hadoop学习总结之四:Map-Reduce的过程解析_第5页
已阅读5页,还剩21页未读 继续免费阅读

下载本文档

版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领

文档简介

1、一、客户Map-Reduce的过程首先从客户端发送任务开始。提交任务主要通过JobClient.runJob(JobConf )静态函数实现publicationstaticrunningjobrunjob (jobconfjob ) throwsioexception;/第一位老师成为作业客户端对象作业客户端JC=新作业客户端(job )调用submitJob来发送任务running=jc.submitJob(job )JobID jobId=running.getID ();while (true) /while循环持续获取此任务的状态,并打印在客户端控制台上以下返回运行;以下其中,job

2、客户端的submitJob函数实现如下publicunningjobsubmitjob (jobconfjob ) throws文件not found异常InvalidJobConfException,IOException 从作业跟踪器中获取当前任务的idjobid jobid=jobsubmittlient.getnewjobid ();/准备将执行任务所需的元素写入HDFS :/任务执行程序所在的jar封装在job.jar中/任务处理的input split信息被写入job.split/任务执行的构成项目被汇总写入到job.xml中pathsubjectdir=new path (get

3、 system dir ()、jobId.toString ();pathsubmitjarfile=new path (submitjobdir, job.jar );pathsubsubplitfile=new path (subtkjobdir, job.split );/将在此处-libjars命令行中指定的jar上载到HDFSconfigureecommandndlineoptions (job,submitJobDir,submitJarFile )pathsubtjjobfile=new path (subtkjobdir, job.xml );以输入格式的形式获取适当的输入剥离

4、。 默认类型为FileSplit输入split splits=job.getinputformat ().get splits (job,job.getNumMapTasks () );生成将input split信息写入job.split文件的写入流fsdataoutputstreamout=文件系统. create (fssubmitSplitFile,newfsspermission (job _ file _ permission );try被写入job.split文件的信息依次写入split文件标头、split文件版本号、split的数量、接着按每个input split的信息。对于

5、每个输入split,split类型名(默认的FileSplit )、split的大小、split的内容(在FileSplit的情况下,写入文件名,该split位于文件的开头)、split的位置信息(其writeSplitsFile(splits,out ) finally out.close ();以下job.set(mapred.job.split.file ,submitSplitFile.toString ();根据split的个数设定映射任务的个数job.setnummaptasks (splits.length )把/job的配置信息写入job.xml文件out=FileSystem

6、.create(fs,submitJobFile,newfs permission (job _ file _ permission );tryjob.writeXml(out ) finally out.close ();以下/真的调用JobTracker来发送任务jobstatustatus=jobsubmtclient.submit job (jobid )以下二、作业跟踪器JobTracker作为单独的JVM运行,主函数的主要调用由以下两部分组成通过调用静态函数startTracker(new JobConf () )创建JobTracker对象调用JobTracker.offerSe

7、rvice ()函数来提供服务JobTracker构造函数通过生成taskScheduler成员变量来调度作业。 缺省值为JobQueueTaskScheduler。 也就是说,以FIFO方式调度作业。off service函数调用taskScheduler.start ()。 此函数在作业跟踪器(任务排程器的taskracker管理器)中注册了两个监听器jobqueueejobinprogresslistenerjobqueuejobinprogresslistener用于监视作业的执行状态eagteraskinitionalizationlistenereaertaskinitialize

8、listener用于初始化作业eaertaskinitializelistener具有线程JobInitThread,它获取jobInitQueue的JobInProgress对象,并调用JobInProgress对象的initTasks函数来执行任务在上一部分中,客户端调用了JobTracker.submitJob函数。 此函数在调用addJob函数之前先使第一个老师成为JobInProgress对象。 此函数包含以下逻辑:输入同步(jobs ) )已同步(taskscheduler )jobs.put (job.get profile ().get jobid (),job;为JobTra

9、cker的每个监听器调用jobAdded函数for (jobinprogresslistenerlistener : jobinprogresslistners ) 监听器.作业添加(作业)以下以下以下eaertaskinitializelistener的jobAdded函数通过向jobInitQueue添加JobInProgress对象,自然开始初始化该作业,JobInProgress将initTasks函数添加到publicsynchronizedvoidinittasks () throwsioexception从HDFS导入job.split文件,并生成input splitsstri

10、ngjobfile=profile.getjobfile ();path sysdir=new path (this.job tracker.get system dir ();文件系统fs=sysdir.getfile system (conf )数据输入剥离文件=新路径(conf.get ( mapred.job.split.file ) );作业客户端. rawsplit splits;trysplits=job client.readsplitfile (split file ) finally splitFile.close ();以下/map task的个数是input split

11、的个数。numMapTasks=splits.length;为每个映射任务生成TaskInProgress以处理输入剥离maps=newtaskinprogress nummaptasks ;for(int i=0; i numMapTasks; 表示I )输入长度=splits I .get datalength ();maps I =newtaskkinprogress (作业文件,splits是作业跟踪器,conf,this,I;以下/Maptask时,放入nonRunningMapCache中。 这是一个map,在map task中,它被指定给包含输入剥离的节点。 nonRunning

12、MapCache用于作业跟踪器将映射任务分配给任务器。if (numMapTasks 0) nonrunningmapcache=create cache (splits,maxLevel )以下创建reduce taskthis.reductions=newtaskinprogress numreductiontasks ;for (int i=0; numreductiontasks; 表示I )red ces I =newstaskinprogress (作业文件,numMapTasks,I作业跟踪器、conf、this;/reduce task位于nonRunningReduces中,

13、用于作业跟踪器将reduce task分配给TaskTracker。nonrunningreduces.add (reduces I );以下/创建两个清除up任务,一个清除map,另一个清除reduce清除=新任务程序2;清除0=newtaskkinprogress (作业,作业文件,splits0 )作业跟踪器、conf、this、numMapTasks;cleanup0.setJobCleanupTask ();清除1=newtaskkinprogress (jobid,作业文件,numMapTasks )数字跟踪,作业跟踪,conf,this;cleanup1.setJobCleanu

14、pTask ();/创建两个初始化任务、一个初始化映射和一个初始化结果setup=new TaskInProgress2;setup 0=newtaskkinprogress (作业,作业文件,splits0 )作业跟踪器、conf、this、numMapTasks 1;setup0.setJobSetupTask ();setup 1=newtaskkinprogress (jobid,job文件,numMapTasks )numreductiontasks1,作业跟踪器,conf,this;setup1.setJobSetupTask ();tasksInited.set(true) /初始化完成以下三、TaskTrackerTaskTracker也用作单独的JVM,其主函数调用new TaskTracker(conf).run (),而run函数主要调用如下stateofforservice () throws exception长最后一个头部=0;/TaskTracker的进展一直存在while (运行! shuttertingdown )“请长now=system.current time mills ();/每隔一段时间向作业跟踪器发送heartbeat长等待时间=heartbeaterval-(no

温馨提示

  • 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
  • 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
  • 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
  • 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
  • 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
  • 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
  • 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。

评论

0/150

提交评论