You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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("")
    }
  }
}

关键细节说明

  1. 键拆分处理:使用split("\\s+")而非split(" "),可以兼容输入中多个连续空格的情况,保证拆分结果准确。
  2. 状态存储:利用Reporter的计数器上下文存储同一时间下的级别数据,解决Reducer实例被Hadoop复用导致的状态丢失问题。
  3. 输出时机控制:通过!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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 22:01:17