简介:这是一份基于Hadoop的疾病信息统计平台毕业设计项目,面向计算机、通信、人工智能等专业的学生和从业者,适合作为课程大作业或毕设参考。系统围绕疾病数据采集、存储与统计分析场景,完整提供源代码及配套文档说明,代码经调试测试可运行,可用于理解Hadoop生态在实际业务中的落地流程。压缩包共41个文件,主要包含25个Java源码、6个XML配置、2个properties配置文件、2个JAR依赖包,以及YAML、ARFF数据集、Maven脚本等,整体约10.87MB,目录结构清晰便于定位核心模块。目前已有104人学习下载,适合需要快速搭建同类统计平台或完成期末项目的读者。项目具备较高借鉴价值,基础较好的开发者可在现有结构上扩展数据可视化或算法分析功能。
1. 基于Hadoop的疾病信息统计平台:毕业设计为什么都在选这个方向
每到毕业设计选题季,基于Hadoop的疾病信息统计平台这类题目总被反复拿出来。不是因为医院数据真有多大规模,而是它把HDFS存储、MapReduce计算、YARN调度这一整套Hadoop核心机制,落到了一个需求足够明确、数据足够可控的场景里:把一批疾病记录文件放进分布式文件系统,按病种、地区、年龄等维度跑统计作业,再把结果做成图表展示。整套链路开题能写、中期有进度、答辩有实物。适合想稳妥完成毕业设计的学生,也适合想快速了解Hadoop全流程的开发者照着本文去复现。源代码加文档说明的交付形式,也让这个题目天然适合作为课程设计和本科毕设的蓝本。
2. 先跑通Hadoop伪分布式:从零开始的核心配置与项目目录规划
2.1 伪分布式还是完全分布式:为什么毕业设计只够用伪分布式
先回答最常被问的问题:这个题目要不要搭三台机器的完全分布式集群。我的观点很直接,不要。本科毕业设计的时间线通常只有三到四个月,其中还要写文档、做PPT、准备答辩,花两周去调集群网络和节点同步,性价比太低。伪分布式模式下,NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager各自独立进程都跑在同一台机器上,HDFS的块、副本、YARN调度这些核心机制一个不少,对理解分布式原理完全够用。
完全分布式和伪分布式的差别要能说清楚,这是答辩老师必问的点。完全分布式下数据块会按 replica 参数分布到不同节点,伪分布式下所有副本都在本机,并不具备真正的故障容错能力。所以伪分布式环境里副本数建议直接设成1,既省磁盘,也不会在演示的时候因为副本同步问题出幺蛾子。如果只是想让文档里多一张“三副本”的截图,临时把 dfs.replication 改成 2 再跑一次上传,截图完改回来,不必长期开着。
内存方面,伪分布式对笔记本的配置要求不高,4G内存能跑,8G会舒服很多。我一般建议用 Apache Hadoop 3.3.x 的稳定版本配 JDK 1.8,这两个版本组合经过大量实践验证,yarn 和 mapreduce 的兼容告警最少。不要为了追新用 JDK 17 或更高版本,跑 YARN 时会出现一堆序列化相关的兼容问题,属于给自己找事。
2.2 五份核心配置文件与启动命令:能跑起来的完整版本
Hadoop 安装包解压后,所有配置都在 etc/hadoop 目录下。核心要改的文件是 hadoop-env.sh、core-site.xml、hdfs-site.xml、yarn-site.xml、mapred-site.xml 这五个。先从环境变量开始。
# hadoop-env.sh 中必须显式指定 JAVA_HOME # 不要指望它自动读取系统变量,不同发行版的行为不一致 export JAVA_HOME=/usr/lib/jvm/java-1.8.0-openjdk export HADOOP_HOME=/opt/hadoop-3.3.6 export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin逻辑说明:JAVA_HOME 要写绝对路径,指向实际 JDK 安装位置。很多第一次搭建的人在这里省略,导致后续 start-dfs.sh 报错找不到 java 命令。
<!-- core-site.xml:指定 NameNode 地址 --> <configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/opt/hadoop/data/tmp</value> </property> </configuration>hadoop.tmp.dir 是 NameNode 和 DataNode 存放元数据与数据块的根目录,必须改出默认的 /tmp,否则系统重启一次,你的 HDFS 数据就没了。这是伪分布式下最容易被忽略的坑。
<!-- hdfs-site.xml:单机模式关键配置 --> <configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>file:///opt/hadoop/data/nn</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>file:///opt/hadoop/data/dn</value> </property> </configuration>把 namenode 和 datanode 的数据目录分开指定,后续出问题需要重置时,删目录也删得干净。
<!-- yarn-site.xml:单机跑 MR 的关键 --> <configuration> <property> <name>yarn.nodemanager.aux-services</name> <value>mapreduce_shuffle</value> </property> <property> <name>yarn.nodemanager.resource.memory-mb</name> <value>2048</value> </property> <property> <name>yarn.nodemanager.vmem-check-enabled</name> <value>false</value> </property> </configuration>vmem-check-enabled 设成 false 是伪分布式血泪经验。物理内存只有 4G 到 8G 的机器上,YARN 默认虚拟内存检查经常误杀 Container,关掉它至少不会让你在演示时突然看到作业失败。
<!-- mapred-site.xml:指定计算框架走 YARN --> <configuration> <property> <name>mapreduce.framework.name</name> <value>yarn</value> </property> </configuration>配置完成后,首次启动顺序是固定的:
# 1. 格式化 NameNode,只在第一次启动前执行 hdfs namenode -format # 2. 启动 HDFS 与 YARN start-dfs.sh start-yarn.sh # 3. 验证五个进程是否都在 jpsjps 输出里应该看到 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager 这五个进程,少一个都不正常。进程没起来时,第一时间去 $HADOOP_HOME/logs 目录看对应日志文件,绝大部分问题日志里都写得很明白。
2.3 源代码与文档目录怎么规划:从src到docs都有讲究
基于Hadoop的疾病信息统计平台既然带了源代码和文档说明,项目目录结构最好一开始就按可交付的标准组织,避免答辩前到处找文件。
disease-stat-platform/ ├── bin/ # 启动、上传、统计的一键脚本 │ ├── start_all.sh │ ├── upload_data.sh │ └── run_jobs.sh ├── conf/ # 项目级配置 ├── input/ # 样例疾病数据(CSV) ├── output/ # MapReduce 统计结果 ├── src/main/java/ # 源码 │ └── com/disease/stat/ │ ├── mapreduce/ # 各维度统计作业 │ ├── etl/ # 数据清洗相关 │ └── web/ # 可视化后端接口 └── docs/ # 毕业设计文档 ├── 需求分析.md ├── 概要设计.md ├── 详细设计.md └── 测试报告.md源代码放 src,文档放 docs,数据放 input,结果放 output,脚本放 bin。这个结构的好处是每个目录职责单一,文档里写“系统模块划分”时可以直接引用目录名,答辩老师问你代码在哪,打开目录就能讲。bin 下的脚本建议从第一天就开始写,每次手动执行的命令都要沉淀成脚本,后期跑演示时会省非常多时间。
3. 把疾病数据装进HDFS:字段设计、清洗与上传的完整流程
3.1 疾病记录表怎么设计:字段、样例数据与保留原始字段的原因
疾病信息统计平台的数据源是结构化的疾病记录。设计表结构时有一个原则:原始表只做标准化存储,不提前做聚合。统计是 MapReduce 的事,原始数据里字段越全,后续按病种、地区、年龄、时间切片就越灵活。
| 字段名 | 类型 | 样例值 | 说明 |
|---|---|---|---|
| case_id | string | C0001 | 病例唯一标识,用于去重 |
| disease | string | 高血压 | 病种名称,统计的主维度之一 |
| gender | string | 男 | 性别 |
| age | int | 62 | 年龄,清洗时转数值 |
| region | string | 杭州市 | 地区维度 |
| diagnose_date | string | 2024-03-12 | 诊断日期,可再做月份统计 |
| status | string | confirmed | 记录状态:confirmed/treating/cured |
样例数据先造出一份,数据量不用大,五六十条能覆盖各维度即可,跑通流程后再决定要不要扩到上万条:
case_id,disease,gender,age,region,diagnose_date,status C0001,高血压,男,62,杭州市,2024-03-12,confirmed C0002,2型糖尿病,女,54,宁波市,2024-03-13,confirmed C0003,急性上呼吸道感染,男,28,温州市,2024-03-14,cured C0004,高血压,女,71,绍兴市,2024-03-14,confirmed C0005,冠心病,男,66,嘉兴市,2024-03-15,treating C0006,哮喘,女,35,湖州市,2024-03-16,confirmed注意第一行是表头,MapReduce 作业里要做一次跳过处理,不然会把“disease”也统计进去。âge 字段保留为字符串还是转成 int,取决于你要不要做年龄段分组,后面按年龄段统计时会按 int 计算,建议清洗阶段就完成转换。
3.2 先用Python做ETL再进HDFS:清洗脚本与编码处理
数据进 HDFS 之前必须做一轮清洗。常见做法是先用 Python 脚本在本地完成 ETL,再把清洗后的 CSV 上传,而不是把脏数据直接丢给 MapReduce 去处理。原因很简单:清洗逻辑写一遍就能复用,而且能在上传前直观看到数据质量报告;如果上传后再发现脏数据,你得重跑整个流程,HDFS 又变成一个黑匣子,排查成本高得多。
# etl_clean.py # 清洗原始疾病数据:去BOM、去空值、去异常年龄、统一编码 import pandas as pd raw_path = "input/disease_records_raw.csv" clean_path = "input/disease_records.csv" df = pd.read_csv(raw_path, encoding="utf-8-sig", dtype=str, na_filter=False) # 去空格:病种和地区经常混入肉眼不可见的空白字符 df["disease"] = df["disease"].str.strip() df["region"] = df["region"].str.strip() # 年龄转数值,非法值直接剔除 df["age"] = pd.to_numeric(df["age"], errors="coerce") df = df[df["age"].notna()] df["age"] = df["age"].astype(int) # 病种为空的记录剔除 df = df[df["disease"] != ""] # 统一写出为 UTF-8,保持字段顺序 df.to_csv(clean_path, index=False, encoding="utf-8") print(f"清洗完成:原始 {len(pd.read_csv(raw_path, encoding='utf-8-sig'))} 行,剩余 {len(df)} 行")逻辑说明:read_csv 用 encoding="utf-8-sig" 是为了自动吃掉 UTF-8 BOM 头。Windows 下用 Excel 另存的 CSV 经常带 BOM,如果不处理,Hadoop 读到的第一列第一个字段前面会多一个不可见字符,按病种统计时“高血压”和“\ufeff高血压”会被当成两个完全不同的 key,这是后续乱码问题的头号来源。
清洗前后行数差异要打印出来,这个数字会写进毕业设计文档的测试报告里,证明数据预处理环节是有效的。
3.3 上传到HDFS的两套入口:命令行PUT与Java API
清洗后的 CSV 上传到 HDFS,有两条路可以走。命令行最快,适合开发和演示时手动操作;Java API 适合集成到平台代码里,做成一个“数据上传”功能按钮。
# 在 HDFS 上创建目录结构 hdfs dfs -mkdir -p /disease/input # 上传清洗后的数据 hdfs dfs -put input/disease_records.csv /disease/input/ # 验证:-ls 看文件列表,-blocks 看数据块分布 hdfs dfs -ls /disease/input/ hdfs fsck /disease/input/disease_records.csv -files -blocks参数说明:-put 是上传,-mkdir -p 是递归建目录,fsck 命令会列出文件被切成几个 block、每个 block 存在哪个节点上。答辩时演示 fsck 的输出是展示你对 HDFS 底层机制有理解的最直观方式。
Java API 的写法对应的是平台内部的数据接入模块:
// UploadToHdfs.java // 把本地文件上传到 HDFS 指定路径,对应平台“数据接入”功能 import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; import org.apache.hadoop.fs.Path; String localPath = "/tmp/disease_records.csv"; String hdfsPath = "/disease/input/disease_records.csv"; Configuration conf = new Configuration(); conf.set("fs.defaultFS", "hdfs://localhost:9000"); FileSystem fs = FileSystem.get(conf); fs.copyFromLocalFile(new Path(localPath), new Path(hdfsPath)); fs.close(); System.out.println("上传完成:" + hdfsPath);逻辑说明:Configuration 里的 fs.defaultFS 要和 core-site.xml 保持一致,FileSystem.get 会读取 conf 并建立连接,copyFromLocalFile 完成本地到 HDFS 的拷贝。实际项目中这段代码会包在一个 Service 类里,由后端接口触发。
4. 统计逻辑落地:MapReduce多维计数作业的完整代码与调优参数
4.1 按病种统计发病数:一个可以直接编译的MapReduce作业
这是整个平台的核心作业,也是最简单的 MapReduce 入门形态。Mapper 读入 CSV 行,抽出病种字段作为 key,value 固定为 1;Reducer 把同一个病种的所有 1 相加,输出病种和总数。
// DiseaseCountJob.java // 按病种统计发病数:Mapper抽取病种,Reducer聚合计数 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.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.io.IOException; public class DiseaseCountJob { // Mapper:每行输入,输出 <病种, 1> public static class DiseaseMapper extends Mapper<Object, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text disease = new Text(); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); if (line.startsWith("case_id")) { return; // 跳过表头 } String[] fields = line.split(",", -1); if (fields.length < 6) { return; // 字段数不足,视为脏数据丢弃 } disease.set(fields[1].trim()); // 第二列是病种 context.write(disease, one); } } // Reducer:同一个病种的计数累加 public static class CountReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable result = new IntWritable(); @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } result.set(sum); context.write(key, result); } } public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "disease count by type"); job.setJarByClass(DiseaseCountJob.class); job.setMapperClass(DiseaseMapper.class); job.setCombinerClass(CountReducer.class); job.setReducerClass(CountReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); // 结果文件数直接决定输出 part 文件个数,这里设 1 便于后续导入 job.setNumReduceTasks(1); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }逻辑说明:split 用 ",", -1 而不是 split(","),区别在于 -1 会保留末尾的空字段,避免某行数据以逗号结尾时字段被截断。Combiner 这里直接复用 Reducer 的类,因为求和操作满足交换律和结合律;这个技巧要能讲明白原因,答辩常问。
参数说明:setNumReduceTasks(1) 会强制所有输出写进一个 part-r-00000 文件,方便后续和 MySQL 导入、图表展示对接。但要注意,如果数据量超过几百万行,Reduce 数量设 1 会成为性能瓶颈,那时应按数据量设置 3 到 5 个 Reduce,并且用下一节的方法处理多个输出文件。
4.2 组合Key与MultipleOutputs:一次扫描输出多张统计表
按病种统计只是最简单的一个维度。平台真正要展示的能力是“多维统计”:按病种、按地区、按年龄段、按月份。如果每个维度单独写一个作业跑一遍,逻辑没错,但文档里会显得设计得很笨。常见做法是用组合 Key 或者 MultipleOutputs 一次扫描同时输出多张结果表。
组合 Key 的思路是这样:Mapper 输出时,key 用 “病种#地区” 这样的组合字符串,Reducer 内部再拆开。代价是维度组合要提前写死。MultipleOutputs 更灵活,它允许在 Reducer 里根据不同条件把结果写入不同的输出文件。
// MultiDimStatJob.java 关键片段 // 在 Reducer 中同时输出按病种、按地区两张统计表 import org.apache.hadoop.mapreduce.lib.output.MultipleOutputs; import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat; public class MultiDimReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private MultipleOutputs<Text, IntWritable> mos; private IntWritable result = new IntWritable(); @Override protected void setup(Context context) { mos = new MultipleOutputs<>(context); } @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } result.set(sum); // key 形如 "高血压#杭州市",按分隔符拆开分别输出 String[] dims = key.toString().split("#"); if (dims.length == 2) { mos.write("by_type", new Text(dims[0]), result); mos.write("by_region", new Text(dims[1]), result); } } @Override protected void cleanup(Context context) throws IOException, InterruptedException { mos.close(); // 不关闭会丢失缓冲区内容 } }Driver 里通过 addNamedOutput 声明输出通道:
MultipleOutputs.addNamedOutput(job, "by_type", TextOutputFormat.class, Text.class, IntWritable.class); MultipleOutputs.addNamedOutput(job, "by_region", TextOutputFormat.class, Text.class, IntWritable.class);逻辑说明:map 阶段把地区和病种拼成组合 key,reduce 阶段拆开并分流到不同命名输出。输出目录里会出现 by_type-r-00000、by_region-r-00000 这样的文件。好处是只跑一次全量扫描,得到多张统计表;代价是组合 key 的分组方式要仔细设计,选错维度组合会导致 Reduce 阶段的数据倾斜。
年龄段分组在 map 阶段直接做离散化:读取 age 后判断落在哪个区间,输出 key 里拼上区间名,比如 “高血压#60-70岁”。这个逻辑放在 Mapper 里而不是 Reducer 里,能让分组计算提前完成,Reduce 端只做简单拼接。
4.3 提交到YARN的三个必调参数与调度流程复盘
作业代码写完后,打包成 jar,通过 YARN 提交。提交命令按下面这个模板来:
# 打包后用 yarn jar 提交 mvn clean package -DskipTests yarn jar target/disease-stat-1.0.jar \ com.disease.stat.mapreduce.DiseaseCountJob \ /disease/input \ /disease/output/count_type参数说明:第一个参数是输入目录,第二个是输出目录。输出目录在提交前必须不存在,否则 HDFS 会直接抛出 FileAlreadyExistsException,这个坑后面避坑章会专门展开。
作业提交到 YARN 的完整调度流程要能复盘:客户端把 jar 和输入分片信息发给 ResourceManager,RM 找一个 NodeManager 启动 ApplicationMaster,AM 向 RM 申请容器,容器里才真正跑 Map Task 和 Reduce Task。整个流程是答辩时的高频问题,建议把四步调度链路画进设计文档。
开发阶段还有三个必调参数要确认:
<!-- mapred-site.xml 中补充配置 --> <property> <name>mapreduce.map.memory.mb</name> <value>1024</value> </property> <property> <name>mapreduce.reduce.memory.mb</name> <value>1024</value> </property> <property> <name>mapreduce.job.reduces</name> <value>1</value> </property>mapreduce.map.memory.mb 和 reduce 的同名参数决定每个 Map/Reduce Task 能申请到的容器内存。伪分布式下物理内存有限,设太大容易溢出;设太小作业会在 shuffle 阶段频繁 GC。mapreduce.job.reduces 是在不写代码的情况下直接指定 Reduce 数量的开关,优先级低于 Job#setNumReduceTasks。
5. 避坑专场:伪分布式日常翻车的5条血泪经验
5.1 NodeManager总是自动退出:内存参数没配对
现象:start-yarn.sh 执行后,ResourceManager 在 jps 里能看见,NodeManager 进程启动几秒后自动消失,或者作业提交后一直卡在 ACCEPTED 状态不调度。
原因:NodeManager 启动时会检查 yarn.nodemanager.resource.memory-mb 的值,如果设得比物理机可用内存大,或与 ResourceManager 的调度内存不匹配,NM 会拒绝启动。更常见的场景是机器内存只有 4G,配置却按网上教程抄了 8G 的值。
解决:把 yarn-site.xml 里的内存参数调到物理机真实内存的一半以下,比如 4G 机器设 2048,同时把 vmem-check-enabled 设成 false。改完配置文件后,要执行 yarn stop 和 start-yarn.sh 完整重启,只重启 NodeManager 有时不会重新加载配置。
# 改完配置后的标准重启顺序 yarn stop start-yarn.sh jps5.2 NameNode启动失败:重复format的后果与重置方法
现象:hdfs namenode -format 执行完后,start-dfs.sh 启动,NameNode 日志里报错提示 clusterID 不一致,或直接打印 “NameNode format must be executed only once”。
原因:同一个 Hadoop 数据目录被 format 了多次,旧元数据里的 clusterID 和新写入的 clusterID 对不上。伪分布式下很多人习惯不好,遇到问题就 format,format 其实是初始化操作,不是维修操作。
解决:先把所有进程停掉,然后删除 hadoop.tmp.dir、dfs.namenode.name.dir、dfs.datanode.data.dir 三个配置指向的目录,全新 format 一次再启动。这个操作会清空 HDFS 上的所有数据,所以数据要提前备份到本地。
# 彻底重置 NameNode 数据目录 stop-dfs.sh rm -rf /opt/hadoop/data/tmp /opt/hadoop/data/nn /opt/hadoop/data/dn hdfs namenode -format start-dfs.sh5.3 统计结果中文乱码:BOM头与编码统一
现象:按病种统计的结果文件里,“高血压”显示正常,但同样的病种出现在两个不同的 key 下,其中一个 key 前面带着一个看不见的字符;或者导入 MySQL 后中文全部变成问号。
原因:CSV 文件是 UTF-8 with BOM 格式,清洗阶段没有用 utf-8-sig 读取,导致第一列第一个字段前面残留了 BOM 头。Hadoop 的 TextInputFormat 默认按 UTF-8 读取,BOM 会被当成普通字符拼进 key。MySQL 侧乱码则是表字符集和连接字符集不一致导致。
解决:清洗脚本里强制用 encoding="utf-8-sig" 读取,写出时用 encoding="utf-8";MySQL 建表时指定 utf8mb4 字符集,导入命令加 --default-character-set=utf8mb4。这两处统一后,乱码问题基本不会再出现。
5.4 Windows下Permission Denied:HADOOP_USER_NAME与winutils
现象:Windows 上用 IDEA 写代码,本地调试 MapReduce 作业时,报错信息里出现 Permission denied,或者 HDFS 写入时提示用户名不存在。
原因:HDFS 的权限模型认为当前用户是 Windows 系统用户名,比如 Administrator,而这个用户在 HDFS 上没有任何权限,访问 /disease 目录直接被拒。
解决:在 IDEA 的运行配置里加一个环境变量 HADOOP_USER_NAME=hdfs,让 HDFS 相认当前用户。如果还需要本地访问 HDFS 文件系统,winutils.exe 的版本要和 Hadoop 版本一致,这个工具负责在 Windows 下提供 HDFS 原生 POSIX 权限的实现,版本不匹配会出现各种莫名其妙的 native 方法报错。
IDEA Run Configuration 环境变量: HADOOP_USER_NAME=hdfs5.5 作业报FileAlreadyExistsException:输出目录清理的顺手指势
现象:同一个作业第二次跑,控制台报 org.apache.hadoop.mapred.FileAlreadyExistsException,输出目录已存在。
原因:MapReduce 的输出目录不允许预先存在,这是 HDFS 防止结果被覆盖的安全机制。很多人跑完一次不清理,直接重跑,就踩中这个错。
解决:每次跑批前删除输出目录。不要把这个操作写进 Java 代码里,在 bin 下的 run_jobs.sh 脚本里管理更清晰:
# run_jobs.sh:每次执行前先清理历史输出 hdfs dfs -rm -r /disease/output/count_type yarn jar target/disease-stat-1.0.jar \ com.disease.stat.mapreduce.DiseaseCountJob \ /disease/input /disease/output/count_type-rm -r会递归删除整个输出目录,配合前面第 4 章提到的 output 目录管理制度,整个跑批过程能做到无脑重复执行。
6. 结果入库与抽查验证:用MySQL和原始数据反推MR输出
6.1 用LOAD DATA把part文件导入MySQL:一条命令与三个参数
MapReduce 统计完成后,结果落在 HDFS 的 part 文件里。下一步是把这些结果导入 MySQL,供后端接口查询和图表展示。最直接的路径是先用 hdfs dfs -cat 拉到本地,再用 MySQL 的 LOAD DATA 命令导入,整条链路没有多余依赖。
# 1. 导出到本地 hdfs dfs -cat /disease/output/count_type/part-r-00000 > /tmp/count_type.tsv # 2. 导入 MySQL(表预先建好,字符集 utf8mb4) mysql -uroot -p --default-character-set=utf8mb4 disease_stat \ -e "LOAD DATA LOCAL INFILE '/tmp/count_type.tsv' \ INTO TABLE stat_result \ FIELDS TERMINATED BY '\\t' \ (disease_name, total_count);"这里三个参数最容易出错。FIELDS TERMINATED BY 必须写 '\t',因为 MapReduce 的 TextOutputFormat 默认用制表符分隔 key 和 value;LINES 默认按换行符识别,不用改;字符集必须用 utf8mb4,否则中文病种名入库后全部变成问号。导入后建议立刻执行 SELECT 抽查计数和 SUM 总计,核对是否和 part 文件一致。
6.2 抽查法验证:用原始数据反推MR输出是否可信
验证 MapReduce 结果对不对,最朴素的方法是从原始数据反推。选一个病种,直接在 HDFS 的原始 CSV 上用 grep 和 wc 统计,再和 MR 输出对比。
# 从原始数据统计高血压出现次数 hdfs dfs -cat /disease/input/disease_records.csv | grep "高血压" | wc -l # 查看 MR 统计输出 hdfs dfs -cat /disease/output/count_type/part-r-00000 | grep "高血压"两个数字一致,这个维度的统计就是可信的。如果数据量上了几十万行,grep 统计会很慢,但验证的目的本来就是抽查,不要求全量。
还有一个更高级的核验手段:计数恒等式。无论按哪个维度切分,病例总数应该相等。按病种统计的各病种数量之和、按地区统计的各地区数量之和,理论上都等于去重后的 total case_id 数。这个恒等式可以写成一个 SQL 校验脚本,放进测试报告里作为结论性证据。
-- 校验:各维度 SUM 是否一致 SELECT 'by_type' AS dim, SUM(total_count) AS total FROM stat_result_by_type UNION ALL SELECT 'by_region', SUM(total_count) FROM stat_result_by_region UNION ALL SELECT 'raw_data', COUNT(DISTINCT case_id) FROM disease_records;三个值相等,你的统计平台在逻辑上就是自洽的。我自己的习惯是每改一次 MapReduce 作业,先取 500 行小样本跑一遍,对照 part 文件行数和抽查数字,确认无误再跑全量。这套核验动作虽然简单,但能在答辩前拦下绝大多数低级错误,也让你面对“结果凭什么可信”这个问题时,能不心虚地把验证流程讲完整。希望帮到你。
本文还有配套的精品资源,点击获取