首页 文章 精选 留言 我的

精选列表

搜索[流程管理],共10009篇文章
优秀的个人博客,低调大师

hive执行流程(3)-Driver类分析1Driver类整体流程

Driver类是对 1 org.apache.hadoop.hive.ql.processors.CommandProcessor.java 接口的实现,重写了run方法,定义了常见sql的执行方式. 1 public class Driver implements CommandProcessor 具体的方法调用顺序: 1 2 run--->runInternal--->(createTxnManager+recordValidTxns)----->compileInternal---> compile--analyzer(BaseSemanticAnalyzer)--->execute 其中compile和execute是两个比较重要的方法: compile用来完成语法和语义的分析,生成执行计划 execute执行物理计划,即提交相应的mapredjob 通过打印perflog可以看到Driver类的简单地时序图: 下面来看下Driver类的几个常用的方法实现: 1)createTxnManager 用来获取目前设置的用于实现lock的类,比如: 1 org.apache.hadoop.hive.ql.lockmgr.DummyTxnManager 2)checkConcurrency 用来判断当前hive设置是否支持并发控制: 1 boolean supportConcurrency=conf.getBoolVar(HiveConf.ConfVars.HIVE_SUPPORT_CONCURRENCY); 主要是通过判断hive.support.concurrency参数,默认是false 3)getClusterStatus 调用JobClient类的getClusterStatus方法来获取集群的状态: 1 2 3 4 5 6 7 8 9 10 11 12 13 public ClusterStatusgetClusterStatus() throws Exception{ ClusterStatuscs; try { JobConfjob= new JobConf(conf,ExecDriver. class ); JobClientjc= new JobClient(job); cs=jc.getClusterStatus(); } catch (Exceptione){ e.printStackTrace(); throw e; } LOG.info( "Returningclusterstatus:" +cs.toString()); return cs; } 4)getSchema //返回表的schema信息 5) 1 doAuthorization/doAuthorizationV2/getHivePrivObjects 用来在开启权限验证情况下对sql的权限检测操作 6) 1 getLockObjects/acquireReadWriteLocks/releaseLocks 都是和锁相关的方法 ,其中getLockObjects用来获取锁的对象(锁的路径,锁的模式等),最终返回一个包含所有锁的list,acquireReadWriteLocks用来控制获取锁,releaseLocks用来释放锁: 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 getLockObjects: private List<HiveLockObj>getLockObjects(Databased,Tablet,Partitionp,HiveLockModemode) throws SemanticException{ List<HiveLockObj>locks= new LinkedList<HiveLockObj>(); HiveLockObjectDatalockData= new HiveLockObjectData(plan.getQueryId(), String.valueOf(System.currentTimeMillis()), "IMPLICIT" , plan.getQueryStr()); if (d!= null ){ locks.add( new HiveLockObj( new HiveLockObject(d.getName(),lockData),mode)); //数据库层面的锁 return locks; } if (t!= null ){ //表层面的锁 locks.add( new HiveLockObj( new HiveLockObject(t.getDbName(),lockData),mode)); locks.add( new HiveLockObj( new HiveLockObject(t,lockData),mode)); mode=HiveLockMode.SHARED; locks.add( new HiveLockObj( new HiveLockObject(t.getDbName(),lockData),mode)); return locks; } if (p!= null ){ //分区层面的锁 locks.add( new HiveLockObj( new HiveLockObject(p.getTable().getDbName(),lockData),mode)); if (!(p instanceof DummyPartition)){ locks.add( new HiveLockObj( new HiveLockObject(p,lockData),mode)); } //Alltheparentsarelockedinsharedmode mode=HiveLockMode.SHARED; //Fordummypartitions,onlypartitionnameisneeded Stringname=p.getName(); if (p instanceof DummyPartition){ name=p.getName().split( "@" )[ 2 ]; } StringpartialName= "" ; String[]partns=name.split( "/" ); int len=p instanceof DummyPartition?partns.length:partns.length- 1 ; Map<String,String>partialSpec= new LinkedHashMap<String,String>(); for ( int idx= 0 ;idx<len;idx++){ Stringpartn=partns[idx]; partialName+=partn; String[]nameValue=partn.split( "=" ); assert (nameValue.length== 2 ); partialSpec.put(nameValue[ 0 ],nameValue[ 1 ]); try { locks.add( new HiveLockObj( new HiveLockObject( new DummyPartition(p.getTable(),p.getTable().getDbName() + "/" +p.getTable().getTableName() + "/" +partialName, partialSpec),lockData),mode)); partialName+= "/" ; } catch (HiveExceptione){ throw new SemanticException(e.getMessage()); } } locks.add( new HiveLockObj( new HiveLockObject(p.getTable(),lockData),mode)); locks.add( new HiveLockObj( new HiveLockObject(p.getTable().getDbName(),lockData),mode)); } return locks; } acquireReadWriteLocks调用了锁具体实现类的acquireLocks方法 releaseLocks调用了锁具体实现类的releaseLocks方法 7) run方法是Driver类的入口方法,调用了runInternal方法,我们主要来看runInternal的方法,大体步骤: 1 2 3 4 5 运行hive.exec.driver.run.hooks中设置的hook, 运行HiveDriverRunHook相关类的的preDriverRun方法---->检测是否支持并发,并获取并发实现的类 --->compileInternal---->运行锁相关的操作(判断是否只对mapredjob进行锁,获取锁等) ---->调用execute---->释放锁--->运行HiveDriverRunHook相关类的的postDriverRun方法 ---->返回CommandProcessorResponse对象 相关代码: 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 private CommandProcessorResponserunInternal(Stringcommand, boolean alreadyCompiled) throws CommandNeedRetryException{ errorMessage= null ; SQLState= null ; downstreamError= null ; if (!validateConfVariables()){ return new CommandProcessorResponse( 12 ,errorMessage,SQLState); } HiveDriverRunHookContexthookContext= new HiveDriverRunHookContextImpl(conf,command); //Getallthedriverrunhooksandpre-executethem. List<HiveDriverRunHook>driverRunHooks; try { //运行hive.exec.driver.run.hooks中设置的hook driverRunHooks=getHooks(HiveConf.ConfVars.HIVE_DRIVER_RUN_HOOKS, HiveDriverRunHook. class ); for (HiveDriverRunHookdriverRunHook:driverRunHooks){ driverRunHook.preDriverRun(hookContext); //运行HiveDriverRunHook相关类的的preDriverRun方法 } } catch (Exceptione){ errorMessage= "FAILED:HiveInternalError:" +Utilities.getNameMessage(e); SQLState=ErrorMsg.findSQLState(e.getMessage()); downstreamError=e; console.printError(errorMessage+ "\n" +org.apache.hadoop.util.StringUtils.stringifyException(e)); return new CommandProcessorResponse( 12 ,errorMessage,SQLState); } //Resettheperflogger PerfLoggerperfLogger=PerfLogger.getPerfLogger( true ); perfLogger.PerfLogBegin(CLASS_NAME,PerfLogger.DRIVER_RUN); perfLogger.PerfLogBegin(CLASS_NAME,PerfLogger.TIME_TO_SUBMIT); int ret; boolean requireLock= false ; boolean ckLock= false ; try { ckLock=checkConcurrency(); //检测是否支持并发,并获取并发实现的类,比如常用的org.apache.hadoop.hive.ql.lockmgr.DummyTxnManager createTxnManager(); } catch (SemanticExceptione){ errorMessage= "FAILED:Errorinsemanticanalysis:" +e.getMessage(); SQLState=ErrorMsg.findSQLState(e.getMessage()); downstreamError=e; console.printError(errorMessage, "\n" +org.apache.hadoop.util.StringUtils.stringifyException(e)); ret= 10 ; return new CommandProcessorResponse(ret,errorMessage,SQLState); } ret=recordValidTxns(); if (ret!= 0 ) return new CommandProcessorResponse(ret,errorMessage,SQLState); if (!alreadyCompiled){ ret=compileInternal(command); //调用compileInternal方法 if (ret!= 0 ){ return new CommandProcessorResponse(ret,errorMessage,SQLState); } } //thereasonthatwesetthetxnmanagerforthecxthereisbecauseeach //queryhasitsownctxobject.Thetxnmgrissharedacrossthe //sameinstanceofDriver,whichcanrunmultiplequeries. ctx.setHiveTxnManager(txnMgr); if (ckLock){ //断是否只对mapredjob进行锁,参数hive.lock.mapred.only.operation,默认为false boolean lockOnlyMapred=HiveConf.getBoolVar(conf,HiveConf.ConfVars.HIVE_LOCK_MAPRED_ONLY); if (lockOnlyMapred){ Queue<Task<? extends Serializable>>taskQueue= new LinkedList<Task<? extends Serializable>>(); taskQueue.addAll(plan.getRootTasks()); while (taskQueue.peek()!= null ){ Task<? extends Serializable>tsk=taskQueue.remove(); requireLock=requireLock||tsk.requireLock(); if (requireLock){ break ; } if (tsk instanceof ConditionalTask){ taskQueue.addAll(((ConditionalTask)tsk).getListTasks()); } if (tsk.getChildTasks()!= null ){ taskQueue.addAll(tsk.getChildTasks()); } //doesnotaddbackuptaskhere,becausebackuptaskshouldbethesame //typeoftheoriginaltask. } } else { requireLock= true ; } } if (requireLock){ //获取锁 ret=acquireReadWriteLocks(); if (ret!= 0 ){ try { releaseLocks(ctx.getHiveLocks()); } catch (LockExceptione){ //Notmuchtodohere } return new CommandProcessorResponse(ret,errorMessage,SQLState); } } ret=execute(); //job运行 if (ret!= 0 ){ //ifneedRequireLockisfalse,thereleaseherewilldonothingbecausethereisnolock try { releaseLocks(ctx.getHiveLocks()); } catch (LockExceptione){ //Nothingtodohere } return new CommandProcessorResponse(ret,errorMessage,SQLState); } //ifneedRequireLockisfalse,thereleaseherewilldonothingbecausethereisnolock try { releaseLocks(ctx.getHiveLocks()); } catch (LockExceptione){ errorMessage= "FAILED:HiveInternalError:" +Utilities.getNameMessage(e); SQLState=ErrorMsg.findSQLState(e.getMessage()); downstreamError=e; console.printError(errorMessage+ "\n" +org.apache.hadoop.util.StringUtils.stringifyException(e)); return new CommandProcessorResponse( 12 ,errorMessage,SQLState); } perfLogger.PerfLogEnd(CLASS_NAME,PerfLogger.DRIVER_RUN); perfLogger.close(LOG,plan); //Takeallthedriverrunhooksandpost-executethem. try { for (HiveDriverRunHookdriverRunHook:driverRunHooks){ //运行HiveDriverRunHook相关类的的postDriverRun方法 driverRunHook.postDriverRun(hookContext); } } catch (Exceptione){ errorMessage= "FAILED:HiveInternalError:" +Utilities.getNameMessage(e); SQLState=ErrorMsg.findSQLState(e.getMessage()); downstreamError=e; console.printError(errorMessage+ "\n" +org.apache.hadoop.util.StringUtils.stringifyException(e)); return new CommandProcessorResponse( 12 ,errorMessage,SQLState); } return new CommandProcessorResponse(ret); } 8) 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 再来看下compileInternal方法 private static final ObjectcompileMonitor= new Object(); private int compileInternal(Stringcommand){ int ret; synchronized (compileMonitor){ ret=compile(command); //调用compile方法 } if (ret!= 0 ){ try { releaseLocks(ctx.getHiveLocks()); } catch (LockExceptione){ LOG.warn( "Exceptioninreleasinglocks." +org.apache.hadoop.util.StringUtils.stringifyException(e)); } } return ret; } 调用了compile方法,compile方法分析命令,生成Task,关于compile的具体实现后面详细讲解 9.execute方法,提交task并等待task运行完毕,并打印task运行的信息,比如消耗的时间等 (这里信息也比较多,后面单独讲解 本文转自菜菜光 51CTO博客,原文链接:http://blog.51cto.com/caiguangguang/1571890,如需转载请自行联系原作者

