基于Hadoop与Spark的河南省空气质量预测系统

发布时间:2026/9/1 12:55:49
基于Hadoop与Spark的河南省空气质量预测系统 “河南省空气质量数据分析与预测系统”这个毕业设计选题核心是把 Hadoop 的分布式存储、Spark 的批量数据处理和 CatBoost 的回归预测串在同一条流水线里。它解决的不是“跑通某一个算法”而是“海量监测数据从哪里来、怎么存、怎么清洗、怎么分析、怎么用模型预测未来浓度”这一整条链路。适合两类人一类是计算机或大数据专业做毕设想找一个既有工程广度、又有算法深度的题目另一类是刚学完 Spark 基础想用真实业务把 Hadoop、Spark、机器学习串起来练手的人。最值得关注的不是组件多酷而是你能否在普通笔记本上把这条链路完整跑通并解释清楚每一步的输入、输出和参数。同时要注意如果只是简单把 CSV 读进来运行一个模型那不叫基于 Spark 的大数据分析系统。真正的重点是数据经过 HDFS、Spark 之后再进入 CatBoost 训练这中间每一步都有可能出现文件格式、编码、内存、分区不一致的问题。下面我把这套系统的落地过程按实际开发顺序拆开讲给你一条可以直接参考的路线。1. 这个选题到底在做什么先想清楚再动手很多同学看到“基于 Spark 的空气质量分析与预测”这个题目第一反应是“先装 Hadoop再装 Spark最后跑 CatBoost”。但这是典型的把顺序搞反了。真正应该先想清楚的是这套系统要回答什么问题。1.1 系统解决的实际问题河南省的空气监测站点会持续产生小时级甚至分钟级数据包括 PM2.5、PM10、SO2、NO2、O3、CO 等污染物浓度也会叠加温度、湿度、风速、风向等气象记录。数据量大、来源多、格式不一致单机 Excel 处理不了传统数据库也不适合做大规模清洗和特征计算。这时候需要 Hadoop 负责存储原始数据Spark 负责批量读取和处理再训练一个 CatBoost 模型去预测未来某个时刻的污染物浓度。从系统功能上看至少要覆盖四个模块数据采集与存储把历史数据上传到 HDFS。数据清洗处理缺失值、重复值、异常值和格式问题。数据分析和特征工程统计城市、月份、小时维度的污染规律并为模型构造特征。模型训练与预测用 CatBoost 训练回归模型输出未来 PM2.5 或 AQI 的预测值。这个系统最终要能回答几类问题河南哪个城市污染最严重冬天和夏天的浓度差距有多大一天当中哪个时段浓度最高明天某一站点的 PM2.5 大概是多少第一个和第二个问题靠 Spark 统计解决第三个问题需要把滞后特征、时间特征和气象特征组合起来交给 CatBoost。1.2 技术栈中每个组件的边界Hadoop、Spark、CatBoost 并不是越复杂越好每个组件都有自己的职责边界。Hadoop 在这个项目里的核心是 HDFS而不是 MapReduce。因为数据存在 HDFS 上后续 Spark 可以直接从 HDFS 读取不需要把 CSV 复制到每台机器上。如果你只是本地运行也可以在 Linux 上使用 Hadoop 伪分布式模式先跑通流程再考虑多台机器组成集群。相关热搜里经常提到“hadoop 集群搭建”“hadoop 安装与配置”“hadoop 和 zookeeper 整合实战”这些都是环境层面的投入。但从毕设角度我建议先只在单机伪分布式中验证完整数据流不要一上来就搭三台服务器否则后期调试会非常痛苦。Spark 负责的是分布式数据清洗、统计和特征工程。常见做法是利用 PySpark 的 DataFrame 和 Spark SQL把原始数据从 HDFS 读出经过过滤、去重、补全、聚合、窗口计算后再输出成特征文件。这里要注意Spark 不是用来训练 CatBoost 的。CatBoost 是单机内存中的梯度提升树框架通常的做法是让 Spark 生成特征数据集导出成 CSV 或 Parquet再交给 CatBoost 训练。这里有一个容易出现的误解为了体现“大数据”非要把 CatBoost 也放到 Spark 里去分布式训练。对于毕设项目来说完全没必要反而增加了部署难度。你只要在文献综述和答辩说明里讲清楚Spark 负责海量数据的预处理和特征工程CatBoost 负责高精度回归预测两者通过文件接口衔接这就够了。1.3 完整数据流我画过很多次这套系统的架构图最容易讲清楚的是一条顺序流水线监测数据 / 公开数据集 ↓ HDFS 原始数据存储 ↓ Spark 数据清洗去重、缺失值、异常值、格式统一 ↓ Spark 统计分析城市排名、时间趋势、相关系数 ↓ Spark 特征工程时间特征、滞后特征、滚动统计 ↓ 导出特征 CSV / Parquet ↓ CatBoost 训练与预测 ↓ 结果写回 HDFS / MySQL / 可视化页面建议先按这条线把每一段的输入输出写清楚再去写代码。如果你在项目启动第一天就把环境装好然后直接训练模型很可能出现“数据文件有乱码”“时间列没有解析”“滞后特征全是空值”这些问题。更合理的顺序是先准备一小份样例数据完整走一遍流水线确认每一步输出正确再扩大数据量。2. 环境准备与工程目录别急着写代码环境准备是最容易让人心态崩溃的部分。尤其是第一次接触 Hadoop 和 Spark 的人经常会卡在“jar does not exist or is not a normal file: /usr/local/hadoop/share/hadoop/m...”这类路径错误上。这类问题通常不是功能不支持而是环境变量、目录权限或文件路径不对。2.1 本地开发还是集群环境先说结论如果你的目标是先跑通功能优先选择本地伪分布式或单机 Spark。如果项目要求必须做集群部署再准备三台虚拟机或云主机。伪分布式的意思是 Hadoop 的 NameNode、DataNode、ResourceManager、NodeManager 都跑在同一台机器上虽然不体现多机分布式能力但 HDFS 的操作方式、Spark 读取 HDFS 的代码路径、权限配置和真实集群基本一致。对毕设来说这样的环境足够完成数据上传、清洗、统计和结果回写。硬件方面8G 内存起步16G 更稳。磁盘至少留 100G 空间因为 Hadoop 默认会保存副本再加上中间结果和模型输出占空间比你想象得快。安装顺序一般是JDK、Hadoop、Spark、Python然后再安装 PySpark 和 CatBoost 的 Python 包。不要跳过 JDK因为没有 JDK 的话 Hadoop 和 Spark 都起不来。如果你还需要在 Windows 上调试建议把 Spark 和 Hadoop 安装到虚拟机或 WSL 中Windows 下直接安装 Spark 虽然也能跑但文件路径和权限问题会多不少。常见环境配置项包括JAVA_HOME指向 JDK 安装目录。HADOOP_HOME指向 Hadoop 安装目录。SPARK_HOME指向 Spark 安装目录。Python 中安装pyspark、pandas、catboost、scikit-learn。装好之后先用一个最简单的 Spark 程序验证环境from pyspark.sql import SparkSession spark SparkSession.builder.appName(env_check).master(local[*]).getOrCreate() df spark.createDataFrame([(1, a), (2, b)], [id, value]) df.show()能跑出来就说明基础环境没问题。很多同学不先做这一小步直接跑几十万条数据一旦报错也分不清是环境问题还是代码问题。2.2 数据集怎么准备河南省空气质量数据可以从公开数据源获取比如中国环境监测总站的公开历史数据、Kaggle 上的城市空气质量数据集或者学校提供的实验数据。关键是数据里要有时间、站点、污染物浓度最好还有气象字段。没有气象字段模型预测效果一般不会太好。我建议先把数据整理成统一 CSV 格式。字段可以参考下面这张表字段名类型示例说明timestampdatetime2024-01-01 08:00:00监测时间citystring郑州城市名称station_codestringzz_001站点编号pm25float75.5PM2.5 浓度pm10float120.3PM10 浓度so2float10.2二氧化硫no2float45.8二氧化氮o3float80.6臭氧cofloat0.9一氧化碳temperaturefloat12.5温度humidityfloat58.0湿度wind_speedfloat3.2风速wind_directionstring东北风风向如果你从爬虫或接口采集数据原始字段可能不是这个顺序甚至站点名称和城市名称混在一起。不要在原文件里手工改最好用脚本先做一次格式标准化再上传 HDFS。因为后续每一步都依赖字段名和时间格式统一得越早后面越省事。2.3 工程目录与模块划分很多毕设项目最后代码全堆在一个 Notebook 或一个 Python 文件里答辩时很难讲清楚。我建议按模块拆分air_quality_system/ data/ raw/ # 原始 CSV cleaned/ # 清洗后的数据 features/ # 特征数据 etl/ spark_clean.py # 数据清洗 analysis/ spark_analysis.py # 统计分析 model/ train_catboost.py # 模型训练 predict.py # 预测脚本 web/ app.py # 可视化后端 templates/ # 页面模板 scripts/ upload_hdfs.sh # 上传脚本这样划分的最大好处是每层都有清晰输入输出报错时能快速定位。比如模型效果差可能是特征文件的问题先把特征文件单独检查一遍比如 CSV 读进来中文乱码就直接去 etl 层看编码参数。答辩时你也能按目录顺序讲述整个系统而不是“这是我的代码文件夹”。3. 数据采集与清洗先让数据可信数据清洗往往占整个项目 60% 以上工作量。不要小看这一步模型后面效果不好十有八九是数据没洗干净。3.1 数据怎么进入 HDFS先把本地 CSV 上传到 HDFS。在伪分布式环境下HDFS 命令和在真实集群上基本一致hdfs dfs -mkdir -p /air/raw hdfs dfs -put ./data/raw/henan_air_2020_2024.csv /air/raw/上传完成后可以检查一下hdfs dfs -ls /air/raw这里有一个关键点最好先把 CSV 转换成无 BOM 的 UTF-8 格式再上传。否则 Spark 读取时第一列列名可能带\ufeff导致后续select(timestamp)报找不到列。如果你想性能更好也可以先用 pandas 把 CSV 转成 Parquet 格式再上传Spark 读取 Parquet 会快很多而且类型自动保留。但这会增加一步转换逻辑新手阶段可以先直接用 CSV。3.2 清洗规则与代码实现Spark 读取 HDFS 上的 CSV一般用option(header, True)和option(inferSchema, True)。注意inferSchema挺方便但大数据量时推断代价高建议先用小数据确定字段类型再显式schema指定。示例代码如下from pyspark.sql import SparkSession spark SparkSession.builder.appName(air_quality_clean).getOrCreate() df spark.read \ .option(header, True) \ .option(encoding, UTF-8) \ .csv(/air/raw/henan_air_2020_2024.csv) # 去重同一站点、同一时刻只能有一条记录 df df.dropDuplicates([city, station_code, timestamp]) # 过滤明显异常值浓度不能为负数CO 单位要一致 df df.filter( (df[pm25] 0) (df[pm10] 0) (df[so2] 0) (df[no2] 0) ) df.show(10)清洗时还要处理空值。时间序列数据里某一天某个站点可能因为设备维护缺了一条记录。常见做法有两种如果只是孤立缺失可以用前后时刻均值填充。如果缺失比例很高比如某站点一个月都没有数据建议直接把该站点整段剔除。不能用整列平均值去填充所有缺失值因为空气污染有明显的季节和日夜变化全局均值会让特征带有误导性。处理缺失值时我会按城市和站点分组再用窗口函数做前后值平均这样更符合实际情况。如果数据量比较大建议把清洗结果写回 HDFS方便下一步使用df.write.mode(overwrite).parquet(/air/cleaned)3.3 常见清洗坑我把自己踩过的一些坑列出来你遇到时先按这个顺序排查先看文件编码再看时间格式再看数据粒度。中文乱码CSV 文件编码不是 UTF-8会导致城市名称乱码。时间解析失败数据里既有2024/1/1 8:00又有2024-01-01 08:00:00需要统一成 Spark 的timestamp类型。重复数据一台站点一天可能上报多条相同记录必须先按唯一键去重。数值单位不一致比如 CO 有的数据是mg/m3有的是ug/m3不统一会导致特征失真。小文件问题清洗结果如果每个分区都写大量小 CSVHDFS 上会出现成千上万个小文件后续 Spark 读取会变慢。建议统一写成 Parquet或者coalesce(1)导出。注意不要一上来直接用coalesce(1)把所有数据缩成一个文件。如果数据量很大单个文件反而会让下游处理变慢。先看数据量级几百 MB 以内可以几 GB 以上就保持适度分区。4. Spark 分析从统计口径到特征工程数据清洗完之后进入分析阶段。这个阶段既要做统计报表又要为模型准备特征。很多人把这两个任务混在一起写最后代码很难维护。我建议先做统计再做特征。4.1 核心分析任务怎么拆你可以用 Spark SQL 完成大部分统计。先注册成临时表df.createOrReplaceTempView(air)然后按城市和月份统计平均浓度SELECT city, month(timestamp) AS month, round(avg(pm25), 2) AS avg_pm25, round(avg(pm10), 2) AS avg_pm10, count(*) AS record_count FROM air GROUP BY city, month(timestamp) ORDER BY city, month看到record_count就能快速判断哪个时段哪个站点数据缺失严重。如果某个城市某个月只有几十条记录那统计结果不具代表性后面做特征时也要留意。还可以计算污染物之间的相关系数。Spark DataFrame 自带stat.corr方法corr_pm25_pm10 df.stat.corr(pm25, pm10) print(corr_pm25_pm10)这些统计指标一方面用于可视化展示另一方面也能帮你理解数据。比如 PM2.5 和 PM10 高度相关那么建模时特征里可以都保留但不代表两个特征没有冗余。4.2 性能优化分区、缓存、广播变量如果你的数据集只有几十万条在本地 Spark 上不加任何优化也能跑。但如果你上了几千万条需要提前注意资源占用。常见优化顺序先把用不到的列删掉减少 shuffle 数据量。多次复用的 DataFrame 调用cache()。小表关联大表时用小表广播避免 shuffle join。窗口函数分区不能过大否则单个 executor 的内存压力很大。比如做特征时要对每个城市分别排序计算滞后值如果PARTITION BY city只有 18 个城市分区很小看起来很快。但如果改成按站点分区站点数量可能几百个每个分区的数据量会不同部分站点数据少、部分站点数据多容易出现数据倾斜。再比如使用cache时要注意如果一个 DataFrame 只被使用一次缓存没有意义反而浪费内存。你可以这样判断后续代码中同一个df会被groupBy、filter、join多次就cache()只用一次就不缓存。4.3 特征工程滞后特征和滚动统计是重点CatBoost 虽然是强模型但它不会自动理解“上一时刻的 PM2.5 和当前 PM2.5 的时序关系”。因此必须手动构造特征。常用特征包括时间特征小时、星期、月份、季节、是否节假日。滞后特征目标城市过去 1 小时、6 小时、24 小时的 PM2.5 浓度。滚动统计过去 24 小时 PM2.5 均值、最大值、标准差。气象特征温度、湿度、风速、风向。Spark SQL 的窗口函数非常适合做这件事。下面是一个示例按城市和时间排序计算当前时刻过去 24 小时的均值以及滞后 24 小时的历史值。SELECT city, station_code, timestamp, pm25, hour(timestamp) AS hour, dayofweek(timestamp) AS dow, month(timestamp) AS month, lag(pm25, 24) OVER w AS pm25_lag24, avg(pm25) OVER w AS pm25_avg24 FROM air WINDOW w AS (PARTITION BY city ORDER BY timestamp ROWS BETWEEN 24 PRECEDING AND 1 PRECEDING)注意这里ROWS BETWEEN 24 PRECEDING AND 1 PRECEDING表示从过去第 24 行到前一行不包含当前行可以避免使用未来数据和当前数据导致的数据泄漏。如果目标是预测未来 24 小时你需要把目标列pm25_t24也通过lead(pm25, 24) OVER w构造出来然后用滞后特征去预测它。这样的目标构造方法才是合理的时间序列预测。5. CatBoost 预测模型训练与结果评估特征表生成之后系统开始进入模型阶段。这里有一个衔接动作把 Spark 输出的特征 CSV 或 Parquet 下载到本地然后用 Pandas 读取交给 CatBoost。如果你的数据总量不大这一步完全可以本地完成。5.1 为什么选 CatBoostCatBoost 是梯度提升树框架和 XGBoost、LightGBM 属于同一类模型。相比其他两个它最突出的特点是原生支持类别特征不需要手动 One-Hot 编码。比如城市名、季节、风向这些字段直接传给 CatBoost并指定cat_features即可。对毕设项目来说这能省掉不少特征编码的代码。另外 CatBoost 的默认参数通常能给出不错的结果很适合作为基础模型。如果你想证明自己在算法上做了比较可以再加一两个对比模型比如线性回归或随机森林然后用 RMSE、MAE、R2 做对比。注意不要为了对比而对比重点是解释为什么 CatBoost 更适合这个数据。5.2 训练集与测试集划分时间序列预测不能像普通分类问题那样随机切分训练集和测试集否则会把未来的信息泄漏到训练集里测试指标虚高。我建议按时间排序后用前 80% 的数据训练后 20% 的数据测试。import pandas as pd from catboost import CatBoostRegressor, Pool features pd.read_csv(data/features/air_features.csv, parse_dates[timestamp]) features features.sort_values([city, timestamp]) # 用全局时间切分而不是每个城市随机切分 cutoff features[timestamp].quantile(0.8) train features[features[timestamp] cutoff] test features[features[timestamp] cutoff] feature_cols [hour, dow, month, is_holiday, temperature, humidity, wind_speed, pm25_lag1, pm25_lag6, pm25_lag24, pm25_avg24, pm25_max24] cat_cols [city, season, wind_direction] train_pool Pool(train[feature_cols], train[target], cat_featurescat_cols) test_pool Pool(test[feature_cols], test[target], cat_featurescat_cols) model CatBoostRegressor( iterations1000, learning_rate0.05, depth6, loss_functionRMSE, verbose200 ) model.fit(train_pool, eval_settest_pool, early_stopping_rounds100)这里target不是当前时刻的 pm25而是未来 24 小时后的 pm25。如果你用当前 pm25 既当特征又当目标模型训练时看似效果很好实际预测时却无法复现。5.3 参数调整和评估指标CatBoost 常用参数范围可以参考下面这张表但实际参数要以你自己的数据为准不要照搬参数常见范围说明iterations500 - 2000树的数量配合早停learning_rate0.01 - 0.1学习率越小越慢但更稳depth4 - 8树深度过深容易过拟合l2_leaf_reg3 - 10L2 正则降低过拟合early_stopping_rounds50 - 100验证集连续无提升时停止评估指标建议同时看 RMSE、MAE 和 R2。RMSE 对大误差更敏感MAE 更直观R2 能反映模型相对均值预测的改善程度。代码里可以在model.fit后使用model.best_score_查看验证集结果。预测完还要画出真实值和预测值的散点图或时间序列图这是答辩时最直观的展示材料。5.4 预测结果落库与可视化模型预测结果不能只停留在 notebook 里。你需要把预测结果写回 HDFS 或 MySQL再交给可视化后端。我比较推荐的方案是用 PySpark 读取预测结果文件写入 HDFS 的/air/predictions目录同时导出一份 CSV 给 Flask 后端用 ECharts 画图。这样既体现了大数据系统的闭环又满足了可视化展示需求。可视化页面重点展示四块内容河南省各地市近期空气质量排名。某站点过去一年 PM2.5 时间变化曲线。污染物相关性和特征重要性柱状图。未来 24 小时预测曲线和真实值对比。可视化不要做得太复杂关键是能把“清洗之后的数据规律”和“模型预测效果”讲清楚。6. 从单机到集群部署和调优的关键点整体跑通之后如果项目要求部署到集群或者你需要把 Spark 作业提交到 YARN 上就要关注资源参数和排查链路。很多 Spark 作业不是代码写错了而是提交参数不合理。6.1 Spark 提交任务与资源参数在本地代码中SparkSession 常使用master(local[*])。到了集群环境要改成从提交命令中传入spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 3 \ etl/spark_clean.py参数含义如下--driver-memoryDriver 内存负责解析任务和调度。--executor-memory每个 Executor 内存真正执行计算。--executor-cores每个 Executor 的 CPU 核数。--num-executorsExecutor 数量。不要盲目把内存和核数拉到最大。第一次调试时我建议按集群总资源的 30% 到 50% 配置先跑通再逐步调大。比如三台节点每台 8G 内存先把 executor 内存设为 2g核数设 1跑一个小文件验证流程。如果一上来就申请超大资源很可能因为其他任务占用导致等待甚至直接 OOM。6.2 常见报错和排查链路基于这套系统我整理了几个高频问题和优先排查顺序现象优先排查Spark 作业一直卡在某个 stage先看 YARN 资源是否充足再看 stage 内部是否发生数据倾斜Executor OOM缩小处理分区数或减少窗口函数数据量不能只加内存读取 HDFS 文件报错 file not found检查路径、HDFS 权限、文件名是否带后缀中文乱码检查 CSV 编码和option(encoding, UTF-8)CatBoost 报特征类型错误检查特征列是否存在 NaN、字符串列是否被误识别为数值列预测结果基本是常数优先排查目标列构造是否使用了未来数据导致时间泄漏实际排查时建议按“现象 - 输入 - 环境 - 参数 - 模型”的顺序来。比如 Spark 作业卡住第一步不是改代码而是先看 YARN 页面有没有资源、Executor 日志有没有 GC再决定是不是需要调spark.sql.shuffle.partitions。如果是读取 HDFS 的文件错误先检查 HDFS 路径和文件名再看本地文件是否已经上传成功。很多人会把这类问题当成 Spark 代码问题改了半天下载才知道是 HDFS 路径写错。6.3 答辩展示和验收建议如果是毕设答辩建议不要只展示代码运行成功的截图。你应该准备一条完整可讲述的证据链一张系统架构图说明数据如何从采集到预测。一份清洗前后的数据对比说明你处理了哪些脏数据。一组统计结果比如各城市月平均 PM2.5 变化。一张特征重要性图说明模型主要依赖哪些特征。一条真实站点的预测曲线同一张图中包含真实值和预测值。还要能诚实回答系统局限。比如“只用了历史监测数据和部分气象数据没有融合气象预报数据”“某些站点缺失较多预测结果可能偏差较大”。这些不是减分项反而能证明你对系统边界有清晰认知。最后这套系统真正落地时最该盯住的不是功能列表而是输入格式、资源占用和失败重试。如果你也是拿它做毕设我建议先把单条数据流跑通再考虑集群和调优否则很容易一个环节卡住整体停滞。数据清洗和特征构造做得越扎实后面的 CatBoost 预测才有说服力。