ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Hadoop MapReduce实战:气象数据年最高气温统计完整指南

Hadoop MapReduce实战:气象数据年最高气温统计完整指南 简介这是一套基于HDFS与MapReduce的气象数据处理完整项目代码适合大数据初学者、Java开发人员及课程设计者用于实战参考。项目围绕气象数据的分布式存储与并行分析展开Map阶段完成原始数据的读入、清洗和键值对转换Reduce阶段负责聚合计算区域平均气温、最高最低温度、降雨量等统计指标整个流程体现了大数据处理中典型的拆分—聚合思想。同时资源整合了SSM框架将处理结果通过Web页面进行可视化展示并给出了清晰的前后端交互接口还附有war包、数据库脚本与配置文件便于快速部署和运行。压缩包共562个文件涵盖Java源码、编译后的class文件、第三方依赖jar包、SQL初始化脚本、jsp/css/js前端资源、xml/properties配置文件以及png/jpg等图片素材整体大小约34.88MB目录按功能模块划分可快速定位到Hadoop核心代码或前端页面。目前已有10563人学习下载读者既能对照源码理解MapReduce作业的编写与执行流程也能参考class文件验证运行结果还可借鉴SSM整合Hadoop的接口封装与页面渲染方式掌握大文件上传、解析及结果展示等常见问题的处理思路适合作为大数据实习、课程设计或实际项目中的参考也可为后续扩展时间序列分析与预测功能提供基础。 讲个真实的事情你随便搜Hadoop 课程设计 题目十个里有八个都带气象数据。为什么因为气象数据天生就是 MapReduce 的教科书案例——数据量大、格式规整、按年份聚合的统计语义非常清晰。我在带新人练手和做面试准备时也一直推荐这个场景一套代码跑通MapReduce 的 Map、Shuffle、Reduce 全流程你基本就吃透了。这篇文章我会直接给你一版真正能跑通的完整代码从环境怎么搭、数据长什么样、每一行代码为什么这么写到打包、上传、跑任务、看结果的完整操作再到我踩过的几个坑和优化思路都写清楚。适合正在做 Hadoop 课程设计的人也适合准备面试想拿 MapReduce 练手的人。1. 项目定位为什么是气象数据为什么用 Hadoop1.1 气象数据的“高吞吐”特点先明确一点我们用 Hadoop 分析气象数据不是因为它炫技而是因为这个场景的底层逻辑和 Hadoop 的设计目标完全吻合。气象数据有几个典型特征第一数据量极大一个全球范围的历史气象观测数据集轻松到 GB 甚至 TB 级单机内存根本hold不住第二数据格式高度固定每一条记录都是定长字段非常适合做字段截取解析第三计算逻辑相对简单大多数场景就是按时间年、月做极值统计、均值聚合这种“先分组再计算”的需求就是 MapReduce 的看家本领。你可能要问这年头 Spark 那么火为什么不用 Spark答案很简单——Hadoop 生态里HDFS MapReduce 是底层基础。你先把 MapReduce 的编程模型搞明白再去看 Spark 的 RDD、DataFrame 那套东西会发现很多概念都是相通的比如 map、reduceByKey。而且面试官问 Hadoop 八股的时候MapReduce 一定是重灾区不会手写 Mapper 和 Reducer光背理论是撑不过去的。1.2 项目要解决的“最小核心问题”我们这个项目的核心目标可以浓缩成一句话从海量气象观测记录中统计出每一年的最高气温Max Temperature。听起来很简单但它完整覆盖了几个关键步骤从 HDFS 上读取原始数据文件在 Mapper 阶段解析每一行文本提取出“年份”和“气温”两个字段在 Shuffle 阶段让相同年份的记录自动汇聚到同一个 Reducer在 Reducer 阶段比较该年份所有气温值输出最大值。这个链路跑通之后你想改成求最低气温、平均气温、按月份统计、按站点统计都只是改几行代码的事。也就是说这个项目是一个“骨架”你在这个骨架上加肉就是各种变体。2. 环境准备与数据选型2.1 Hadoop 环境伪分布式完全够用很多新手上来就想着搭一个三节点、五节点的集群结果光调网络和免密登录就耗了两天。我的建议是跑这种离线批处理 demo伪分布式完全够用。所谓伪分布式就是在一台机器上同时启动 NameNode、DataNode、ResourceManager、NodeManager 这几个守护进程HDFS、YARN 的功能都具备适合开发和验证。只有当你需要测数据分布、带宽、节点容错这些集群特性时才需要上多节点。搭建要点我简单梳理一下这里不展开所有细节但核心配置要清楚core-site.xmlconfiguration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configurationhdfs-site.xmlconfiguration property namedfs.replication/name value1/value /property !-- 伪分布式副本数必须设为1否则多个DataNode不在同一台机器时会报副本不足 -- /configurationyarn-site.xmlconfiguration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property /configuration关键一步是启动前要格式化 NameNodehdfs namenode -format然后分别启动 HDFS 和 YARNstart-dfs.sh start-yarn.sh用jps命令如果能看到 NameNode、DataNode、ResourceManager、NodeManager 四个进程说明环境就绪了。这一步卡住的人最多后面专门讲。2.2 气象数据格式先把 NCDC 格式啃明白我们要处理的数据最经典的是 NCDC美国国家气候数据中心的公开数据集。这种数据每一行是一条定长记录看起来像天书002902907099999190101010600464333023450FM-12000599999V0202701N015919999999N0000001N9-0078199999102001ADDGF108991999999999999999999但你只需要记住我们要用的两个字段年份位于第 15~19 个字符按 0 开始的下标就是substring(15, 19)气温位于第 87~92 个字符其中第 88 个字符是正负号标志或-后面是四位数值实际温度要除以 10。我用 Python 模拟了几行带不同年份的数据方便你验证效果lines [ 002902907099999190101010600464333023450FM-12000599999V0202701N015919999999N0000001N9-0078199999102001ADDGF108991999999999999999999, 002902907099999190201010600464333023450FM-12000599999V0202701N015919999999N0000001N90021199999102001ADDGF108991999999999999999999, 002902907099999200301010600464333023450FM-12000599999V0202701N015919999999N0000001N90085299999102001ADDGF108991999999999999999999 ] for line in lines: year line[15:19] # 温度在 87~92符号位 87 temp_str line[87:92] if temp_str.startswith(): temp int(temp_str[1:]) / 10.0 elif temp_str.startswith(-): temp -int(temp_str[1:]) / 10.0 else: temp None print(year, temp)运行结果就是第二年、第三年的摄氏度气温值。实际作业里数据量大了你不可能手动解析这活就是要交给 Mapper 去干。3. 核心代码实现与思路拆解3.1 Mapper从一行原始数据里“抠”出关键字段Mapper 的职责很简单从原始文本行中提取我们需要的年份和气温输出为年份, 气温的键值对。但“简单”不等于“随便写”。这里有三个容易忽视的细节第 87 位置的符号位表示正温-表示负温Java 的Integer.parseInt能处理负号却不能处理正号所以必须分开截取气温字段如果出现 9999表示数据缺失这种脏数据要在 Mapper 里过滤掉气温暖数据末尾还有一个质量码字段表示该记录是否可信通常0、1、4、5、9是可用数据后面的异常数据不要进入统计。直接上代码package com.example.weather; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.IOException; public class WeatherMapper extends MapperLongWritable, Text, Text, IntWritable { private static final int MISSING 9999; Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); if (line null || line.trim().isEmpty()) { return; } // 年份第15~19个字符Java下标 15~18 String year line.substring(15, 19); // 温度第87~92个字符第87位是符号位 int airTemperature; if (line.charAt(87) ) { airTemperature Integer.parseInt(line.substring(88, 92)); } else { airTemperature Integer.parseInt(line.substring(87, 92)); } // 过滤缺失值和异常质量码 String quality line.substring(92, 93); if (airTemperature ! MISSING quality.matches([01459])) { context.write(new Text(year), new IntWritable(airTemperature)); } } }这里有个小细节值得说明为什么substring(15, 19)而不是substring(16, 20)因为 Java 的substring是左闭右开substring(15, 19)取的是第 16、17、18、19 这四个字符对应数据格式里的第 15~19 个字符如果从 0 开始数。新手最容易在这类偏移量上翻车我的习惯是把一行原始数据先打印出来数好位置再写代码而不是靠记忆。3.2 Reducer按年份聚合求极值Reducer 的逻辑比 Mapper 还简单拿到同一个年份下的一串气温值遍历一遍求最大。但如果只是写出这个逻辑你会漏掉一个非常经典的优化点——Combiner。Combiner 是在 Map 端做的“预聚合”。它的本质是一个运行在 Mapper 输出端的“小型 Reducer”先把每个 Map 任务内部的重复按键合并一部分从而减少 Shuffle 阶段要传输的数据量。像“求最大值”这种操作是满足交换律和结合律的完全可以安全地用 Combiner 预聚合而且不会影响最终结果。package com.example.weather; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; public class MaxTemperatureReducer extends ReducerText, IntWritable, Text, IntWritable { Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int maxValue Integer.MIN_VALUE; for (IntWritable value : values) { if (value.get() maxValue) { maxValue value.get(); } } context.write(key, new IntWritable(maxValue)); } }这个 Reducer 同时兼做 Combiner——你只需要在 Driver 里同时设置job.setCombinerClass(MaxTemperatureReducer.class)Hadoop 就会在 Map 端自动调用它。这也是为什么我把它单独写成一个类而不是匿名内部类因为要复用。3.3 Driver串联整个作业Driver 是整个作业的“组装车间”负责配置 Job 的各种参数输入输出路径、Mapper/Reducer/Combiner 类、输出键值类型等。这里的核心要点有两个一是输出键值类型必须和 Reducer 的输出类型一致。很多人 Mapper 输出的是Text和IntWritableReducer 输出类型也是这两个但如果中间设置不一致运行时会直接报类型不匹配的错。二是输出路径不能预先存在。Hadoop 会直接报FileAlreadyExistsException这是新手高频错误没有之一。我看到过不少人在命令行里反复执行同一任务结果一直卡在这一步。package com.example.weather; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; 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 MaxTemperature { public static void main(String[] args) throws Exception { if (args.length ! 2) { System.err.println(Usage: MaxTemperature input path output path); System.exit(-1); } Configuration conf new Configuration(); Job job Job.getInstance(conf, Max Temperature); job.setJarByClass(MaxTemperature.class); job.setMapperClass(WeatherMapper.class); job.setCombinerClass(MaxTemperatureReducer.class); job.setReducerClass(MaxTemperatureReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }有几个额外的配置项我想提醒你留意。如果你的数据文件特别多比如上千个小文件可以给 Job 设置setNumReduceTasks来指定 Reducer 数量默认一个 Reducer 处理所有分组数据数据量大时会有性能问题。另外如果你对输入文件的压缩格式有要求比如 gz、snappyHadoop 会根据扩展名自动判断解压方式不需要额外代码但如果你用的是自定义压缩格式就需要在 Configuration 里手动声明了。4. 完整实操从上传数据到跑出结果4.1 准备数据并上传到 HDFS代码写完之后先不要急着打包。第一步是准备数据。你可以从 NCDC 官网下载真实数据也可以像我一样自己拼一些测试数据关键是字段位置要对。新建一个文本文件比如weather.txt每一行都按 NCDC 格式来。我上面给出的三行样本数据就可以直接用其中包含了 1901 年、1902 年、2003 年三个年份的不同气温。然后把数据上传到 HDFS# 创建输入目录 hdfs dfs -mkdir -p /weather/input # 上传数据 hdfs dfs -put weather.txt /weather/input/ # 确认上传成功 hdfs dfs -ls /weather/input这里顺带提一个新手容易犯的错误HDFS 的根目录是/user/你的用户名如果你直接-put到相对路径文件可能并不是你想放的地方。建议上传前先用hdfs dfs -ls /看清楚目录结构避免后面路径对不上。4.2 编译打包、提交作业、验证结果我用 Maven 管理项目pom.xml里引入 Hadoop 依赖后直接打包。如果你问我为什么用 Maven 而不是手动javac——第一Hadoop 生态的依赖关系复杂手动引 jar 包容易漏第二Maven 打包时可以直接生成可运行的 fat jar省去一堆 classpath 配置的烦恼。关键的pom.xml依赖如下版本号根据你本机 Hadoop 版本调整dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version3.3.6/version /dependency打包mvn clean package -DskipTests然后提交到 Hadoophadoop jar target/weather-demo-1.0-SNAPSHOT.jar \ com.example.weather.MaxTemperature \ /weather/input /weather/output跑的时候你会看到 MapReduce 的进度信息从map 0%到reduce 100%。第一次跑通的时候这个过程其实是有点激动的因为它证明了你写的代码真的在分布式框架上跑起来了。跑完之后查看结果hdfs dfs -cat /weather/output/part-r-00000期望输出类似1901 -78 1902 21 2003 85注意这里输出的气温是原始值的整数形式也就是实际温度乘以 10因为我们在 Mapper 里没有除以 10。如果你想直接输出摄氏度可以在 Reducer 输出前转成 FloatWritable 处理回顾一下我在数据格式那节写的 Python 解析代码逻辑是一样的。5. 常见问题排查与性能优化经验5.1 运行中一定会遇到的六个坑我一个一个列出来很多都是我在实操中被折磨过才总结出来的。现象原因解决办法作业提交后报FileAlreadyExistsException输出目录已存在每次运行前删掉旧输出目录或换个新目录名jps看不到 DataNode且日志显示 clusterID 不一致格式化 NameNode 后DataNode 的 clusterID 和 NameNode 不一致停掉所有进程删除dfs.namenode.name.dir和dfs.datanode.data.dir下的数据重新namenode -format任务卡在map 100% reduce 0%很久有数据倾斜或 Reducer 数量太少检查 key 分布是否极端比如所有数据都压在同一个年份上适当增加 Reducer 数NumberFormatException某行数据格式异常字段位置不对在 Mapper 里加 try-catch 或字段长度校验先定位是哪一行导致网页访问不了 NameNode 管理界面防火墙未关闭或端口不对Hadoop 3.x 默认 9870确认端口号关闭防火墙或用curl localhost:9870测试日志显示Container ... is running beyond physical memory limitsYARN 给容器分配的内存不够在yarn-site.xml里调大yarn.nodemanager.resource.memory-mb和yarn.scheduler.maximum-allocation-mb重启 YARN关于 clusterID 不一致这个问题我再展开说两句。这是伪分布式新手最容易遇到的怪现象第一次格式化 NameNode 后一切正常但如果你手贱重新执行了hdfs namenode -formatDataNode 启动时就会报错因为 NameNode 的 clusterID 变了而 DataNode 的本地存储中还保留着旧的 clusterID。解决办法不是再格式化一次而是把 NameNode 和 DataNode 的数据目录都清掉然后重新格式化、重新启动。所以格式化这个操作能不做就不做做了就别后悔。5.2 三个真正有价值的优化思路当你把基础流程跑通可以试着做下面几件事对你的理解和面试帮助都很大。第一个优化是给作业设置合理的 Reducer 数量。默认情况下 Hadoop 只启动一个 Reducer 处理所有分组数据数据量一旦上来单个 Reducer 就成了瓶颈。你可以在 Driver 中加入job.setNumReduceTasks(4);这样输出结果会变为多个part-r-00000到part-r-00003的文件每个文件对应一个 Reducer 的输出。具体数值不是越大越好通常与数据量和集群核数有关。第二个优化是做一个自定义Writable数据类型。假设你要同时统计某年的最高气温、最低气温、平均气温、记录条数你可以在 Mapper 端把所有相关信息封装成一个自定义对象输出Reducer 端只处理这个对象。虽然写起来比IntWritable复杂但这是 Hadoop 面试中如何在 MapReduce 中传递复合类型的标准答案。第三个优化是数据倾斜的预处理。如果数据按年份倾斜严重比如某一年数据特别多可以在 Mapper 输出的 key 上加上随机前缀把原本集中在一个 Reducer 的压力打散到多个 Reducer然后在下一个 Job 中做二次聚合。这种两阶段聚合的思路在 Flink 和 Spark 里也是通用解决方案。6. 通过调试日志看 MapReduce 的执行内幕很多新手跑完任务后只知道“出结果了”但对 MapReduce 内部发生了什么一无所知。我建议你学会看任务日志这是从“会用”到“懂用”的关键一步。在作业运行过程中你会看到类似这样的日志行INFO mapreduce.Job: map 0% reduce 0% INFO mapreduce.Job: map 100% reduce 0% INFO mapreduce.Job: map 100% reduce 100% INFO mapreduce.Job: Job job_1700000000000_0001 completed successfully如果你用的是 Hadoop 3.x 且有 YARN 的 ResourceManager 界面可以打开http://localhost:8088在 Application 列表中找到你运行的作业然后点进 Logs 查看每个 Map 任务和 Reduce 任务的详细日志。这里你能看到每个任务处理了多少条记录Map input records...、Shuffle 传输了多少字节Shuffle Bytes...等关键指标。为什么要强调这个因为在面试或调优时面试官最喜欢问“你的任务有多少 Map、多少 Reduce为什么这么分配”。如果你能从日志中直接读出每个 Map 处理的字节数、每个 Reduce 处理的分组数回答会非常有力。这也是演示你“真的跑过任务”而不是只会背八股的最好方法。有一次我调一个真实业务任务发现 90% 的时间都花在 Shuffle 阶段点开日志一看发现 Mapper 输出非常大而 Reducer 数只有 1 个所有数据都往一个节点上汇聚不慢才怪。后来把 Reducer 数调到 8加上 Combiner整个任务耗时直接下降了 60%——这种排障经验你在书本上很难学得到。最后再分享一点个人体会。如果你是自己练习建议把这个项目拆成两个阶段来跑第一遍不加 Combiner第二遍加上 Combiner对比两次的 Shuffle 数据量和运行时间。你会发现一个 Combiner 带来的提升比你想象中大得多。这种“对比实验”式的学习方法比单纯照着代码敲一遍要有效得多。等你把这个流程玩熟再去尝试按月份分组统计、输出多个统计指标、用自定义 Writable 传递复合数据你会发现自己已经不是在背 Hadoop 了而是真的在用它解决业务问题了。本文还有配套的精品资源点击获取
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进