
在计算机专业毕业设计里基于 Hadoop 和 Spark 的空气质量分析系统一直是个热门选题。它把大数据存储、分布式计算、特征工程、机器学习预测和可视化展示串在同一条链路上比起单纯做 Web 增删改查能更直观地体现“数据工程”的完整流程。本文基于“河南省空气质量数据分析与预测系统”这个选题讲解如何用 Spark 完成海量监测数据的清洗与特征工程用 Hadoop 管理文件存储并引入 CatBoost 完成空气质量等级的预测最终形成一个可运行、可讲解、可扩展的大数据分析闭环系统。这篇文章面向准备做毕业设计的本科生、刚开始接触大数据分析项目的开发者以及对 Spark 与 CatBoost 集成方式感兴趣的人。读者不需要已经具备完整的分布式系统经验但至少应该熟悉 Java、Python 和 MySQL 等基础开发技能。文中给出的代码、配置和命令都以“能够跑通、能够讲清楚设计逻辑”为目标生产环境需要的额外保障会在对应章节单独说明。1. 系统整体设计与技术选型1.1 这个系统解决什么问题河南省地域范围广空气质量监测站点分布在多个城市每天会产生大量包含 PM2.5、PM10、SO2、NO2、O3、CO 浓度以及温度、湿度、风速等气象要素的时序数据。人工查看这些数据只能回答“某天哪个城市空气质量差”很难回答“未来空气质量会出现什么变化趋势”“哪些因素对 PM2.5 浓度影响最大”这类需要建模计算的问题。这个毕业设计系统要完成两件事一是对历史监测数据做统计分析和可视化比如按城市、按月份汇总 AQI 均值二是基于历史特征训练预测模型用未来几天的气象预报数据预测空气质量等级或具体污染物浓度为环保相关展示场景提供辅助判断。1.2 为什么选择 Hadoop、Spark 和 CatBoost这三者不是并列关系而是职责不同组件定位在这个系统中的作用Hadoop HDFS分布式文件存储存放原始监测数据、清洗后数据、特征数据和预测结果Hadoop YARN资源调度为 Spark 计算任务分配 CPU 和内存资源Spark分布式计算引擎批量读取文件、清洗数据、聚合统计、构造特征CatBoost梯度提升决策树模型库基于 Spark 构造好的特征表训练分类或回归模型完成 AQI 等级预测选 Spark 而不只写 Python 脚本核心原因是监测数据量到一定规模后单机 Pandas 读取会慢清洗逻辑也会受内存限制。Spark 支持把数据分片加载到集群内存中执行 map、filter、groupBy、join 等操作代码结构和 Pandas 类似但能天然扩展到多台机器。选 CatBoost 而不是 Spark MLlib 自带的 GBDT主要考虑三点一是题目明确要使用 CatBoost二是 CatBoost 对缺失值和类别特征处理比较友好不需要手动做大量编码三是 CatBoost 在中小特征集上的训练速度和调参体验适合毕业设计阶段快速迭代。需要说明的是如果数据量非常大单机训练 CatBoost 会成为瓶颈这时可以改用 Spark 分布式训练或先将数据采样后再训练本篇以“特征在 Spark 侧完成、模型在 Python 侧训练”作为主线。1.3 系统功能模块划分整个系统可以拆成五个模块数据接入模块从 CSV、JSON 或数据库导出文件中读取河南省各城市空气质量监测数据。数据存储模块原始数据上传至 HDFS按日期或城市分区保存。数据清洗与分析模块Spark 读取 HDFS 文件完成缺失值处理、类型转换、时间字段拆分、城市维度聚合统计。特征工程与建模模块Spark 生成训练特征表Python CatBoost 读取特征表训练模型输出预测结果。结果展示模块将统计结果和预测结果写入 MySQL 或直接导出为 JSON供后端接口与前端图表展示。模块划分的好处是每一部分都能单独验证方便在答辩时按流程演示。比如可以先展示 HDFS 上已经存在哪些原始文件再展示 Spark 清洗后生成的特征表最后切到 Python 环境执行模型训练脚本看到 loss 曲线和预测结果。2. 环境准备集群、运行模式和依赖版本规划2.1 硬件与软件版本规划毕业设计环境通常不会像生产集群那么大但版本选型仍然要提前对齐否则后面会出现 Jar 包冲突、Python 回调失败、HDFS 权限报错等问题。下面是一套可复现的本地/单机环境组合适合先把功能跑通软件推荐版本说明JDK1.8Hadoop 3.x 和 Spark 3.x 对 JDK8 支持最稳妥Hadoop3.3.x比 2.x 简单支持 Java 8兼容性好Spark3.4.x与 Hadoop 3.x 配合良好支持 PySparkPython3.8 或 3.9CatBoost 和 PySpark 都支持CatBoost1.2.x安装时注意与 Python 版本匹配MySQL5.7 或 8.0存储统计结果与预测结果Redis5.x 以上可选用于缓存热点查询结果版本并不是越高越好关键是“验证过的组合”。如果自己所在学校实验室已经有一套 Hadoop 集群要先确认 Spark 版本是否与 Hadoop 的 YARN 兼容再决定提交模式。2.2 Hadoop 伪分布式搭建的关键点在没有多台服务器的情况下伪分布式模式已经足够完成毕业设计的数据存取验证。HDFS 的 NameNode 和 DataNode 都在本机YARN 的 ResourceManager 和 NodeManager 也在本机但目录结构和进程管理方式和集群版一致。安装 Hadoop 后需要修改的配置文件主要有四个!-- 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 /configuration!-- yarn-site.xml -- configuration property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property /configuration# mapred-site.xml 在 Hadoop 3.x 中需要从模板复制 cp etc/hadoop/mapred-site.xml.template etc/hadoop/mapred-site.xml还需要在hadoop-env.sh中显式指定 JDK 路径export JAVA_HOME/usr/local/jdk1.8这里最常见的坑是直接使用系统默认的 openjdk 路径却忘记设置JAVA_HOME启动 hdfs 时会出现类似Error: JAVA_HOME is not set and could not be found的提示。启动前先格式化 NameNodehdfs namenode -format然后执行start-dfs.sh start-yarn.sh用jps查看进程正常情况下能看到NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager五个进程。注意格式化操作会清空 NameNode 上的元数据不要在集群已经保存了重要数据后随意执行。2.3 Spark 安装与部署模式选择Spark 可以以三种模式运行local 模式不连接 YARN直接在本地进程内跑适合调试代码。standalone 模式使用 Spark 自带的 Master/Worker 集群部署简单。yarn 模式由 Hadoop YARN 统一分配资源更接近生产场景。毕业设计建议先使用 yarn 模式因为可以和 HDFS 形成完整链路。安装 Spark 时不需要额外修改太多配置但要确保SPARK_HOME和HADOOP_CONF_DIR环境变量正确export SPARK_HOME/usr/local/spark export HADOOP_CONF_DIR/usr/local/hadoop/etc/hadoop export PATH$PATH:$SPARK_HOME/bin提交任务时如果代码同时读取 HDFS 上的文件需要指定 HDFS 路径spark-submit \ --master yarn \ --deploy-mode client \ --num-executors 2 \ --executor-memory 2g \ --driver-memory 1g \ hdfs://localhost:9000/user/air/scripts/etl_job.py2.4 Python 与 CatBoost 环境在 Python 侧需要安装的基础库包括pip install pyspark3.4.0 pip install catboost pip install pandas numpy scikit-learn pip install pymysql redis这里要特别注意如果 Spark 安装在远程服务器上而 Python 在本地 Windows 运行那么 Spark 提交的 Python 脚本路径和文件路径都要仔细核对。更推荐的做法是 Spark 和 Python 都存放在同一台 Linux 机器或同一套开发环境中减少环境变量不一致带来的问题。3. 数据清洗与特征工程Spark 侧核心实现3.1 原始数据结构设计空气质量数据在设计上通常至少包含以下字段字段名类型说明citystring城市名如郑州、洛阳、开封station_codestring监测站点编码monitor_timestring监测时间格式 yyyy-MM-dd HH:mm:sspm25doublePM2.5 浓度单位 ug/m3pm10doublePM10 浓度单位 ug/m3so2double二氧化硫浓度no2double二氧化氮浓度o3double臭氧浓度codouble一氧化碳浓度temperaturedouble温度单位摄氏度humiditydouble相对湿度单位百分比wind_directionstring风向wind_speeddouble风速如果原始材料没有给出具体数据格式可以按这个通用字段结构准备数据并在论文中说明数据来源与格式约定。实际使用时要先确认站点编码、时间单位等细节避免后续特征拼接时出现口径不一致。3.2 使用 Spark 读取、清洗和聚合数据下面的 PySpark 代码是一个可运行的清洗脚本示例。它读取 HDFS 上的 CSV 文件完成时间字段解析、缺失值过滤、以及按城市和月份统计 PM2.5 均值。from pyspark.sql import SparkSession from pyspark.sql.functions import col, split, concat_ws, avg, month, year spark SparkSession.builder \ .appName(HenanAirQualityETL) \ .enableHiveSupport() \ .getOrCreate() df spark.read \ .option(header, True) \ .option(inferSchema, True) \ .csv(hdfs://localhost:9000/user/air/data/original/henan_air_2023.csv) # 清洗过滤掉关键字段为空的行 df_clean df.filter( col(city).isNotNull() col(pm25).isNotNull() col(monitor_time).isNotNull() ) # 解析时间字段抽取年月 df_clean df_clean.withColumn(year, year(col(monitor_time))) \ .withColumn(month, month(col(monitor_time))) # 按城市和月份聚合 city_month_stats df_clean.groupBy(city, year, month) \ .agg(avg(pm25).alias(avg_pm25)) \ .orderBy(city, year, month) city_month_stats.show(20) # 写入清洗后的结果 df_clean.write.mode(overwrite) \ .parquet(hdfs://localhost:9000/user/air/data/clean/henan_air_2023.parquet)代码包含三个关键设计读取文件时使用inferSchemaTrue让 Spark 自动推断列类型省去手动定义 schema 的代码。但在生产环境中建议显式指定 schema因为自动推断在文件变大或字段缺失时可能出现数据类型漂移。过滤缺失值时使用isNotNull()不能只判断pm25 ! 0因为 PM2.5 在极端天气下可能出现非常低的值但不能用 0 来表示缺失。最终以 Parquet 格式写回 HDFS。Parquet 是列式存储格式后续 Spark 读取时只加载需要的列能明显减少磁盘和网络开销。3.3 构造模型训练特征表模型预测不能只使用当前小时的污染物浓度还需要构造延迟特征和统计特征。例如预测第二天某个城市的空气质量等级可以使用过去 24 小时的平均 PM2.5、最大 PM2.5、风速均值、湿度均值等。实现思路是把原始数据按城市和时间排序用 Spark SQL 的窗口函数生成滞后特征。from pyspark.sql.window import Window from pyspark.sql.functions import lag window_spec Window.partitionBy(city).orderBy(monitor_time) feature_df df_clean.withColumn(pm25_lag1, lag(pm25, 1).over(window_spec)) \ .withColumn(pm25_lag24, lag(pm25, 24).over(window_spec)) \ .withColumn(tmp_lag1, lag(temperature, 1).over(window_spec)) \ .withColumn(hum_lag1, lag(humidity, 1).over(window_spec))窗口函数lag取同一个城市内前 1 行或前 24 行的值。这样每条记录就携带了历史状态CatBoost 训练时可以利用这些时序特征判断趋势。构造完特征后输出到 HDFS 或本地供 Python 读取feature_df.select( city, year, month, day, pm25, pm25_lag1, pm25_lag24, tmp_lag1, hum_lag1, wind_speed ).write.mode(overwrite) \ .csv(/tmp/air_features.csv, headerTrue)输出时要注意如果 HDFS 上的输出目录已经存在直接写会报错所以上面使用了mode(overwrite)覆盖写入。实际项目中如果要保留历史特征建议改为按时间分区写入。4. CatBoost 建模与预测从特征到结果4.1 为什么用 CatBoost 而不是普通 XGBoostCatBoost 是俄罗斯搜索公司 Yandex 开源的梯度提升决策树库。它在处理类别特征、缺失值和防止过拟合方面做得比较完善。对空气质量数据来说城市名、风向这类字段都是类别特征如果手动用 one-hot 编码会比较繁琐而 CatBoost 可以直接声明类别特征列并在内部做有序编码处理简化了特征工程代码。在模型中可以把问题定义为多分类预测 AQI 等级等级包括优、良、轻度污染、中度污染、重度污染、严重污染。也可以定义为回归直接预测 PM2.5 浓度。本文以回归为例讲解分类任务只需要更换损失函数和评估指标。4.2 读取特征数据并训练模型在 Python 侧只需要用 Pandas 读取 Spark 输出的特征文件然后交给 CatBoost 训练。import pandas as pd import catboost as cb from catboost import CatBoostRegressor, Pool from sklearn.model_selection import train_test_split from sklearn.metrics import mean_squared_error, mean_absolute_error df pd.read_csv(/tmp/air_features.csv) # 只选取模型需要的特征列 features [ pm25_lag1, pm25_lag24, tmp_lag1, hum_lag1, wind_speed, city ] X df[features] y df[pm25] X_train, X_val, y_train, y_val train_test_split( X, y, test_size0.2, random_state42 ) train_pool Pool(X_train, y_train, cat_features[city]) val_pool Pool(X_val, y_val, cat_features[city]) model CatBoostRegressor( iterations500, learning_rate0.05, depth6, loss_functionRMSE, eval_metricMAE, random_seed42, verbose100 ) model.fit( train_pool, eval_setval_pool, use_best_modelTrue, early_stopping_rounds50 )这里要注意cat_features[city]。如果不声明类别特征CatBoost 会把城市名当作普通字符串处理导致训练时间变长或模型效果变差。训练完成后用验证集查看误差y_pred model.predict(val_pool) rmse mean_squared_error(y_val, y_pred, squaredFalse) mae mean_absolute_error(y_val, y_pred) print(fRMSE: {rmse:.2f}) print(fMAE: {mae:.2f})RMSE 对较大误差更敏感MAE 反映平均误差水平。对空气质量预测来说如果 MAE 在 10 到 15 ug/m3 左右已经具备较强的参考价值但如果误差超过 30则要检查特征是否包含足够的历史信息比如是否缺少前一天同一时刻的污染物浓度这一重要特征。4.3 模型保存与预测输出将训练好的模型保存下来后续可以独立运行预测脚本不需要重新加载训练数据。model.save_model(/tmp/catboost_air_model.cbm)预测脚本可以单独写成predict.pyimport pandas as pd import catboost as cb model cb.CatBoostRegressor() model.load_model(/tmp/catboost_air_model.cbm) # 假设 Spark 已经生成了未来一天的特征集 future_features pd.DataFrame([ { pm25_lag1: 52.0, pm25_lag24: 60.0, tmp_lag1: 8.0, hum_lag1: 45.0, wind_speed: 2.1, city: 郑州 } ]) future_features[city] future_features[city].astype(str) pred model.predict(future_features) print(f预测 PM2.5 浓度: {pred[0]:.2f} ug/m3)预测结果可以写回 MySQL也可以输出为 JSON 文件供前端展示。这里把预测脚本和训练脚本分离是为了保持模块独立避免在预测阶段重复调用 Spark Session。5. 运行验证从 HDFS 文件到预测结果5.1 完整运行流程整个系统的运行流程可以归纳为一条命令链上传原始数据到 HDFShdfs dfs -mkdir -p /user/air/data/original hdfs dfs -put /home/student/data/henan_air_2023.csv /user/air/data/original/提交 Spark 清洗任务spark-submit --master yarn eda_and_clean.py确认 HDFS 上生成了 Parquet 清洗文件和 CSV 特征文件hdfs dfs -ls /user/air/data/clean/ hdfs dfs -ls /tmp/air_features.csv如果特征文件只用于本地 Python 训练可以不用写 HDFS而是使用 Spark 在本地输出时指定file:///tmp/air_features.csv。这里要留意 Spark 在 yarn 模式下的输出路径直接写/tmp/air_features.csv可能不会落在你想找的服务器路径上最好显式加上file://前缀。执行 Python 训练脚本python3 train_catboost.py执行预测脚本python3 predict.py5.2 验证点与预期结果完成环境搭建后不要只确认“能启动”“能跑完”还要针对每个阶段确认输出是否合理验证阶段验证命令或方法预期结果HDFS 状态hdfs dfsadmin -reportDataNode 正常存储容量可见Spark SQL 聚合结果脚本中的show(20)能展示郑州、洛阳等城市的月份均值特征文件字段数wc -l /tmp/air_features.csv与原始记录数接近CatBoost 训练训练日志每次迭代打印 MAEMAE 随迭代次数下降预测结果python3 predict.py输出合理 PM2.5 浓度区间比如 40-80 之间如果预测结果出现负数或数量级异常首先要检查特征是否被错误缩放。空气质量浓度本身是物理量不可能小于 0模型出现负值说明训练集存在异常值或数据泄漏。5.3 可视化展示思路毕业设计通常需要展示页面。可以采用的轻量方案是Spark 聚合统计结果写回 MySQL后端使用 Spring Boot 提供查询接口前端使用 ECharts 绘制折线图、柱状图和地图。如果不希望引入 Web 框架也可以直接用 Jupyter Notebook 展示图表但答辩时建议至少准备一张系统结构图和一张预测效果对比图例如横轴为日期纵轴为 PM2.5 实际值与预测值的折线对比。6. 常见问题排查环境、提交和模型训练阶段6.1 Spark 提交任务时提示 jar 不存在报错示例Error: Could not find or load main class org.apache.spark.launcher.Main Caused by: java.io.IOException: jar does not exist or is not a normal file: /usr/local/hadoop/share/hadoop/m...这个报错经常出现在切换 Hadoop 或 Spark 版本后。原因是SPARK_HOME或HADOOP_CONF_DIR指向了不存在的路径或者 spark-env.sh 中配置的HADOOP_HOME路径与实际安装路径不一致。检查顺序执行echo $HADOOP_HOME和echo $SPARK_HOME确认路径存在。查看 Spark 目录下的conf/spark-env.sh确认没有写死错误的HADOOP_HOME。如果使用spark-submit --master yarn还需要确认 Hadoop 的 classpath 配置正确。可以在 hadoop 目录下执行hadoop classpath再把输出内容追加到 spark-env.sh 的SPARK_DIST_CLASSPATH中。6.2 Spark OOM 问题Spark 任务处理大数据量时常见报错包括ExecutorLostFailure、java.lang.OutOfMemoryError或Container killed by YARN for exceeding memory limits。排查路径确认 driver 内存和 executor 内存是否分配过小。单机 8GB 内存时executor 设置 2gdriver 设置 1g不要超过物理内存。检查代码中是否有.collect()把全量数据拉回 driver。collect会把所有分区的数据集中到一台机器上数据量大时必然 OOM。检查是否存在数据倾斜。如果某个城市的数据量远大于其他城市聚合时该分区的负载过高可以增加spark.sql.shuffle.partitions或使用repartition(city)重新分布数据。6.3 CatBoost 训练速度慢或内存占用高在数据量较大时CatBoost 默认参数可能较慢。可以采用以下措施设置iterations500而不是默认的更高值。减小depth6或depth4模型复杂度下降训练和推理都会变快。设置thread_count4限制 CPU 并发避免挤占 Spark 资源。如果特征表非常大先做一次采样再训练。毕业设计阶段并不需要把所有数据全部投入训练抽样后模型效果通常仍可接受。6.4 HDFS 写入权限不足当使用hdfs dfs -put上传文件到根目录时可能出现Permission denied报错。最简单的方式是使用当前用户创建自己的目录hdfs dfs -mkdir -p /user/air hdfs dfs -chown -R $USER /user/air注意不要在生产环境随意把目录权限改成 777。开发环境为了调试方便可以这样做但生产环境应该按用户分组授权。7. 最佳实践与扩展方向7.1 代码与工程结构规范按照模块化方式组织代码后续调试和写论文都更方便。一个推荐的项目结构是air-quality-project/ ├── data/ │ ├── original/ # 原始数据 │ └── feature/ # 特征数据与预测结果 ├── scripts/ │ ├── etl_job.py # Spark 清洗任务 │ ├── feature_engineer.py # Spark 特征工程 │ ├── train_catboost.py # 模型训练 │ └── predict.py # 预测脚本 ├── sql/ │ └── init_table.sql # MySQL 建表语句 └── docs/ └── architecture.md # 系统设计说明这种结构的好处是每个脚本只负责一段职责可以从 ETL、特征、训练、预测四个环节分别测试。写论文时也能直接引用每个文件的实现说明。7.2 部署与运行的环境隔离开发时可以在本地运行 Spark local 模式部署时切到 yarn 模式。两类环境下的路径和资源配置差异较大建议做一个配置常量文件# config.py IS_YARN True if IS_YARN: HDFS_BASE hdfs://localhost:9000/user/air SPARK_MASTER yarn else: HDFS_BASE file:///tmp/air SPARK_MASTER local[*]这样切换环境时不用改每一行路径减少漏改导致的文件找不到问题。7.3 扩展方向完成基础版之后可以从以下几个方向继续扩展预测范围扩展从单一城市扩展到多个城市并使用模型解释工具分析哪些气象特征对 PM2.5 影响最大。实时预测将批处理改成短周期调度比如每 1 小时从 Kafka 读取新数据Spark Structured Streaming 做实时特征更新模型定期重新训练。引入更多数据源加入气象预报接口或卫星遥感数据提高预测准确率。结果存储优化统计结果和预测结果写入 MySQL 后对高频查询增加 Redis 缓存减轻数据库压力。模型服务化使用 Flask 或 FastAPI 封装模型接口让 Web 端通过 REST API 直接获取预测结果比每次执行 Python 脚本更易维护。7.4 对毕业设计答辩的建议答辩时不需要把全部代码讲一遍重点讲清楚三件事数据从哪来、数据怎么变成特征、模型如何基于特征做预测。演示时按照“原始文件 - Spark 清洗 - 特征表 - CatBoost 预测结果 - 可视化图表”这条链路走每个环节保留一条命令和一次输出截图即可。遇到评委提问时重点说明 Spark 与 CatBoost 的分工以及为什么特征工程发生在 Spark 侧而不是 Pandas 侧这比背概念更有说服力。总体来看这套基于 Hadoop、Spark 和 CatBoost 的空气质量分析预测系统技术链路完整、模块边界清晰、扩展空间大非常适合作为大数据方向毕业设计的选题骨架。只要把原始数据结构摸清楚再把环境版本固定住整个项目的实现周期可以控制在一个月以内。做完之后你不仅能讲清楚 Spark 的 RDD 和 DataFrame 区别、HDFS 的读写机制、CatBoost 的类别特征处理逻辑还能在这套代码基础上继续做实时预测或模型优化后续职业方向转向数据工程或机器学习都会有直接帮助。