
空气质量类的毕业设计在往年并不是一个热门方向。但近几年情况完全不同很多学校的大数据课程设计和毕业设计选题里“空气质量数据分析与预测”开始高频率出现。原因也很直白数据公开、指标标准、业务场景清晰、可视化效果好。更重要的是它能把 Hadoop、Spark、机器学习算法这些硬核技术串成一条完整的链路从数据采集到存储从清洗分析到模型预测每一层都有可展示的成果。这个题目真正卡住人的地方不是技术本身而是“不知道从哪里开始”。很多同学装了 Hadoop 又装 Spark最后发现只会在命令行敲几个启动脚本数据还在本地文件里躺着分析结果也只是一张静态图表。这个项目看起来是一个“基于 Spark CatBoost 的预测系统”但真正做到位的其实是数据工程和机器学习之间那条“接口”。本文就按照一个可以落地的毕业设计标准把整体架构、环境搭建、数据梳理、Spark 分析、CatBoost 预测和工程化建议完整拆开讲一遍。1. 这个选题能解决什么问题先下一个明确判断这个选题不是“套壳大数据”而是一套能真正跑起来的技术闭环。很多毕业设计项目存在三个问题。第一题目太空比如“基于大数据的某系统”最后只是用 Python 处理了一个 Excel 文件完全没有大数据处理的味道。第二技术栈太旧还在用最原始的方式做数据统计体现不出 Spark 分布式计算的价值。第三没有模型产出只有“分析报表”没有“预测能力”答辩时容易被追问到底。河南省空气质量分析与预测系统恰好可以覆盖这些盲点。数据维度河南 18 个地级市郑州、洛阳、南阳、周口、商丘等的空气监测站点数据字段包含城市、日期、AQI、PM2.5、PM10、SO2、NO2、CO、O3数据量级可以达到几十万条以上适合用 HDFS 存储和 Spark 分布式处理。分析维度年度/季度/月度空气质量变化趋势、各地市 PM2.5 排名、首要污染物分布、重污染天气频率统计、冬季与夏季污染差异。预测维度基于历史监测数据用 CatBoost 模型预测未来某一天某个城市的 AQI 或 PM2.5 浓度既可以是回归任务也可以转成分类任务优、良、轻度污染、中度污染、重度污染、严重污染。展示维度最终可以对接 ECharts 或 Superset输出趋势曲线、热力图、地图下钻等可视化页面。这个项目最值得做的不是“算法调参”而是数据全链路的稳定性。换句话说如果你的数据清洗没做好、特征构造不合理再好的模型也白搭。答辩时真正能拿高分的是你说清楚每一层在做什么、为什么这么做、遇到什么问题、怎么解决的。从适用的读者来看这类文章最适合三类人计算机/软件工程/大数据相关专业的学生正在选毕业设计题目需要完成 Hadoop、Spark 课程设计或实训项目的在校生想从“会调 API”走向“会做数据项目”的初学者。2. 系统总体架构与核心概念2.1 整体架构设计这个项目的技术架构可以分为五层层次组件作用数据源层公开空气质量监测数据 / 模拟数据原始 CSV、JSON 数据存储层Hadoop HDFS分布式存储原始数据和分析结果计算层Spark SQL / Spark DataFrame数据清洗、统计分析、特征工程模型层CatBoost、scikit-learnAQI/PM2.5 回归预测应用层Flask/FastAPI ECharts/Superset结果展示、可视化大屏整个项目不需要一开始就搭企业级集群。开发阶段用 Hadoop 伪分布式 Spark Local 模式就能完成全部流程有条件再用三台虚拟机或云主机搭集群。关键在于代码是跟着分布式思路走的而不是临时写一套单机脚本。2.2 Hadoop 与 HDFS 在项目中的角色Hadoop 在项目里主要提供 HDFS 存储。HDFS 的特点是“一次写入、多次读取”非常适合存储大规模的历史监测数据。它会自动把数据切成块Block并做多副本冗余默认副本数通常为 3在伪分布式模式下可以调成 1避免磁盘浪费。在毕业设计场景中HDFS 的作用不是“显得高级”而是让数据存储和计算分离。原始监测数据上传到 HDFSSpark 程序从 HDFS 读取数据分析结果再写回 HDFS 或数据库。这个过程体现的是大数据处理的标准范式。2.3 Spark 为什么比 Pandas 更合适Pandas 是单机内存计算数据量大一点就容易内存溢出Spark 是分布式内存计算能同时处理多台机器上的数据。本项目中如果只是几万条数据Pandas 确实够用但毕业设计要体现“大数据分析”的能力就必须有一个理由来使用 Spark。更实际的原因是Spark SQL 可以直接用 SQL 语法对 DataFrame 做复杂查询支持分区、广播变量、缓存等优化手段而且 Python 的 Pandas 在 Spark 3.x 中可以通过 Pandas API on Spark 或 spark.sql 紧密配合。这样既降低了编程难度又保留了分布式处理能力。2.4 CatBoost 算法为什么适合做预测CatBoost 是 Yandex 开源的一种梯度提升树GBDT算法。相比 XGBoost 和 LightGBM它有四个特点非常适合这个项目。第一它原生支持类别特征。空气质量数据里有“城市”这种分类字段在 CatBoost 中不需要手动做 One-Hot 编码声明成 cat_features 即可。第二它对缺失值有内置处理策略不需要为缺失的监测数值专门做复杂插值。第三它使用对称树Oblivious Tree结构训练速度快不容易过拟合而且默认参数就能取得不错的效果。第四它支持特征重要性分析可以输出哪些特征对 AQI 预测影响最大这也适合在论文里做解释。当然并不等于说 CatBoost 一定比随机森林、XGBoost 好而是它在“结构化表格数据 少量类别特征 有缺失值”的场景下是性价比最高的选择非常适合本科阶段的毕业设计实际落地。3. 环境准备与前置条件3.1 环境版本建议由于不同机器的环境差异较大这里不写死具体版本号只给一个各组件兼容的建议范围实际安装时以官方文档为准。组件建议版本说明JDK1.8 或 11Hadoop 3.x 必须使用 JDK 8 或 11Hadoop3.x本文基于 Hadoop 3.x 的伪分布式部署Spark3.x需要与 Hadoop 版本兼容Python3.8推荐使用 Conda 管理环境CatBoost1.xpip install catboost操作系统Ubuntu 20.04 / CentOS 7 / Windows WSL2建议 Linux 环境Windows 需要额外处理 Hadoop 本地库3.2 Hadoop 伪分布式安装要点安装 Hadoop 的过程网上资料很多但最容易踩坑的是三处SSH 免密登录、Namenode 格式化、环境变量。# 1. 配置 SSH 免密登录 ssh-keygen -t rsa -P -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys ssh localhost # 2. 配置 Hadoop 环境变量写入 ~/.bashrc export HADOOP_HOME/usr/local/hadoop export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin export HADOOP_CONF_DIR$HADOOP_HOME/etc/hadoop source ~/.bashrc修改$HADOOP_HOME/etc/hadoop/core-site.xml和hdfs-site.xml!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name valuefile:///usr/local/hadoop/tmp/name/value /property property namedfs.datanode.data.dir/name valuefile:///usr/local/hadoop/tmp/data/value /property /configuration注意修改配置后第一次启动前必须格式化 Namenode否则会报java.io.IOException: NameNode is not formatted。hdfs namenode -format start-dfs.sh jpsjps命令能看到 NameNode、DataNode、SecondaryNameNode 三个进程说明 HDFS 启动成功。3.3 Spark 安装与 Python 环境配置Spark 安装更简单下载二进制包解压后配置环境变量即可。export SPARK_HOME/usr/local/spark export PATH$PATH:$SPARK_HOME/bin:$SPARK_HOME/sbin开发阶段用 PySpark 运行建议创建一个独立 Conda 环境conda create -n air_quality python3.8 -y conda activate air_quality pip install pyspark catboost pandas scikit-learn matplotlib有一点要提前说明PySpark 在 Python 3.8 以上需要注意 Java 进程和 Python 进程之间的通信限制数据结构如果太大直接用collect()容易把 Driver 内存打满。本文后面的代码会刻意避免把大量数据收集到本地而是用 Spark 自带的聚合和输出功能。4. 数据说明与预处理流程4.1 数据源与字段设计本项目可以使用两类数据。一类是公开的空气质量监测历史数据通常可以从中国环境监测总站等公开渠道获取字段一般包括监测站点、城市、日期、时点、AQI、PM2.5、PM10、SO2、NO2、CO、O3、首要污染物、空气质量等级。另一类是真找不到数据时用脚本生成的模拟数据。但生成模拟数据有一个前提必须能模拟出时间趋势和城市差异不能是纯随机数。建议做法是以郑州为基准冬季 PM2.5 均值设定为 90 左右夏季均值设定为 45 左右再叠加高斯噪声和不同城市的偏移量。无论使用哪种数据统一整理成以下结构字段名类型说明citystring城市名称如郑州datestring日期如 2024-01-05aqiint空气质量指数pm25doublePM2.5 浓度pm10doublePM10 浓度so2double二氧化硫浓度no2double二氧化氮浓度codouble一氧化碳浓度o3double臭氧浓度levelstring空气质量等级????????????这里需要提醒如果在模型中直接使用“level”作为预测目标它就是一个多分类问题如果用“aqi”作为预测目标就是回归问题。从毕业设计的角度来讲两者都可以做但建议先做回归再做分类因为 AQI 的数值预测可以继续计算 MAE、RMSE 这类指标写论文时更有说服力。4.2 数据清洗的关键动作数据清洗是整个项目里最容易被低估、也最影响后续效果的一步。在 Spark 中建议按照以下顺序处理去重同一城市同一天的记录只能有一条。时间解析把 date 字段转成标准日期格式抽取 year、month、day、weekday 等特征。缺失值处理先统计每列的缺失率。如果缺失比例低于 5%可以直接删除或填充均值如果某一列缺失比例过高要考虑是否从特征中剔除。异常值处理AQI 不可能为负值PM2.5 也不可能超过 1000。使用filter做范围过滤避免异常值污染模型。类别特征统一城市名称的大小写、空格优先转成统一格式。以下是一个用 PySpark 做基础清洗的示例# 文件路径scripts/clean_data.py from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_date, year, month, dayofmonth, when spark SparkSession.builder.appName(AirQualityClean).getOrCreate() # 读取 HDFS 上的原始数据 df spark.read.option(header, True).csv(hdfs://localhost:9000/data/air_quality_raw.csv) # 去重 df df.dropDuplicates([city, date]) # 日期解析与特征抽取 df df.withColumn(date, to_date(col(date), yyyy-MM-dd)) df df.withColumn(year, year(col(date))) df df.withColumn(month, month(col(date))) df df.withColumn(day, dayofmonth(col(date))) # 过滤不可能出现的负值 df df.filter(col(aqi) 0) df df.filter(col(pm25) 0) df df.filter(col(pm10) 0) # 缺失值填充这里只填充数值列按城市分组取均值 numeric_cols [aqi, pm25, pm10, so2, no2, co, o3] for c in numeric_cols: df df.withColumn( c, when(col(c).isNull(), 0).otherwise(col(c)) ) # 输出清洗后结果 df.write.mode(overwrite).parquet(hdfs://localhost:9000/data/air_quality_clean.parquet) spark.stop()这段代码需要注意的是缺失值直接填充 0 对预测模型有很大负面影响。实际项目中更推荐做法是用Spark的groupBy(city).avg(c)计算各城市均值再用这个均值做填充。上面代码只演示流程不要直接照抄到正式项目里。5. 基于 Spark SQL 的空气质量分析完成了清洗下面进入统计分析阶段。这一步不仅是“做题”更是为后续建模准备特征。5.1 上传数据与创建临时表先把原始数据放到 HDFS 中。hdfs dfs -mkdir -p /data hdfs dfs -put air_quality_raw.csv /data/然后启动 PySpark读取清洗后的 Parquet 数据并注册成临时表# 文件路径scripts/analysis.py from pyspark.sql import SparkSession spark SparkSession.builder.appName(AirQualityAnalysis).getOrCreate() df spark.read.parquet(hdfs://localhost:9000/data/air_quality_clean.parquet) df.createOrReplaceTempView(air_quality)5.2 统计每年河南省 PM2.5 平均浓度SELECT year, ROUND(AVG(pm25), 2) AS avg_pm25, COUNT(*) AS record_count FROM air_quality GROUP BY year ORDER BY year;在 PySpark 中执行yearly_pm25 spark.sql( SELECT year, ROUND(AVG(pm25), 2) AS avg_pm25, COUNT(*) AS record_count FROM air_quality GROUP BY year ORDER BY year ) yearly_pm25.show()预期输出是一张按年份排列的统计表能明显看出冬季峰值。这里record_count可以用来检查每年数据量是否均衡防止某一年数据缺失导致结论偏差。5.3 各地市污染天数排名污染天数定义为 AQI 100即轻度污染及以上的天数。SELECT city, SUM(CASE WHEN aqi 100 THEN 1 ELSE 0 END) AS polluted_days, ROUND(AVG(aqi), 2) AS avg_aqi FROM air_quality GROUP BY city ORDER BY polluted_days DESC;通过这个查询可以快速看出河南哪些城市重污染天气更频繁。答辩时可以结合地图可视化解释地理和工业布局因素但不要做过度推测。5.4 特征工程构造滞后特征与滑动平均在建模之前只拿pm25本身预测未来的pm25是不够的。还要加入“时间滞后”和“统计聚合”特征让模型能够感知历史趋势。from pyspark.sql.window import Window from pyspark.sql.functions import col, lag, avg # 按城市分组按日期排序 window_spec Window.partitionBy(city).orderBy(date) # 滞后1天、7天的 PM2.5 df df.withColumn(pm25_lag1, lag(pm25, 1).over(window_spec)) df df.withColumn(pm25_lag7, lag(pm25, 7).over(window_spec)) # 过去7天 PM2.5 滑动平均 df df.withColumn(pm25_7d_avg, avg(pm25).over( window_spec.rowsBetween(-7, -1) )) # 滞后特征有空值先过滤 df_model df.dropna(subset[pm25_lag1, pm25_lag7, pm25_7d_avg])这一步是建模前最重要的一步。如果滞后特征为空一定要dropna否则后面训练模型会报错或产生严重偏差。6. CatBoost 预测模型实现6.1 确定预测目标与评估指标预测目标设为“明天的 PM2.5 浓度”。这是一个回归任务评价指标采用 RMSE均方根误差和 MAE平均绝对误差同时输出 R² 判断模型解释力。评估指标的选择要看具体场景RMSE对较大误差更敏感适合衡量极端污染事件的预测能力MAE更直观代表平均误差多少μg/m³R²越接近 1 越好说明模型能解释数据中大部分方差。6.2 将 Spark 处理好的数据转成训练集数据量少的时候可以先把结果collect()到本地再交给 CatBoost 训练。数据量大时需要走 Spark 分布式训练或把结果写入特征表再用 Pandas 加载。这里用较小的“代表性数据集”演示流程。# 文件路径scripts/train_catboost.py import pandas as pd from catboost import CatBoostRegressor from sklearn.model_selection import train_test_split from sklearn.metrics import mean_absolute_error, mean_squared_error, r2_score import numpy as np # 假设 df_model 已经被转换为 pandas DataFrame data df_model.select( city, year, month, day, pm25, pm25_lag1, pm25_lag7, pm25_7d_avg, so2, no2, co, o3, aqi ).toPandas() # 特征与目标 feature_cols [ city, year, month, day, pm25_lag1, pm25_lag7, pm25_7d_avg, so2, no2, co, o3, aqi ] X data[feature_cols] y data[pm25] # 城市作为类别特征 cat_features [city] X_train, X_test, y_train, y_test train_test_split( X, y, test_size0.2, random_state42 ) model CatBoostRegressor( iterations1000, learning_rate0.05, depth6, loss_functionRMSE, eval_metricRMSE, random_seed42, verbose100 ) model.fit( X_train, y_train, cat_featurescat_features, eval_set(X_test, y_test), use_best_modelTrue )这里use_best_modelTrue表示在验证集指标不再提升时自动停止训练能有效防止过拟合。这是一个很适合初学者的默认选择。6.3 模型评估与保存pred model.predict(X_test) mae mean_absolute_error(y_test, pred) rmse np.sqrt(mean_squared_error(y_test, pred)) r2 r2_score(y_test, pred) print(fMAE: {mae:.2f} μg/m³) print(fRMSE: {rmse:.2f} μg/m³) print(fR2: {r2:.4f}) # 保存模型 model.save_model(catboost_pm25_model.cbm)在正常的空气质量预测项目中如果只有日尺度数据和有限特征RMSE 在 20~35 μg/m³ 之间是比较常见的范围R² 在 0.6~0.85 之间说明模型基本可用。不要期望达到 0.99那是数据泄露的典型信号。6.4 特征重要性分析CatBoost 自带get_feature_importance方法可以直接输出特征重要性。importance model.get_feature_importance() for name, imp in zip(feature_cols, importance): print(f{name}: {imp:.4f})这个输出非常适合放进毕业论文的模型解释章节。通常可以看到pm25_lag1、pm25_7d_avg、month这些特征权重更高这是合理的污染有持续性和季节性。7. 运行结果与效果验证7.1 运行流程总览整个项目的运行顺序是# 1. 启动 Hadoop HDFS start-dfs.sh # 2. 上传原始数据 hdfs dfs -mkdir -p /data hdfs dfs -put air_quality_raw.csv /data/ # 3. 执行数据清洗与分析 spark-submit scripts/clean_data.py spark-submit scripts/analysis.py # 4. 执行特征工程与模型训练 python scripts/feature_engineering.py python scripts/train_catboost.py # 5. 启动 Web 可视化可选 python app.py7.2 验证方法每一步需要检查以下内容阶段检查内容判断标准HDFS 启动jps输出NameNode、DataNode 都存在数据上传hdfs dfs -ls /data文件大小与本地一致清洗完成Parquet 文件是否生成有输出文件且行数合理分析结果Spark SQL 结果表平均值在合理范围模型训练控制台日志RMSE 不报 NaN且有下降趋势预测结果网页或脚本输出预测值在合理区间如果训练时出现 NaN首先检查原始数据是否有极端值再检查滞后特征是否有大量空值。7.3 如果失败先看哪里如果start-dfs.sh后jps没有 DataNode优先看/usr/local/hadoop/logs下的日志常见原因是格式化后 data 目录冲突。如果spark-submit出现jar does not exist or is not a normal file通常是因为SPARK_DIST_CLASSPATH或 Hadoop 环境变量没配对。如果toPandas()报内存错误说明数据量超过了 Driver 内存改用 Spark 算子先聚合或者只抽取部分城市的数据。8. 常见问题与排查思路问题现象可能原因排查方式解决方案Hadoop 启动后 DataNode 进程消失NameNode 格式化和 DataNode 目录冲突查看日志/usr/local/hadoop/logs/hadoop-*.log清空 tmp 目录重新格式化后重启Spark 任务 OOM数据量过大Executor 内存不足查看 Spark UI 中的 Executor 内存增加spark.executor.memory或增加分区数提交中文路径或字段乱码编码不统一检查文件编码和控制台统一使用 UTF-8CSV 加charsetutf-8CatBoost 训练报分类特征错误特征列类型不对打印X[categorical]的 dtype转换为str或category类型toPandas()直接卡死数据量太大查看 Spark UI 内存占用先用limit采样测试或用spark.sql聚合后导出预测结果全部是平均值模型欠拟合或数据无趋势检查特征重要性增加时间滞后特征、月份特征提高迭代次数HDFS 上传慢或失败副本数设置过高检查磁盘空间和配置伪分布式设置dfs.replication19. 最佳实践与工程建议这个选题虽然是毕业设计但建议按真实项目的标准来做。以下几点是很容易拉开档次的细节。9.1 数据分层管理不要只建一个 HDFS 目录。推荐采用数据分层/data/raw 原始数据 /data/clean 清洗后数据 /data/analysis 分析结果 /data/feature 特征数据这样做的好处是每个流程的输出都清晰可见答辩时直接展示目录层级评委一眼就能看出设计思路。9.2 配置项采用参数化Spark 配置和模型参数不要散落在代码里。推荐使用 YAML 或 JSON 配置文件统一管理。# 文件路径config.yml hdfs: raw_path: hdfs://localhost:9000/data/raw clean_path: hdfs://localhost:9000/data/clean feature_path: hdfs://localhost:9000/data/feature spark: app_name: AirQualityPredictor executor_memory: 2g shuffle_partitions: 4 model: iterations: 1000 learning_rate: 0.05 depth: 6 test_size: 0.2 random_seed: 42然后在代码里读取import yaml with open(config.yml, r, encodingutf-8) as f: config yaml.safe_load(f) spark SparkSession.builder \ .appName(config[spark][app_name]) \ .config(spark.executor.memory, config[spark][executor_memory]) \ .getOrCreate()9.3 模型版本管理训练出的.cbm文件要保存版本号和评估指标不要只保留一个最终版本。models/ catboost_pm25_v1.cbm catboost_pm25_v2.cbm model_metrics.txt建议在训练脚本中自动生成model_metrics.txt记录当前模型的特征列表、训练时间、MAE、RMSE、R²。这样在论文中写实验对比时有据可查。9.4 安全与权限注意如果使用了真实监测数据注意数据公开授权范围不要在论文中泄露通过非正规渠道获取的数据来源。在 HDFS 上实践中应该给不同用户设置访问权限不要使用全局写权限。hdfs dfs -chmod -R 755 /data只适用于单机开发环境生产环境必须遵循最小权限原则。9.5 可视化不要只做静态图建议用 FastAPI 或 Flask 搭建一个简单接口接收前端请求后由 Spark 或预加载的模型结果返回数据ECharts 负责渲染。这样系统有前后端交互也更容易在答辩现场演示。在演示时可以先跑 Spark SQL 统计再把预测结果叠加到趋势图上展示。10. 总结与后续学习方向这篇文章完整展示了基于 Spark 的河南省空气质量分析与预测系统的实现路径。从 Hadoop 环境搭建、数据清洗、Spark SQL 统计分析到 CatBoost 建模与特征重要性分析每一步都对应真实项目里需要用到的技能。如果能动手从头跑一遍比只看教程记概念有效得多。当前这个版本已经能够支撑毕业设计答辩但还可以往三个方向继续深入一是实时化。现在的预测基于历史数据如果引入消息队列和流处理框架可以对监测站点上报的实时数据做分钟级预测这就进阶到实时数仓方向了。二是数据仓库化。分析依赖 Spark SQL可以把清洗后的数据进一步建模成数仓分层加入 Hive 或其他 OLAP 工具便于多维分析和报表查询。三是模型深化。除了预测 PM2.5 浓度可以加入空间特征比如把相邻城市距离和风向数据作为特征设计时空图网络模型这属于研究型内容更适合有条件的同学放手尝试。这个题目真正的价值在于它让你在一套完整链路里同时用上了 Hadoop、Spark、机器学习这些核心工具而不是去比较“Spark 和 Hadoop 谁更厉害”或者“哪个算法碾压谁”。把一条路走通比只了解一堆技术名词更能体现在大数据方向上的工程能力。