Hadoop 3.0 MapReduce如何同时输出Mapper与Reducer结果至HDFS及MySQL
单MapReduce任务实现多输出(Mapper到HDFS+MySQL、Reducer到HDFS)
嘿,这个需求我之前做Hadoop数据处理的时候碰到过,完全不用拆成两次任务来跑!下面给你两个实操性很强的方案,都是基于Hadoop 3.0的特性实现的:
方案一:使用MultipleOutputs实现多路径输出 + Mapper内JDBC写入MySQL
MultipleOutputs是Hadoop官方提供的多输出工具,支持在Mapper和Reducer阶段同时输出到不同路径,刚好能满足你把Mapper结果单独存HDFS的需求,再配合JDBC在Mapper里直接写MySQL,一步到位。
具体步骤:
Job配置阶段注册多输出
初始化Job时,除了设置常规的Reducer输出,还要注册Mapper的额外输出路径:Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "MultiOutputTask"); // 配置Reducer的常规输出(到HDFS) job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); job.setOutputFormatClass(TextOutputFormat.class); FileOutputFormat.setOutputPath(job, new Path("/hdfs/reducer-result")); // 注册Mapper的额外输出(单独存到HDFS的指定路径) MultipleOutputs.addNamedOutput(job, "mapper-hdfs-output", TextOutputFormat.class, Text.class, IntWritable.class); MultipleOutputs.setCountersEnabled(job, true);Mapper类实现多输出逻辑
在Mapper的setup()里初始化JDBC连接和MultipleOutputs,map()方法同时完成三件事:给Reducer传数据、把Mapper结果写HDFS、把数据插入MySQL,最后在cleanup()里关闭所有资源:public class MultiTaskMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private MultipleOutputs<Text, IntWritable> multipleOutputs; private Connection jdbcConn; private PreparedStatement insertStmt; @Override protected void setup(Context context) throws IOException, InterruptedException { // 初始化多输出工具 multipleOutputs = new MultipleOutputs<>(context); // 初始化MySQL连接(建议用连接池优化,避免频繁创建连接) try { Class.forName("com.mysql.cj.jdbc.Driver"); jdbcConn = DriverManager.getConnection("jdbc:mysql://localhost:3306/your_db", "username", "password"); insertStmt = jdbcConn.prepareStatement("INSERT INTO mapper_data (key_col, value_col) VALUES (?, ?)"); } catch (Exception e) { throw new IOException("JDBC初始化失败", e); } } @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 模拟业务逻辑:拆分输入数据得到键值对 String[] dataParts = value.toString().split("\t"); Text outputKey = new Text(dataParts[0]); IntWritable outputVal = new IntWritable(Integer.parseInt(dataParts[1])); // 1. 正常输出给Reducer context.write(outputKey, outputVal); // 2. 把Mapper结果写入HDFS的指定路径 multipleOutputs.write("mapper-hdfs-output", outputKey, outputVal, "mapper-result/part"); // 3. 插入MySQL数据库 try { insertStmt.setString(1, outputKey.toString()); insertStmt.setInt(2, outputVal.get()); insertStmt.executeUpdate(); } catch (SQLException e) { throw new IOException("MySQL插入失败", e); } } @Override protected void cleanup(Context context) throws IOException, InterruptedException { // 关闭资源 if (insertStmt != null) { try { insertStmt.close(); } catch (SQLException ignored) {} } if (jdbcConn != null) { try { jdbcConn.close(); } catch (SQLException ignored) {} } multipleOutputs.close(); } }Reducer保持常规逻辑
Reducer只需要处理来自Mapper的输入,输出到预先配置的HDFS路径即可,不需要额外修改。
方案二:自定义CompositeOutputFormat(更灵活的输出控制)
如果需要更精细的输出格式控制(比如Mapper输出到HDFS用Parquet格式,Reducer用Text格式),可以自定义CompositeOutputFormat,让它同时管理多个输出器:
核心思路:
- 自定义OutputFormat,重写
getRecordWriter()方法,返回一个复合RecordWriter,包含多个子Writer(比如TextRecordWriter写HDFS、自定义JdbcRecordWriter写MySQL)。 - 在Mapper/Reducer的
write()方法中,给键值对添加标记(比如给键加前缀MAPPER_/REDUCER_),让复合Writer自动把数据分发到对应的输出目标。
关键注意点:
- JDBC写入时,建议用连接池(比如HikariCP)优化性能,避免大量Mapper任务同时创建数据库连接导致压力过大。
- HDFS输出路径要确保有写入权限,可设置为临时路径,任务完成后再移动到最终目标路径,避免中间数据残留。
- 若Mapper并行度很高,要注意调整MySQL的最大连接数,或通过
mapreduce.job.maps参数控制Mapper任务数量。
这样提交一次MapReduce任务,就能同时完成三个输出目标:Mapper结果到HDFS、Mapper结果到MySQL、Reducer结果到HDFS,完全不用跑两次任务!
内容的提问来源于stack exchange,提问作者CryptoMonkey
相关产品推荐
相关产品推荐

