Hadoop MapReduce配置3个链式作业仅首个运行生成单输出如何解决
问题根因排查&解决步骤
- 检查运行参数数量是否正确
你的代码逻辑要求传入4个位置参数:
- args[0]:原始输入路径
- args[1]:job1输出路径(同时是job2的输入路径)
- args[2]:job2输出路径
- args[3]:job3输出路径
如果运行时传入参数少于4个,会触发数组越界异常,程序在job1执行完成后直接终止,后续作业不会运行。
正确运行命令示例:
hadoop jar your_jar_file.jar mapreduce.RunMapReduceJob /input/path /output/job1 /output/job2 /output/job3- 检查运行参数数量是否正确
- 提前清理所有输出路径
Hadoop MapReduce默认不允许作业输出路径提前存在,若args[1]/args[2]/args[3]任意一个路径在运行前已存在,对应作业会直接提交失败。你可以手动执行删除命令,也可以在代码中新增自动清理逻辑。
- 提前清理所有输出路径
- 增加作业执行结果校验
原代码未判断前序作业的执行结果,若job1执行失败,后续依赖job1输出的job2必然无法正常运行,建议每个作业执行完成后校验返回状态。
- 增加作业执行结果校验
修改后可直接运行的代码
package mapreduce; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.DoubleWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; public class RunMapReduceJob { public static void main(String[] args) throws Exception { new RunMapReduceJob().run(args); } public void run(String[] args) throws Exception { // 先校验参数数量 if(args.length !=4){ System.err.println("参数错误,请传入4个参数:输入路径 job1输出路径 job2输出路径 job3输出路径"); System.exit(1); } Configuration conf = new Configuration(); FileSystem fs = FileSystem.get(conf); // 自动清理所有已存在的输出路径 Path[] outputPaths = {new Path(args[1]), new Path(args[2]), new Path(args[3])}; for(Path p : outputPaths){ if(fs.exists(p)){ fs.delete(p, true); } } //job 1 Job job1 = Job.getInstance(conf, "hourly"); job1.setJarByClass(RunMapReduceJob.class); job1.setMapperClass(MaxConsumptionMapper.class); job1.setReducerClass(MaxConsumptionReducer.class); job1.setOutputKeyClass(Text.class); job1.setOutputValueClass(DoubleWritable.class); FileInputFormat.addInputPath(job1, new Path(args[0])); FileOutputFormat.setOutputPath(job1, new Path(args[1])); // 校验job1执行结果,失败直接退出 if(!job1.waitForCompletion(true)){ System.err.println("job1执行失败"); System.exit(1); } //job 2 Job job2 = Job.getInstance(conf, "max hourly"); job2.setJarByClass(RunMapReduceJob.class); job2.setMapperClass(MaxConsumptionMapper2.class); job2.setReducerClass(MaxConsumptionReducer2.class); job2.setOutputKeyClass(Text.class); job2.setOutputValueClass(DoubleWritable.class); FileInputFormat.addInputPath(job2, new Path(args[1])); FileOutputFormat.setOutputPath(job2, new Path(args[2])); if(!job2.waitForCompletion(true)){ System.err.println("job2执行失败"); System.exit(1); } //job 3 Job job3 = Job.getInstance(conf, "Avg Daily Consumption"); job3.setJarByClass(RunMapReduceJob.class); job3.setMapperClass(AvgConsumptionMapper.class); job3.setReducerClass(AvgComsumptionReducer.class); job3.setOutputKeyClass(Text.class); job3.setOutputValueClass(DoubleWritable.class); FileInputFormat.addInputPath(job3, new Path(args[0])); FileOutputFormat.setOutputPath(job3, new Path(args[3])); System.exit(job3.waitForCompletion(true) ? 0 : 1); } }
内容的提问来源于stack exchange,提问作者PGNewbie
相关产品推荐
相关产品推荐

