基于Hadoop MapReduce计算累计能耗的最大小时耗电量问题排查
核心问题梳理
你的代码输出全为0.0,根本原因是序列化错误、字符串比较逻辑错误、Iterable遍历滥用以及业务逻辑偏差,下面逐个拆解:
1. Mapper泛型声明错误
你的Mapper类声明是:
public class EnergyMapper extends Mapper<LongWritable, Text, Text, FloatWritable>
但实际输出的键值对是Text和EnergyValues,泛型不匹配会导致Hadoop序列化时出错,数据传递到Reducer时全是无效值。
修复:修改泛型为实际输出类型:
public class EnergyMapper extends Mapper<LongWritable, Text, Text, EnergyValues>
2. EnergyValues序列化/反序列化完全错误
writeChars()是逐个写入字符串的每个字符,而readChar()只读取单个字符,导致time和date在反序列化后只拿到第一个字符(比如05:21:50变成0),后续所有时间、日期比较全失效。
修复:改用writeUTF()和readUTF()处理字符串:
@Override public void write(DataOutput dataOutput) throws IOException { dataOutput.writeInt(houseId); dataOutput.writeUTF(time); // 替换writeChars dataOutput.writeDouble(energyReading); dataOutput.writeUTF(date); // 替换writeChars } @Override public void readFields(DataInput dataInput) throws IOException { houseId = dataInput.readInt(); time = dataInput.readUTF(); // 替换readChar energyReading = dataInput.readDouble(); date = dataInput.readUTF(); // 替换readChar }
另外,EnergyValues(String value2)构造函数未在序列化中处理data字段,建议删除避免混淆。
3. 字符串比较用==而非equals()
Java中==比较对象引用而非字符串内容,比如:
if (key.toString() == amount.getDate()) if (newTime == amount.getTime())
这两个判断永远为false,导致getHourlyComp()找不到匹配数据,直接返回0。
修复:全部改成equals():
if (key.toString().equals(amount.getDate())) if (newTime.equals(amount.getTime()))
4. Reducer中滥用Iterable遍历
MapReduce的Iterable<EnergyValues>是一次性迭代器,在reduce主循环遍历sales后,getHourlyComp()再次遍历同一个sales会直接返回空,拿不到任何数据,所以getHourlyComp()始终返回0。
修复:先把所有数据缓存到集合中:
@Override protected void reduce(Text key, Iterable<EnergyValues> sales, Context context) throws IOException, InterruptedException { // 缓存所有数据到List(注意深拷贝,避免对象复用问题) List<EnergyValues> energyList = new ArrayList<>(); for (EnergyValues ev : sales) { energyList.add(new EnergyValues(ev.getHouseId(), ev.getTime(), ev.getEnergyReading(), ev.getDate())); } double maxConsumption = 0.0; // 按住宅ID分组 Map<Integer, List<EnergyValues>> houseMap = new HashMap<>(); for (EnergyValues ev : energyList) { houseMap.computeIfAbsent(ev.getHouseId(), k -> new ArrayList<>()).add(ev); } // 遍历每个住宅的能耗数据,计算小时耗电量 for (Map.Entry<Integer, List<EnergyValues>> entry : houseMap.entrySet()) { List<EnergyValues> houseEnergy = entry.getValue(); // 按时间排序 houseEnergy.sort(Comparator.comparing(EnergyValues::getTime)); // 按「日期+小时」分组,统计该小时内的最大/最小读数 Map<String, List<Double>> hourlyMap = new HashMap<>(); for (EnergyValues ev : houseEnergy) { String hourKey = ev.getDate() + "_" + ev.getTime().substring(0, 2); hourlyMap.computeIfAbsent(hourKey, k -> new ArrayList<>()).add(ev.getEnergyReading()); } // 计算每个小时的耗电量,更新当前日期的最大值 for (List<Double> readings : hourlyMap.values()) { double hourConsumption = Collections.max(readings) - Collections.min(readings); if (hourConsumption > maxConsumption) { maxConsumption = hourConsumption; } } } // 输出当前日期的最大小时耗电量,后续需全局汇总 context.write(new Text("max"), new FloatWritable((float) maxConsumption)); }
5. 业务逻辑偏差
原代码试图找当前时间加一小时后的读数计算耗电量,但实际数据每10秒一条,很少有刚好差一小时的记录。正确逻辑是:
- 按住宅+日期+小时分组
- 取组内最大能耗读数 - 最小能耗读数,得到该小时耗电量
- 遍历所有组,找到最大值
6. 最终输出单个数值
原Reducer会输出多个日期的结果,需添加全局Reducer汇总最大值:
public class GlobalMaxReducer extends Reducer<Text, FloatWritable, Text, FloatWritable> { @Override protected void reduce(Text key, Iterable<FloatWritable> values, Context context) throws IOException, InterruptedException { double globalMax = 0.0; for (FloatWritable val : values) { if (val.get() > globalMax) { globalMax = val.get(); } } context.write(new Text("final_max_hourly_consumption"), new FloatWritable((float) globalMax)); } }
Job配置中需设置两级Reducer,将第一级输出作为第二级输入,最终得到单个全局最大值。
修复后流程
- Mapper输出
日期为键,EnergyValues(含住宅ID、时间、能耗读数、日期)为值 - 第一级Reducer按日期分组,缓存数据后按住宅+小时计算每小时耗电量,输出
"max"为键、该日期最大小时耗电量为值 - 第二级GlobalMaxReducer汇总所有日期的最大值,输出最终单个数值
内容的提问来源于stack exchange,提问作者coder12345

