Spark+Hadoop+CatBoost:河南省空气质量分析与预测毕设全攻略
如果你是计算机专业的学生正在为毕业设计选题发愁看到“基于Spark的河南省空气质量数据分析与预测系统”这种题目时大概率会有两种反应一是觉得“Spark Hadoop 机器学习”这套技术栈太沉担心自己撑不起来二是觉得这不就是把数据丢给某个模型训练一下没什么技术含量。这两种判断都只看到了局部。这个题目真正有价值的地方在于它把一个看起来很高门槛的大数据项目拆成了三条能独立验证的链路——用Hadoop存数据、用Spark算数据、用CatBoost算法预测数据。每一条链路都有明确产出合在一起又是一个完整的“数据平台 业务分析 算法建模”叙事。对毕业设计来说这比单纯调一个深度学习模型或者做一个增删改查管理系统要稳得多。这篇文章不只是给你一个题目介绍。我会从选题判断、技术选型、系统架构、环境搭建、数据清洗、特征工程、CatBoost训练、结果验证、答辩准备这条完整路线把这个项目讲清楚。文章里会给出可以直接复用的PySpark处理和CatBoost训练代码也会指出哪些地方是评分看点、哪些地方是坑。1. 选题值不值为什么“大数据 预测”是毕业设计的稳妥组合先说结论如果必须在一个月内完成一个大数据方向的毕设这个题目的性价比属于第一梯队。大多数学校的毕业设计评分核心就几件事选题是否有实际意义、工作量是否饱满、技术方案是否合理、结果是否可验证、答辩是否能讲清楚。用这个标准去套这个题目你会发现它天然踩中了所有得分点。空气质量是国家环保公开数据业务意义天然成立不需要额外编故事Spark分布式处理能体现平台类工作量Hadoop HDFS存储能体现大数据生态的完整度CatBoost与基线模型的对比能体现算法层面的思考最终用RMSE、R²和可视化图来展示预测效果结果可验证。这些正好覆盖了评分维度。但这个题目的挑战也很明显不要把它做成“算法调包”项目。很多学生拿到题目就开始训模型忽略了数据平台部分。既然题目里带Spark和Hadoop就要让评委看到分布式处理的必要性。哪怕数据量不到几十GB也要在项目中设计合理的分布式处理链路并解释清楚为什么选择Spark而不是Pandas。这是答辩时最能体现区分度的地方。为什么特别适合河南省这种省级数据因为省级多城市、多站点的数据天然具备地理维度和时间维度做站点对比、区域差异、月度变化时内容会很丰富不至于只有一张折线图。数据规模虽然不会像互联网日志那么大但用来讲清楚分布式处理流程已经足够。2. 项目要解决什么从空气质量分析到空气质量预测把题目拆成两个关键词数据分析和预测。分析部分要回答业务问题河南省不同城市的空气质量逐月趋势如何污染物浓度是否存在周末效应哪些站点PM2.5常年偏高这些问题通过Spark SQL做聚合统计就可以快速得到答案得到的表格和图表可以直接放进论文的“数据分析”章节。预测部分要回答技术问题已知当前时刻的污染物浓度、站点信息和时间特征下一小时的PM2.5是多少这是一个典型的表格型回归问题。目标字段是PM2.5特征是当前各污染物浓度、小时、星期、月份、站点类别以及过去一段时间的统计量。为什么预测PM2.5而不是AQI从工程角度看PM2.5是单项污染物中数据最完整、业务关注度最高的指标AQI是综合指数预测维度更多误差解释也更复杂。做毕业设计时建议先以PM2.5作为预测目标后续可以把AQI和其余污染物作为扩展方向。这一步就足够支撑一篇论文的核心实验了。还有一个容易在答辩时被问倒的问题为什么不直接用ARIMA或者LSTM你要能回答清楚。ARIMA适合单变量平稳序列但很难纳入站点类型、小时、节假日等额外特征LSTM需要大量连续时序数据调参成本高对表格特征的支持不如树模型灵活。CatBoost是梯度提升树模型训练快、支持类别特征、对缺失值容忍度高、可解释性好。对于小时级环境监测数据它往往能在较短时间内获得不错的基线效果。这个对比逻辑答辩时就是你的护城河。3. 技术栈选型Hadoop、Spark、CatBoost是如何分工的很多初学者容易把Hadoop和Spark搞混或者是觉得“有了Spark就不需要Hadoop”。实际上在这个项目里三者各管一段。Hadoop解决的是存储问题。HDFS是数据落地的仓库存放原始CSV文件和处理后的结果文件。它承担的是“数据不会丢、能分布式存”的职责。在毕业设计里体现为“数据文件上传到HDFS”和“分布式存储思想”。Spark解决的是计算问题。它负责读取HDFS上的数据完成数据清洗、缺失值处理、聚合统计和特征计算。Spark的核心优势是内存计算加上与HDFS无缝衔接天然适合做“Hadoop存储 Spark计算”的组合。这个项目里Spark是整个数据管道的中枢。CatBoost解决的是学习问题。它读取Spark产出的特征表训练回归模型然后对下一小时的PM2.5做预测。CatBoost是Yandex开源的梯度提升库在表格型数据上表现稳定。为什么选Spark而不是Pandas可以做一个简单对比对比维度PandasSpark数据规模受单机内存限制可以扩展到多机集群分布式计算不支持原生支持与Hadoop HDFS集成需要额外库原生支持学习成本较低中等但概念值得学适合场景单机数据分析和实验大规模数据处理、生产管道为什么选CatBoost而不是XGBoost或LightGBM这三个都是优秀的GDBT框架但侧重点不同算法类别特征处理缺失值处理训练速度可解释性XGBoost需要手动编码需要预处理较快好LightGBM需要手动编码支持很快好CatBoost原生支持原生支持快好内置特征重要性对于空气质量数据里的“站点名称”这类高基数类别特征CatBoost可以省掉大量独热编码工作同时降低过拟合风险。做毕业设计时时间有限选择CatBoost更容易在短时间内跑出稳定结果答辩时也有“为什么选它”的明确理由。4. 系统总体架构与核心数据流一个规范的毕业设计系统架构不能只是拼凑代码。这里建议采用分层架构让评委一眼看到系统全貌数据层河南省各城市监测站点的小时级污染物数据以CSV或Parquet格式存储上传到HDFS。计算层Spark Session读取原始数据完成清洗、异常值处理、时间特征构造、滑动窗口统计。模型层CatBoost读取处理后的特征数据训练回归模型保存模型文件。应用层评估脚本输出RMSE、MAE、R²绘制真实值与预测值对比图可选通过Flask提供一个预测接口。核心数据流可以这样描述原始监测文件 - HDFS - Spark DataFrame - 数据清洗 - 特征工程 - Parquet/CSV - CatBoost Pool - 模型训练 - 模型文件 - 预测与可视化。架构成型之后论文的“系统设计”章节也就有了骨架。你可以把每一层对应到一章把数据流图画到论文中逻辑非常清晰。这里要提醒一点毕业设计文档和代码不同代码可以写得随意但文档一定要体现“为什么这样分层”。每层之间通过什么接口传递数据数据格式怎么约定这些在答辩时都是加分项。5. 环境准备从Hadoop到Spark再到CatBoost环境搭建是很多人第一个崩溃的环节。我的建议是先从单机开始不要一上来就折腾集群。你是为了完成毕设不是为了考Hadoop运维工程师。单机Local模式足以跑通完整流程等数据链路稳定后再根据时间和机器条件决定要不要扩展成Spark集群。推荐环境清单操作系统Linux发行版例如Ubuntu Server也可以用虚拟机或云服务器。JDK版本以Hadoop官方支持版本为准常见是JDK 8或JDK 11。Hadoop伪分布式模式主要使用HDFS。SparkLocal模式后续可切换YARN模式。Python3.8以上用Anaconda管理环境。依赖库pyspark、catboost、pandas、pyarrow、scikit-learn、flask。环境变量配置可以参考下面这段请根据自己的实际安装路径调整# 文件路径~/.bashrc export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 export HADOOP_HOME/usr/local/hadoop export SPARK_HOME/usr/local/spark export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$SPARK_HOME/bin export PYSPARK_PYTHON/opt/anaconda3/envs/air-quality/bin/python export PYSPARK_DRIVER_PYTHON/opt/anaconda3/envs/air-quality/bin/python配置完成后创建Python环境并安装依赖conda create -n air-quality python3.9 -y conda activate air-quality pip install pyspark catboost pandas pyarrow flask scikit-learn # 启动Hadoop伪分布式集群HDFS start-dfs.sh # 查看HDFS状态 hdfs dfsadmin -report # 创建数据目录并上传数据 hdfs dfs -mkdir -p /data/air hdfs dfs -put henan_air_quality.csv /data/air/如果你对Hadoop不太熟可以先跳过HDFS用local模式跑Spark读取本地CSV文件。这个顺序不丢分反而能让项目在短时间内稳定交付。把数据管道打通之后再把“上传到HDFS”作为加分步骤补上。环境安装阶段最容易出错的是版本兼容问题尤其是PySpark、Hadoop和JDK之间的版本匹配。遇到问题时第一时间去查Spark官方文档中“版本兼容”表而不是盲目升级组件。6. 基于Spark的数据清洗与特征工程实现预处理是整个项目中工作量最大、也最能体现工程能力的部分。常见的监测数据会包含以下字段站点名称、时间、AQI、PM2.5、PM10、SO2、NO2、O3、CO。实际数据可能还包含经纬度等额外字段但建模阶段用不到这么多。字段名请以实际数据为准下面代码以通用格式为例。清洗逻辑要做四件事统一时间格式把字符串时间解析成Spark timestamp类型。删除重复记录同一站点同一时刻只保留一行。处理异常值污染物浓度出现负值可能是设备故障应置为NULL。填充缺失值使用同一站点前一时刻的观测值填充当前缺失。特征工程方面构造这些特征时间特征小时、星期几、月份。滞后特征上一小时PM2.5、过去6小时PM2.5均值、过去24小时PM2.5均值。当前状态当前时刻PM10、SO2、NO2、O3、CO浓度。类别特征站点名称。预测标签下一小时PM2.5值。完整代码如下# 文件路径src/data_processing.py from pyspark.sql import SparkSession from pyspark.sql.functions import ( col, to_timestamp, when, coalesce, lag, lead, avg, hour, dayofweek, month ) from pyspark.sql.window import Window spark SparkSession.builder \ .appName(HenanAirQualityETL) \ .config(spark.sql.session.timeZone, Asia/Shanghai) \ .getOrCreate() # 1. 读取原始数据 df spark.read \ .option(header, True) \ .option(inferSchema, True) \ .csv(data/henan_air_quality.csv) # 2. 统一列名并转成数值类型 cleaned df.select( col(station).alias(station), to_timestamp(col(time), yyyy-M-d H:mm).alias(datetime), col(PM2.5).cast(double).alias(pm25), col(PM10).cast(double).alias(pm10), col(SO2).cast(double).alias(so2), col(NO2).cast(double).alias(no2), col(O3).cast(double).alias(o3), col(CO).cast(double).alias(co) ) # 3. 删除时间为空的行负浓度置为NULL cleaned cleaned.filter(col(datetime).isNotNull()) for c in [pm25, pm10, so2, no2, o3, co]: cleaned cleaned.withColumn( c, when(col(c) 0, None).otherwise(col(c)) ) # 4. 按站点分组按时间排序用前一时刻值填充缺失 win Window.partitionBy(station).orderBy(datetime) fill_cols [pm25, pm10, so2, no2, o3, co] for c in fill_cols: cleaned cleaned.withColumn( c, coalesce(col(c), lag(col(c), 1).over(win)) ) # 5. 去除重复记录 cleaned cleaned.dropDuplicates([station, datetime]) # 6. 构造时间特征 feature_df cleaned.withColumn(hour, hour(col(datetime))) feature_df feature_df.withColumn(day_of_week, dayofweek(col(datetime))) feature_df feature_df.withColumn(month, month(col(datetime))) # 7. 构造滞后特征和滑动平均 win_prev_1h Window.partitionBy(station).orderBy(datetime) win_avg_6h Window.partitionBy(station) \ .orderBy(datetime) \ .rowsBetween(-6, -1) win_avg_24h Window.partitionBy(station) \ .orderBy(datetime) \ .rowsBetween(-24, -1) feature_df feature_df.withColumn(pm25_prev_1h, lag(pm25, 1).over(win_prev_1h)) feature_df feature_df.withColumn(pm25_avg_6h, avg(pm25).over(win_avg_6h)) feature_df feature_df.withColumn(pm25_avg_24h, avg(pm25).over(win_avg_24h)) # 8. 构造预测标签下一小时PM2.5 label_win Window.partitionBy(station).orderBy(datetime) feature_df feature_df.withColumn(pm25_next_hour, lead(pm25, 1).over(label_win)) # 9. 删除标签或特征缺失的行落地为Parquet final_df feature_df.dropna(subset[pm25_next_hour]) final_df.write.mode(overwrite).parquet(data/henan_air_features)这段代码里最值得讲清楚的是窗口函数。Window.partitionBy(station).orderBy(datetime)定义了一个“按站点分组、按时间排序”的滑动窗口lag取上一行lead取下一行avg在窗口内做均值。这样能精确控制“只用过去信息预测未来”避免数据泄漏。需要特别留意如果原始数据存在整小时缺失rowsBetween(-24, -1)取到的是前24行而不是前24个小时。实际项目中要根据数据完整度调整窗口边界或者在预处理阶段先按站点小时做一次补全。这个细节如果能在答辩时主动讲出来会显得你对数据工程有真实理解。数据处理完成后会得到一个Parquet文件data/henan_air_features。这个文件就是CatBoost训练阶段的输入。如果你的机器资源紧张也可以在Spark里直接调用.toPandas()转成Pandas DataFrame但前提是数据量不大且你明确知道这样做的代价是数据全部进入单机内存。7. CatBoost模型训练与预测实现训练阶段的关键不是代码有多花哨而是数据集的划分方式。对时间序列预测严禁随机划分训练集和测试集。正确做法是按时间排序前80%作为训练集后20%作为测试集。这样才能验证模型对“未来未见数据”的预测能力。如果随机划分模型很可能“见到过”测试集附近的样本指标虚高答辩时经不起追问。下面是完整的CatBoost训练脚本# 文件路径src/train_model.py import pandas as pd import numpy as np from catboost import CatBoostRegressor, Pool from sklearn.metrics import mean_absolute_error, mean_squared_error, r2_score # 读取Spark处理后的特征数据 df pd.read_parquet(data/henan_air_features) df df.sort_values([station, datetime]).reset_index(dropTrue) features [ station, hour, day_of_week, month, pm25_prev_1h, pm25_avg_6h, pm25_avg_24h, pm10, so2, no2, o3, co ] cat_features [station] # 按时间顺序划分前80%训练后20%测试 split_idx int(len(df) * 0.8) train_df df.iloc[:split_idx] test_df df.iloc[split_idx:] train_pool Pool( train_df[features], train_df[pm25_next_hour], cat_featurescat_features ) test_pool Pool( test_df[features], test_df[pm25_next_hour], cat_featurescat_features ) # 初始化CatBoost回归模型 model CatBoostRegressor( iterations1000, learning_rate0.05, depth8, loss_functionRMSE, eval_metricRMSE, random_seed42, od_wait50, verbose100 ) # 训练 model.fit(train_pool, eval_settest_pool) # 预测 pred model.predict(test_pool) # 评估 mae mean_absolute_error(test_df[pm25_next_hour], pred) rmse float(np.sqrt(mean_squared_error(test_df[pm25_next_hour], pred))) r2 r2_score(test_df[pm25_next_hour], pred) print(fMAE : {mae:.2f}) print(fRMSE: {rmse:.2f}) print(fR2 : {r2:.4f}) # 保存模型 model.save_model(models/catboost_air_model.cbm)几个关键参数说明iterations树的数量。设为1000配合od_wait可以在验证集指标不再提升时提前停止。learning_rate学习率。0.05是相对保守的默认值越小越稳但需要更多树。depth树深度。8对这类中小规模表格数据已经够用过深容易过拟合。loss_functionRMSE回归任务常用损失。cat_features[station]告诉CatBoost哪些列是类别特征无需手动独热编码。训练完成后控制台会周期性输出RMSE变化。你可以观察验证集RMSE是否持续下降如果连续50轮没有改善训练会自动停止。最终输出的MAE、RMSE、R²就是论文实验部分的核心数据。对于毕业设计R²在0.8以上已经是相当不错的预测效果如果低于0.6优先检查特征构造和数据量而不是盲目调参。比如确认滞后特征是否正确生成、缺失值是否填充合理、训练集和测试集是否有重叠。这些排查思路比单纯改模型参数可靠得多。8. 结果评估与系统展示模型训练完不能只在控制台打印几个数字就结束。一个完整的毕业设计系统需要有可视化结果和可演示的界面。三个评估指标的含义要写清楚MAE平均绝对误差预测值与真实值的平均绝对差距单位是μg/m³直观易懂。RMSE均方根误差对大误差更敏感能反映预测是否在某些时刻严重偏离。R²决定系数模型解释了多少比例的方差越接近1越好。可视化方向可以包括测试集时间段内真实值与预测值的时间序列对比曲线。真实值-预测值散点图越接近对角线说明预测越准。CatBoost输出的特征重要性柱状图说明哪些特征对PM2.5影响最大。此外可以加一个简单的Flask预测接口让系统具备“输入特征、返回预测值”的能力。虽然是加分项但演示效果很好。代码如下# 文件路径src/app.py from flask import Flask, request, jsonify from catboost import CatBoostRegressor import pandas as pd app Flask(__name__) model CatBoostRegressor() model.load_model(models/catboost_air_model.cbm) app.route(/predict, methods[POST]) def predict(): data request.get_json() df pd.DataFrame([data]) pred model.predict(df)[0] return jsonify({predicted_pm25_next_hour: float(pred)}) if __name__ __main__: app.run(host0.0.0.0, port5000)启动后用curl或Postman发送POST请求即可获得预测值curl -X POST http://localhost:5000/predict \ -H Content-Type: application/json \ -d {station:某站点,hour:14,day_of_week:3,month:5,pm25_prev_1h:35.0,pm25_avg_6h:40.0,pm25_avg_24h:45.0,pm10:60.0,so2:8.0,no2:30.0,o3:80.0,co:0.6}如何判断系统是否成功了标准是模型文件能正常加载接口能返回数值且返回值的量级在合理范围内。比如下一小时PM2.5预测为35而不是负值或者几千。如果出现异常预测优先检查请求字段名称是否与训练特征一致以及是否有缺失字段。9. 常见问题与排查方法这个项目在实现过程中最容易踩的坑集中在环境、数据、内存三方面。整理成表格方便对照排查问题现象可能原因排查方式解决方案Spark任务OOMexecutor内存配置不足或数据量超过单机内存查看Spark UI中的Executor内存和GC日志调整spark.executor.memory或先用local模式跑小数据量Hadoop NameNode无法启动未格式化NameNode或格式化后数据目录冲突查看