优秀的个人博客,低调大师

Antenna Magus 天线设计流程

Antenna Magus 是用于加速天线设计和建模过程的软件工具。包含350多个经过验证的天线数据库,用户只需输入天线指标,便能在几分钟内得出满足要求的参数化天线模型,且模型支持导出到CST中做精确仿真。无论对天线设计工程师,还是使用天线模型的电磁兼容工程师以及研究天线布局的系统工程师来说,Magus天线库都是一个极其有用的工具。

优秀的个人博客,低调大师

花呗分期集成流程

本期帖子主要讲花呗支付与花呗分期支付 一、文档地址 官方文档地址:[url]https://docs.open.alipay.com/277/106748/[/url] 二、开发前准备工作 调用步骤:[url]https://openclub.alipay.com/read.php?tid=12194&fid=69[/url] 注意事项:1、不支持沙箱测试;2、具体是否签约需根据场景决定; 如何签约以及签约无法成功等相关签约问题:[url]https://openclub.alipay.com/read.php?tid=276&fid=72[/url] 三、接口集成代码示例 花呗分期单渠道支付示例demo:[url]https://openclub.alipay.com/read.php?tid=13785&fid=56[

优秀的个人博客,低调大师

小程序完整上线流程

小程序需要经过以下几个阶段,方可完全上线: 1、[url=https://docs.alipay.com/mini/introduce/register]入驻开放平台[/url] 2、[url=https://docs.alipay.com/mini/introduce/create]创建小程序[/url] 3、[url=https://openclub.alipay.com/read.php?tid=1847&fid=51]设置小程序[/url] 4、[url=https://openclub.alipay.com/read.php?tid=1839&fid=51]开发与测试[/url] 5、[url=https://openclub.alipay.com/read.php?tid=1841&fid=51]使用体

优秀的个人博客,低调大师

Android渠道包打包流程

1.环境要求 Windows、JDK1.7.0以上、WinRAR 2.打包步骤 (1)从Jenkins打包平台取得最终作为发版外卖apk (2)apk重命名为src.zip(没错,就是改成一个压缩包) (3)打包工具解压,将src.zip解压到打包工具目录,如图 (4)先看一下打包工具目录,build.bat为我们最终执行打包任务的批处理文件,foodfinder.keystore就是我们打包使用的签名,sources.txt中保存着我们要打包使用的渠道号(以及产物apk文件名),最终产物会保存到release文件夹中,至于sources.txt中渠道号与产物文件名填写方式可直接通过后面提供的渠道号Excel表格复制后粘贴进TXT中即可 (5)build.bat文件中需注意的点: path后面winRAR的路径要填写好,其次是保证jdk、Java等环境变量配置OK,否则会报错,双击.bat文件,弹出命令框,等待至提示“按任意键结束”即可 (6)产物构建完之后,记得要及时将release文件夹中产物取出,否则下次构建会直接被当前构建产物覆盖 3.打包渠道号以及注意事项 (1)渠道号列表,左侧为全部应用商店渠道号,右侧为地推用渠道号(目前暂不需要,只需应用商店渠道即可,若PM要求另论) (2)不同应用商店使用的副标题 电子表格 副标题apk需要由RD提供,QA依据不同副标题母包进行渠道包打包操作

资源下载

更多资源
腾讯云软件源

腾讯云软件源

为解决软件依赖安装时官方源访问速度慢的问题,腾讯云为一些软件搭建了缓存服务。您可以通过使用腾讯云软件源站来提升依赖包的安装速度。为了方便用户自由搭建服务架构,目前腾讯云软件源站支持公网访问和内网访问。

Rocky Linux

Rocky Linux

Rocky Linux(中文名:洛基)是由Gregory Kurtzer于2020年12月发起的企业级Linux发行版,作为CentOS稳定版停止维护后与RHEL(Red Hat Enterprise Linux)完全兼容的开源替代方案,由社区拥有并管理,支持x86_64、aarch64等架构。其通过重新编译RHEL源代码提供长期稳定性,采用模块化包装和SELinux安全架构,默认包含GNOME桌面环境及XFS文件系统,支持十年生命周期更新。

Sublime Text

Sublime Text

Sublime Text具有漂亮的用户界面和强大的功能,例如代码缩略图,Python的插件,代码段等。还可自定义键绑定,菜单和工具栏。Sublime Text 的主要功能包括:拼写检查,书签,完整的 Python API , Goto 功能,即时项目切换,多选择,多窗口等等。Sublime Text 是一个跨平台的编辑器,同时支持Windows、Linux、Mac OS X等操作系统。

WebStorm

WebStorm

WebStorm 是jetbrains公司旗下一款JavaScript 开发工具。目前已经被广大中国JS开发者誉为“Web前端开发神器”、“最强大的HTML5编辑器”、“最智能的JavaScript IDE”等。与IntelliJ IDEA同源,继承了IntelliJ IDEA强大的JS部分的功能。

用户登录
用户注册