公司动态
Hadoop DistCp 源码分析:命令入口到核心调用链
前些天发现了一个巨牛的人工智能学习网站通俗易懂风趣幽默忍不住分享一下给大家。点击跳转到网站https://www.captainai.net/dongkelun前言上一篇文章《Hadoop DistCp 使用详解与原理分析》介绍了 DistCp 的用法和基本原理。本文从源码出发分析两个问题hadoop distcp命令从敲下回车到 Java 代码经历了什么DistCp 的 13 个核心类之间是如何协作的基于 Hadoop 3.1.1 源码所有引用均标注文件路径和方法名。命令入口从 Shell 到 JavaLinux Bash 分发bin/hadoopLinux 的bin/hadoop脚本运行时走到hadoopcmd_case函数。distcp不在case列表中明确列出的子命令中因此落到*)通配分支# bin/hadoop : hadoopcmd_case() 中的 *) 分支*)HADOOP_CLASSNAME${subcmd}if!hadoop_validate_classname${HADOOP_CLASSNAME};thenhadoop_exit_with_usage1fi;;这里将HADOOP_CLASSNAME直接设为distcp。然后脚本末尾调用hadoop_generic_java_subcmd_handler该函数拼接类路径并执行java命令。不过distcp实际是一个 Java 类的工具Linux 下更常见的做法是通过hadoop_subcommand_distcp动态函数来处理。但 Hadoop 3.1.1 的bin/hadoop中没有为distcp定义单独的子命令函数因此最终会走到通用处理逻辑通过hadoop_add_classpath把TOOL_PATH加入类路径后以org.apache.hadoop.tools.DistCp启动。Windows CMD 分发bin/hadoop.cmdCMD 脚本的做法更直观维护了一个corecommands列表distcp在其中rem bin/hadoop.cmd corecommands 列表 set corecommandsfs version jar checknative conftest distch distcp daemonlog ... for %%i in ( %corecommands% ) do ( if %hadoop-command% %%i set corecommandtrue ) if defined corecommand ( call :%hadoop-command% )匹配后跳转到:distcp标签直接指定主类rem bin/hadoop.cmd :distcp 标签 :distcp set CLASSorg.apache.hadoop.tools.DistCp set CLASSPATH%CLASSPATH%;%TOOL_PATH% goto :eof最终执行java ... org.apache.hadoop.tools.DistCp 剩余参数备选入口mapred distcpbin/mapredLinux和bin/mapred.cmdWindows同样支持mapred distcp指向同一个类org.apache.hadoop.tools.DistCp。DistCp.main()// DistCp.java main()publicstaticvoidmain(Stringargv[]){DistCpdistCpnewDistCp();// 无参构造用于 ToolRunnerCleanupCLEANUPnewCleanup(distCp);ShutdownHookManager.get().addShutdownHook(CLEANUP,SHUTDOWN_HOOK_PRIORITY);exitCodeToolRunner.run(getDefaultConf(),distCp,argv);System.exit(exitCode);}getDefaultConf()加载distcp-default.xml和distcp-site.xml两个配置文件。ToolRunner 回调// ToolRunner.java run()publicstaticintrun(Configurationconf,Tooltool,String[]args){GenericOptionsParserparsernewGenericOptionsParser(conf,args);// ① 过滤 -D、-conf 等 Hadoop 通用参数tool.setConf(conf);String[]toolArgsparser.getRemainingArgs();returntool.run(toolArgs);// ② 回调 DistCp.run()}回到DistCp.run()// DistCp.java run()publicintrun(String[]argv){// ① 参数解析contextnewDistCpContext(OptionsParser.parse(argv));checkSplitLargeFile();// ② 校验大文件分块支持setTargetPathExists();// ③ 检查目标路径是否存在execute();// ④ 执行拷贝}参数解析OptionsParser.parse()OptionsParser基于 Apache Commons CLI 实现。关键步骤// OptionsParser.java parse()publicstaticDistCpOptionsparse(String[]args){// 1. 注册 CLI 选项遍历 DistCpOptionSwitch 枚举// 2. CustomParser 解析命令行参数// 3. 分离源路径和目标路径// - 最后一个参数 targetPath// - 之前的所有参数 sourcePaths或 -f 指定文件列表// 4. Builder 模式逐个设置选项// 5. builder.build() 触发校验}-D参数在更早的GenericOptionsParser阶段已被过滤掉不会到OptionsParser这一层。DistCpOptions.Builder.validate()参数校验规则DistCpOptions.java Builder.validate()-overwrite与-update互斥-atomic不能与-update/-overwrite一起使用-delete与-diff/-rdiff互斥-append必须配合-update-append不能与-skipcrccheck一起使用-diff/-rdiff必须配合-update-v必须配合-log校验通过后才创建不可变的DistCpOptions对象。DistCpContextDistCpContext包装了不可变的DistCpOptions和可变的运行时状态// DistCpContext.javapublicclassDistCpContext{privatefinalDistCpOptionsoptions;// 不可变参数privateListPathsourcePaths;// 运行时路径列表可被改写privatebooleantargetPathExists;// 目标是否存在privatebooleanpreserveRawXattrs;// 是否保留 raw.* xattr}sourcePaths是可变的——GlobbedCopyListing展开通配符后、DistCpSync处理快照后都会改写它。核心调用链总体流程DistCp.main() → ToolRunner.run() → DistCp.run() → OptionsParser.parse() → DistCpOptions (参数解析) → DistCpContext (运行时上下文) → execute() → createAndSubmitJob() → createJob() 配置 MR 作业 → prepareFileListing() 生成拷贝文件列表 → job.submit() 提交 MR 作业 → waitForJobCompletion() 等待完成createJob() — 配置 MR 作业// DistCp.java createJob()privateJobcreateJob()throwsIOException{JobjobJob.getInstance(getConf());// InputFormat通过配置选择 UniformSizeInputFormat 或 DynamicInputFormatjob.setInputFormatClass(DistCpUtils.getStrategy(getConf(),context));// 纯 Map 作业job.setMapperClass(CopyMapper.class);job.setNumReduceTasks(0);// OutputFormatjob.setOutputFormatClass(CopyOutputFormat.class);configureOutputFormat(job);// 设置工作目录、日志路径context.appendToConf(job.getConfiguration());// 选项写入配置传递到 Mapperreturnjob;}InputFormat 选择策略// DistCpUtils.java getStrategy()publicstaticClass?extendsInputFormatgetStrategy(Configurationconf,DistCpContextcontext){// 配置键distcp.{策略名}.strategy.implStringconfLabeldistcp.StringUtils.toLowerCase(context.getCopyStrategy()).strategy.impl;// 默认回退UniformSizeInputFormatreturnconf.getClass(confLabel,UniformSizeInputFormat.class,InputFormat.class);}context.getCopyStrategy()默认值为uniformsize定义在DistCpOptions.Builder中拼接后为distcp.uniformsize.strategy.impl。对比distcp-default.xml中的配置!-- distcp-default.xml --propertynamedistcp.static.strategy.impl/namevalueorg.apache.hadoop.tools.mapred.UniformSizeInputFormat/value/propertypropertynamedistcp.dynamic.strategy.impl/namevalueorg.apache.hadoop.tools.mapred.lib.DynamicInputFormat/value/property注意代码拼接的是distcp.uniformsize.strategy.impl而 XML 中定义的键是distcp.static.strategy.impl。两者不匹配因此代码中的conf.getClass()会查找不到回退到默认值UniformSizeInputFormat.class。这意味着即使想通过distcp-default.xml替换默认 InputFormat配置键distcp.static.strategy.impl也不会被代码使用需通过distcp.uniformsize.strategy.impl这个键来配置。prepareFileListing() — 文件列表生成核心分歧点// DistCp.java prepareFileListing()privatevoidprepareFileListing(Jobjob)throwsException{if(context.shouldUseSnapshotDiff()){// 路径 A-diff/-rdiff 快照差异模式DistCpSyncdistCpSyncnewDistCpSync(context,getConf());if(distCpSync.sync()){createInputFileListingWithDiff(job,distCpSync);}else{thrownewException(DistCp sync failed);}}else{// 路径 B常规模式createInputFileListing(job);}}路径 A快照差异-diff / -rdiff// DistCp.java createInputFileListingWithDiff()privatePathcreateInputFileListingWithDiff(Jobjob,DistCpSyncdistCpSync){// 直接创建 SimpleCopyListing含 distCpSync 引用不走工厂方法CopyListingcopyListingnewSimpleCopyListing(job.getConfiguration(),job.getCredentials(),distCpSync);copyListing.buildListing(fileListingPath,context);}路径 B常规模式// DistCp.java createInputFileListing()protectedPathcreateInputFileListing(Jobjob)throwsIOException{// 通过工厂方法选择 CopyListing 实现CopyListingcopyListingCopyListing.getCopyListing(job.getConfiguration(),job.getCredentials(),context);copyListing.buildListing(fileListingPath,context);}CopyListing 工厂方法 — 三层选择// CopyListing.java getCopyListing()publicstaticCopyListinggetCopyListing(Configurationconfiguration,Credentialscredentials,DistCpContextcontext)throwsIOException{StringcustomClassconfiguration.get(DistCpConstants.CONF_LABEL_COPY_LISTING_CLASS,);if(!customClass.isEmpty()){// 优先级 1用户通过 distcp.copy.listing.class 自定义returninstantiateClass(customClass);}if(context.getSourceFileListing()null){// 优先级 2命令行源路径 → GlobbedCopyListing支持通配符returnnewGlobbedCopyListing(configuration,credentials);}else{// 优先级 3-f 指定文件列表 → FileBasedCopyListingreturnnewFileBasedCopyListing(configuration,credentials);}}CopyListing 三层委托链FileBasedCopyListing — 从 -f 文件读取路径列表 → 委托 GlobbedCopyListing — 展开通配符globStatus → 委托 SimpleCopyListing — 递归遍历目录树写 SequenceFile每个委托层都会改写context.setSourcePaths()将展开后的路径传递给下一层。SimpleCopyListing.doBuildListing()是核心使用ProducerConsumer模式多线程递归遍历目录线程数由-numListstatusThreads控制默认 1上限 40每个文件/目录写一条 SequenceFile 记录key 相对路径Textvalue 文件状态CopyListingFileStatusCopyListing.buildListing()作为模板方法统一流程// CopyListing.java buildListing()publicfinalvoidbuildListing(PathpathToListFile,DistCpContextdistCpContext){validatePaths(distCpContext);// 1. 校验路径doBuildListing(pathToListFile,distCpContext);// 2. 子类构建列表config.set(CONF_LABEL_LISTING_FILE_PATH,...);// 3. 写入配置config.setLong(CONF_LABEL_TOTAL_BYTES_TO_BE_COPIED,...);config.setLong(CONF_LABEL_TOTAL_NUMBER_OF_RECORDS,...);validateFinalListing(pathToListFile,...);// 4. 校验重复 ACL/XAttr}DistCpSync — 快照差异路径当-diff或-rdiff启用时prepareFileListing()先调DistCpSync.sync()做同步DistCpSync.sync() 1. preSyncCheck() - 只能有一个源目录 - 源和目标都必须是 DistributedFileSystem - 目标端从 fromSnapshot 以来不能有变化 2. getAllDiffs() - 调用 fs.getSnapshotDiffReport(ssDir, from, to) - 按 DiffType 分类MODIFY/CREATE/DELETE/RENAME 3. syncDiff() - moveToTmpDir()将 rename/delete 的文件移到临时目录按 source 倒序 - moveToTarget()从临时目录 rename 到最终位置按 target 正序 4. 改写 sourcePaths 为快照路径随后prepareDiffListForCopyListing()只返回 MODIFY 和 CREATE 条目DELETE/RENAME 已由 sync 处理给SimpleCopyListing.doBuildListingWithSnapshotDiff()使用。UniformSizeInputFormat 与 DynamicInputFormatUniformSizeInputFormat默认按字节数均匀切分分片// UniformSizeInputFormat.java getSplits()longnBytesPerSplit(long)Math.ceil(totalSizeBytes*1.0/numSplits);while(reader.next(srcRelPath,srcFileStatus)){// 累积大小超过阈值时切分不跨文件if(currentSplitSizesrcFileStatus.getChunkLength()nBytesPerSplitlastPosition!0){splits.add(newFileSplit(listingFilePath,lastSplitStart,...));lastSplitStartlastPosition;currentSplitSize0;}currentSplitSizesrcFileStatus.getChunkLength();lastPositionreader.getPosition();}DynamicInputFormatWorker 模式——文件列表拆成小 chunkMapper 处理完一个领下一个splitCopyListingIntoChunksWithShuffle() 1. 计算每个 chunk 的条目数numEntriesPerChunk ceil(numRecords / (splitRatio * numMaps)) 2. 遍历 SequenceFile轮询写入多个 chunk 文件 3. 创建 N 个空 split每个关联一个 chunk初始分配 4. Mapper 的 DynamicRecordReader 处理完当前 chunk 后自动领下一个适合文件大小差异大的场景避免 UniformSize 中大文件卡住一个 Mapper其他空闲的问题。CopyMapper 端实际文件拷贝整体流程CopyMapper.setup() → 从 Configuration 读取参数syncFolders、skipCrc、overWrite 等 CopyMapper.map(relPath, sourceFileStatus) → 构造 targetPath targetWorkPath relPath → 获取 sourceCurrStatus实时 FileStatus → 若为目录createTargetDirsWithRetry() → checkUpdate() 决策动作 ├── SKIP目标存在且一致 ├── APPEND尾部追加 └── OVERWRITE覆盖写入 → copyFileWithRetry() 执行拷贝 → DistCpUtils.preserve() 保留文件属性checkUpdate() — 拷贝决策// CopyMapper.java checkUpdate()privateFileActioncheckUpdate(FileSystemsourceFS,CopyListingFileStatussource,Pathtarget,FileStatustargetFileStatus){if(targetFileStatus!null!overWrite){if(canSkip(sourceFS,source,targetFileStatus)){returnFileAction.SKIP;// 目标存在且一致}elseif(append){longtargetLentargetFileStatus.getLen();if(targetLensource.getLen()){// 尾部校验和一致 → APPENDif(sourceChecksum.equals(targetFS.getFileChecksum(target))){returnFileAction.APPEND;}}}}returnFileAction.OVERWRITE;// 覆盖}canSkip() — 能否跳过// CopyMapper.java canSkip()privatebooleancanSkip(...){if(!syncFolders)returntrue;// 非 -update目标存在即跳过// -update 模式比较长度 BlockSize CRCbooleansameLengthtarget.getLen()source.getLen();if(sameLengthsameBlockSize){returnskipCrc||DistCpUtils.checksumsAreEqual(...);}returnfalse;// 大小或时间不一致需要拷贝}RetriableFileCopyCommand — 带重试的文件拷贝RetriableFileCopyCommand.execute() → doCopy() → copyToFile() → OVERWRITE创建临时文件 .distcp.tmp.attempt_xxx → APPEND直接 append → copyBytes() → getInputStream() 打开源包装 ThrottledInputStream 限速 → seekIfRequired() 定位到偏移量 → 循环 read → write → compareFileLengths() 校验长度 → compareCheckSums() 校验 CRC → promoteTmpToTarget() rename 临时文件 → 目标失败后自动重试继承RetriableCommand默认 3 次指数退避。ThrottledInputStream — 令牌桶限速// ThrottledInputStream.java throttle()privatevoidthrottle()throwsIOException{while(getBytesPerSec()maxBytesPerSec){Thread.sleep(SLEEP_DURATION_MS);// 50mstotalSleepTimeSLEEP_DURATION_MS;}}速率计算为整个生命周期平均值短期可能超出但长期趋近设定值。CopyCommitter — 作业后处理commitJob()在 MapReduce 框架回调中执行CopyCommitter.commitJob() 1. concatFileChunks() 大文件分块合并HDFS concat 2. super.commitJob() FileOutputCommitter 标准提交 3. cleanupTempFiles() 清理 .distcp.tmp.attempt_* 残留 4. preserveFileAttributesForDirectories() 保留目录属性 5. 分支互斥 ├─ -deletedeleteMissing() 比对源/目标列表删除目标多余文件 ├─ -atomiccommitData() rename 临时目录 → 最终目录 └─ -track trackMissing() 保存缺失文件列表 6. cleanup() 删除 metaFolder类运行位置汇总类运行端职责DistCpDriver入口配置作业提交DistCpOptionsDriver不可变参数模型通过 Configuration 传递到 MapperOptionsParserDriver命令行解析DistCpContextDriver运行时上下文DistCpConstantsDriver Mapper常量定义DistCpSyncDriver快照差异同步CopyListing(abstract)Driver文件列表抽象SimpleCopyListingDriver递归遍历目录树生成 SequenceFileGlobbedCopyListingDriver展开通配符委托 SimpleCopyListingFileBasedCopyListingDriver从文件读路径委托 GlobbedCopyListingCopyMapperMapper执行文件拷贝RetriableFileCopyCommandMapper带重试的拷贝包装限速流ThrottledInputStreamMapper令牌桶带宽限速UniformSizeInputFormatDriver Mapper按字节数均匀分片DynamicInputFormatDriver Mapper动态 Worker 模式分片CopyCommitterDriver作业后处理MR 框架回调CopyOutputFormatDriver MapperOutputFormat 配置DistCpUtilsDriver Mapper工具方法DiffInfoDriver快照差异数据载体调用流程图hadoop distcp -update src tgt │ ├─ bin/hadoop.cmd → set CLASSorg.apache.hadoop.tools.DistCp [cmd脚本] │ ├─ DistCp.main() [DistCp.java] │ └─ ToolRunner.run(getDefaultConf(), distCp, argv) │ └─ GenericOptionsParser 过滤 -D 参数 │ ├─ DistCp.run() │ ├─ OptionsParser.parse(argv) → DistCpOptions │ │ └─ Builder.build() → validate() 校验参数 │ ├─ new DistCpContext(options) │ └─ execute() │ ├─ DistCp.execute() │ └─ createAndSubmitJob() │ ├─ DistCp.createJob() │ ├─ DistCpUtils.getStrategy() → UniformSizeInputFormat │ ├─ job.setMapperClass(CopyMapper.class) │ ├─ configureOutputFormat() // 工作目录/日志路径 │ └─ context.appendToConf(config) │ ├─ prepareFileListing(job) │ └─ 常规模式 → createInputFileListing(job) │ └─ CopyListing.getCopyListing() 工厂方法 │ └─ GlobbedCopyListing → globStatus → 委托 │ └─ SimpleCopyListing → 递归遍历目录树 │ → 写入 SequenceFilerelativePath, FileStatus │ ├─ job.submit() → MR 作业提交 │ └─ [MR 运行时] ├─ UniformSizeInputFormat.getSplits() 生成分片 └─ CopyMapper.map() 执行拷贝 ├─ checkUpdate() → canSkip() → decision ├─ RetriableFileCopyCommand.execute() │ ├─ ThrottledInputStream.throttle() 限速 │ ├─ copyBytes() → read → write │ └─ promoteTmpToTarget() → rename └─ DistCpUtils.preserve() 保留属性 └─ [作业完成] └─ CopyCommitter.commitJob() ├─ preserveFileAttributesForDirectories() ├─ deleteMissing() / commitData() / trackMissing() └─ cleanup()设计要点装饰器 模板方法FileBasedCopyListing→GlobbedCopyListing→SimpleCopyListing每层只做一件事CopyListing.buildListing()作为模板定义固定流程。Builder 模式DistCpOptions使用 Builder 构建build()时做参数校验构建后不可变。DistCpContext 可变性sourcePaths和targetPathExists可在运行时改写适应 Glob 展开和快照处理。配置键不对齐代码按distcp.{strategy}.strategy.impl查找而distcp-default.xml中定义的是distcp.static.strategy.impl实际上未被使用依赖代码中的回退默认值UniformSizeInputFormat.class。纯 Map 作业job.setNumReduceTasks(0)不设 Reduce 阶段所有文件拷贝在 Mapper 中完成。