Scala实现Hadoop MapReduce键拆分与JSON格式聚合输出
解决方案
要实现你需要的输出格式,需要对Reducer做以下几处关键修改:
- 调整输出类型:原Reducer输出值为
IntWritable,现在需要输出格式化后的字符串,因此泛型要改为Reducer[Text, IntWritable, Text, Text] - 拆分输入键:将输入的
key(如"2:40 INFO")拆分为时间前缀("2:40")和日志级别("INFO") - 收集同时间下的所有级别数据:用Scala的
Map存储同一时间对应的各个日志级别及其数值 - 格式化输出字符串:将收集到的键值对整理成目标格式的字符串
修改后的完整Reducer代码
import org.apache.hadoop.io.{IntWritable, Text} import org.apache.hadoop.mapred.{MapReduceBase, OutputCollector, Reducer, Reporter} import java.util import scala.collection.mutable class Reduce extends MapReduceBase with Reducer[Text, IntWritable, Text, Text] { override def reduce( key: Text, values: util.Iterator[IntWritable], output: OutputCollector[Text, Text], reporter: Reporter ): Unit = { // 拆分输入键,处理多个连续空格的情况 val parts = key.toString.split("\\s+") val timePrefix = parts(0) // 提取时间部分,如"2:40" val logLevel = parts(1) // 提取日志级别,如"INFO" // 获取当前级别对应的数值(输入中每个级别对应单个值,直接取第一个即可) val value = values.asScala.next().get() // 用Reporter上下文存储同时间的级别数据,避免Reducer复用导致状态混乱 val contextKey = s"time_${timePrefix}" val levelMap = Option(reporter.getCounter("ReducerCache", contextKey).getValue.toString) .map(str => mutable.Map(str.split(", ").map(p => { val kv = p.split(" : ") kv(0) -> kv(1).toInt }): _*)) .getOrElse(mutable.Map.empty[String, Int]) // 将当前级别和数值存入Map levelMap.put(logLevel, value) // 更新上下文存储的Map字符串 val mapStr = levelMap.map { case (k, v) => s"$k : $v" }.mkString(", ") reporter.getCounter("ReducerCache", contextKey).setValue(mapStr) // 当处理到当前时间分组的最后一个键值对时,统一输出 if (!values.hasNext) { val outputValue = new Text(s"{${mapStr}}") output.collect(new Text(timePrefix), outputValue) // 清空上下文缓存,避免影响下一个时间的处理 reporter.getCounter("ReducerCache", contextKey).setValue("") } } }
关键细节说明
- 键拆分处理:使用
split("\\s+")而非split(" "),可以兼容输入中多个连续空格的情况,保证拆分结果准确。 - 状态存储:利用
Reporter的计数器上下文存储同一时间下的级别数据,解决Reducer实例被Hadoop复用导致的状态丢失问题。 - 输出时机控制:通过
!values.hasNext判断是否为当前时间分组的最后一个键值对,此时统一输出完整的结果,避免重复输出部分数据。
可选优化:调整Map阶段输出(更简洁的实现)
如果允许修改Map阶段的输出,将键设为时间前缀(如"2:40"),值设为Text类型的"INFO:4",Reducer可以大幅简化:
import org.apache.hadoop.io.Text import org.apache.hadoop.mapred.{MapReduceBase, OutputCollector, Reducer, Reporter} import java.util import scala.collection.mutable class Reduce extends MapReduceBase with Reducer[Text, Text, Text, Text] { override def reduce( key: Text, values: util.Iterator[Text], output: OutputCollector[Text, Text], reporter: Reporter ): Unit = { val levelMap = mutable.Map.empty[String, Int] values.asScala.foreach { value => val parts = value.toString.split(":") levelMap.put(parts(0), parts(1).toInt) } val outputStr = levelMap.map { case (k, v) => s"$k : $v" }.mkString("{ ", ", ", " }") output.collect(key, new Text(outputStr)) } }
这种方式更直观,不需要处理键拆分和状态保留,推荐优先采用。
内容的提问来源于stack exchange,提问作者user20093716
相关产品推荐
相关产品推荐

