如何从Hadoop MapReduce InputFormat中传递额外信息回作业
写入JobContext配置的可行性判断
这种方案仅适用于极窄场景,不推荐通用,限制如下:
- InputFormat的
getSplits()分片逻辑默认在作业提交客户端的JVM中执行,此时修改JobContext关联的Configuration确实可以在客户端上下文生效,但若开启了AM侧分片(配置mapreduce.job.split.metainfo.maxsize = -1),分片逻辑运行在YARN的ApplicationMaster进程中,修改的配置无法同步回客户端,会导致读取不到指标。 - Configuration本身不保证并发修改安全,若开启多线程分片,多个线程同时更新配置值会出现数据覆盖问题。
如果你的作业固定使用客户端分片、且分片逻辑为单线程执行,可以临时用该方案实现需求,但不属于最优解。
回传额外信息的最优方案
根据场景不同可以选择以下三种实现,兼容性和易用性从高到低排序:
- 方案1:强转JobContext获取原生计数器(最符合Hadoop设计规范)
传入getSplits()的JobContext实例在客户端执行分片时,本质就是你提交作业时构造的Job子类对象,可以直接强转后调用计数器API更新指标:
作业运行结束后,直接在客户端通过// 自定义InputFormat的getSplits方法中添加 Job job = (Job) context; job.getCounter("自定义InputFormat指标组", "特定操作执行次数").increment(1);job.getCounters().findCounter("自定义InputFormat指标组", "特定操作执行次数").getValue()即可拿到统计值,该方式完全复用Hadoop原生的计数器汇总能力,不需要额外开发持久化逻辑,是首选方案。 - 方案2:HDFS临时文件持久化指标(兼容性最强)
若不确定分片逻辑运行在客户端还是AM侧,可以在InputFormat中完成指标统计后,将指标序列化写入预先在配置中指定的HDFS路径,作业运行结束后客户端主动读取该HDFS文件解析指标即可,该方案没有上下文限制,支持所有部署场景。 - 方案3:静态变量存储(轻量简单)
若确定分片逻辑固定在客户端执行,可以在自定义InputFormat类中定义线程安全的静态变量(如AtomicLong)统计指标,作业提交完成后直接读取静态变量的值即可,不需要依赖Hadoop的API扩展,实现成本最低。
内容的提问来源于stack exchange,提问作者Garret Wilson
相关产品推荐
相关产品推荐

