每年到毕业设计选题季,很多同学都会纠结一个问题:想做大数据方向,又担心纯理论不好做;想带上机器学习预测,又怕自己基础不够。如果你也处于这个阶段,可以考虑“基于 Spark 的河南省空气质量数据分析与预测系统”这个方向。它把 Spark 海量数据处理、Hadoop 分布式存储、机器学习预测三部分串起来,既能体现工程能力,又能体现建模能力,是一个很典型的“大数据分析 + 算法预测”综合项目。
这篇文章会从项目背景、技术选型、系统设计、环境搭建、代码实现、模型训练到常见 Bug 排查,完整走一遍流程。不管你是准备做毕业设计,还是想入门大数据分析实战,都可以直接参考这套思路落地。
1. 项目背景与核心技术选型
1.1 项目背景
河南省是人口大省,也是工业和农业大省,空气质量受季节、气象、区域传输等多重因素影响。公众关心“明天要不要带口罩”,环保部门关心“哪个区域需要重点管控”,而数据分析项目关心的是:能不能用历史空气质量数据,结合气象特征,提前预测未来一天的 PM2.5 浓度。
传统方式用 Excel 或单机 Pandas 处理几万条数据没有问题,但如果数据量达到几百 GB,或者需要按城市、按日期、按小时做复杂关联统计,单机工具就会力不从心。这时候就可以引入 Spark 来做分布式数据处理,再用 CatBoost 这类梯度提升树模型完成回归预测。
这个选题最大的优势在于:技术栈完整、问题明确、结果可量化。你可以把项目拆成“数据分析”和“预测建模”两条主线,任何一条线都能独立产出结果,适合论文的章节编排。
1.2 为什么选择 Spark + Hadoop + CatBoost
很多同学一开始会困惑:Spark、Hadoop、MapReduce、Hive 到底有什么区别?这里做一个简单区分。
- Hadoop:分布式存储(HDFS)和分布式计算框架,生态中包含 MapReduce、Hive、HBase 等组件。它对海量数据的存储非常友好。
- Spark:基于内存计算的分布式计算引擎,比 MapReduce 更快的迭代计算性能,提供 Spark SQL、Spark Streaming、MLlib 等模块。
- MapReduce:Hadoop 中的批处理计算模型,编程复杂度较高,适合离线清洗,但写起来偏繁琐。
- CatBoost:Yandex 开源的梯度提升决策树(GBDT)框架,擅长处理表格数据,对类别特征和缺失值有原生支持。
在这个项目中,合理的分工是:
| 组件 | 承担职责 |
|---|---|
| HDFS | 存储原始 CSV 数据和分析结果 |
| Spark | 数据清洗、特征加工、统计分析 |
| Pandas | 接收 Spark 处理后的特征数据,做轻量转换 |
| CatBoost | 基于特征构建 PM2.5 浓度预测模型 |
| Matplotlib / Flask(可选) | 结果可视化和简单 Web 展示 |
有人可能会问:为什么不直接用 Spark MLlib 做预测?Spark MLlib 确实支持随机森林、线性回归、GBDT,适合海量数据的并行训练。但毕业设计需要一个更大的加分点,CatBoost 在中小规模表格数据上通常能获得更高精度,而且对类别特征处理更友好,模型解释工具也更完善。所以采用“Spark 负责大数据处理,CatBoost 负责精细建模”的混合架构更为合理。
1.3 系统能力概览
本系统预期具备以下能力:
- 从 HDFS 读取河南省各地市历史空气质量数据。
- 使用 Spark SQL 完成数据清洗、去重、缺失值处理。
- 使用 Spark 窗口函数构造“前一日 PM2.5”“前两日 AQI”等滞后特征。
- 统计各城市空气质量等级分布、月度 AQI 变化趋势。
- 利用 CatBoost 训练次日 PM2.5 回归预测模型。
- 输出模型评估指标,并保存模型文件供后续调用。
2. 系统总体架构设计
2.1 分层架构
本系统按数据流向可以分成四层:数据接入层、数据存储层、数据处理分析层、模型训练与预测层。
- 数据接入层:将采集到的河南各地市空气质量 CSV 文件上传到 HDFS。
- 数据存储层:使用 Hadoop HDFS 存放原始数据,使用本地文件系统存放中间特征结果。
- 数据处理分析层:使用 PySpark 完成数据清洗、聚合统计和特征工程,通过 Spark SQL 输出分析结果。
- 模型训练与预测层:将 Spark 处理好的特征数据导出为 Pandas DataFrame,使用 CatBoost 训练回归模型,最后保存模型文件。
之所以让特征数据从 Spark 落到本地再训练 CatBoost,是为了减少环境复杂度。如果你是本地单机模式运行 Spark,数据量在千万行以内,这个方法最稳。
2.2 数据流说明
完整的数据流如下:
- 原始 CSV 数据通过
hdfs dfs -put命令上传到 HDFS。 - PySpark 读取 HDFS 上的 CSV 文件,利用 DataFrame API 做类型转换和过滤。
- 对清洗后的数据按城市和日期排序,使用
lag窗口函数生成滞后特征。 - 缺失值处理后,选择特征列并转换为 Pandas DataFrame。
- 划分训练集与测试集,训练 CatBoost 回归模型。
- 评估 RMSE、MAE、R2 指标,保存模型和预测结果。
该流程不需要额外的消息队列和实时计算组件,适合作为毕业设计展示,也能说明清楚每一步的输入输出。
2.3 技术栈明细
下面是一份推荐的技术栈版本示意,具体版本需要结合你的机器环境调整:
| 技术组件 | 版本建议 | 说明 |
|---|---|---|
| Ubuntu / CentOS | 20.04 / 7.x | 长期稳定版 |
| JDK | 1.8 或 11 | Hadoop 与 Spark 都需要 |
| Hadoop | 3.x | 相比 2.x 更适合学习 |
| Spark | 3.x | 建议带 PySpark |
| Python | 3.8+ | 配合 PySpark 与机器学习库 |
| CatBoost | 1.x | 通过 pip 安装 |
| pandas | 1.x/2.x | 配合数据处理 |
| scikit-learn | 1.x | 计算评估指标 |
| PyCharm / Jupyter | 任意 | 用于开发和调试 |
需要注意,版本不是越新越好。Hadoop 3.x 与 Spark 3.x 的组合在网络上资料最多,遇到问题容易查到。
3. 环境准备与版本说明
3.1 环境清单
在动手之前,先把环境分成两部分:分布式存储计算环境和 Python 机器学习环境。
分布式存储计算环境需要安装:
- JDK
- Hadoop(可先做伪分布式)
- Spark
Python 机器学习环境需要安装:
- Python 3.8+
- pyspark
- pandas
- catboost
- scikit-learn
- matplotlib
如果你暂时不具备搭建集群的条件,可以先让 Spark 跑在local[*]模式。Hadoop 也只做伪分布式配置,也就是在一个节点上模拟分布式环境。对毕业设计来说,伪分布式足够验证整个流程。
3.2 Hadoop 与 Spark 的安装细节
在 Linux 环境下,通常先把 JDK 安装好,然后配置JAVA_HOME。接着解压 Hadoop 和 Spark 到指定目录,并配置环境变量。
下面是/etc/profile或~/.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配置完成后,执行source ~/.bashrc让环境变量生效。
伪分布式模式下,至少需要修改 Hadoop 的core-site.xml、hdfs-site.xml和yarn-site.xml。下面给出一个常见的core-site.xml配置:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/usr/local/hadoop/tmp</value> </property> </configuration>hdfs-site.xml中建议把副本数设置为 1,避免伪分布式环境下报副本不足的警告:
<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> </configuration>格式化 NameNode 时使用hdfs namenode -format命令。启动 HDFS 后,使用jps命令看到NameNode、DataNode等进程说明启动成功。
Spark 安装相对简单,解压后如果 HADOOP 环境变量已经配置好,就能直接以 local 模式运行。启动pyspark验证环境即可。
很多初学者会遇到一个典型的启动报错:jar does not exist or is not a normal file: /usr/local/hadoop/share/hadoop/m...,这通常是因为 Hadoop 的share/hadoop目录路径不全,或SPARK_DIST_CLASSPATH环境变量没有配置正确。排查时先确认$HADOOP_HOME/share/hadoop下是否存在common、hdfs、mapreduce等子目录,再执行:
export SPARK_DIST_CLASSPATH=$(hadoop classpath)3.3 Python 机器学习依赖
Python 环境建议使用 conda 创建独立虚拟环境,避免与系统 Python 冲突。
conda create -n airquality python=3.9 conda activate airquality pip install pyspark pandas catboost scikit-learn matplotlib安装完成后,可以用 Python 交互环境验证 CatBoost 是否可用:
from catboost import CatBoostRegressor print("CatBoost ready")4. 数据获取与预处理
4.1 数据来源与字段说明
空气质量数据可以使用公开数据集模拟。项目采用河南省各地市的监测数据,核心字段包括:
| 字段名 | 含义 | 示例 |
|---|---|---|
| date | 日期 | 2023-01-01 |
| city | 城市 | 郑州 |
| AQI | 空气质量指数 | 98 |
| PM25 | PM2.5 浓度(μg/m³) | 72 |
| PM10 | PM10 浓度(μg/m³) | 105 |
| SO2 | 二氧化硫浓度(μg/m³) | 12 |
| NO2 | 二氧化氮浓度(μg/m³) | 45 |
| CO | 一氧化碳浓度(mg/m³) | 0.9 |
| O3 | 臭氧浓度(μg/m³) | 68 |
| grade | 空气质量等级 | 良 |
为了简化,可以准备一份henan_airquality.csv,按城市和日期存放逐日监测结果。字段中不要使用PM2.5这种带点号的列名,否则 Spark SQL 处理时需要用反引号包裹,略显麻烦,统一改为PM25更省事。
4.2 将数据上传到 HDFS
启动 HDFS 后,先创建项目目录,再把数据文件上传到 HDFS:
hdfs dfs -mkdir -p /airquality/input hdfs dfs -put /your_local_path/henan_airquality.csv /airquality/input/查看上传结果:
hdfs dfs -ls /airquality/input如果你还没有 HDFS,也可以在 Spark 中直接读取本地文件,代码中把路径改为file:///your_local_path/henan_airquality.csv即可。学习阶段先用本地文件打通流程,再切到 HDFS 更稳妥。
4.3 Spark 读取数据与基本清洗
使用 PySpark 读取 CSV 时,需要开启表头解析和类型推断:
# 文件路径:spark_preprocess.py from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_date spark = SparkSession.builder \ .appName("HenanAirQualityAnalysis") \ .master("local[*]") \ .getOrCreate() df = spark.read.csv( "hdfs://localhost:9000/airquality/input/", header=True, inferSchema=True ) df = df.withColumn("date", to_date(col("date"), "yyyy-MM-dd")) # 过滤掉关键字段为空的数据 df = df.filter( col("PM25").isNotNull() & col("city").isNotNull() & col("date").isNotNull() ) print("清洗后数据量:", df.count()) df.show(5)这里说明一下to_date的作用:CSV 中的 date 字段往往是字符串,如果不转换,后续按时间排序和求滞后特征时会出现逻辑错误。inferSchema=True虽然能自动推断字段类型,但日期类型不一定能自动识别,所以最好手动指定。
清洗后的数据还可以继续做去重:
df = df.dropDuplicates(["date", "city"])同一城市同一天如果存在重复监测记录,只保留一条,避免污染统计结果和预测标签。
5. Spark 核心数据分析
5.1 城市级空气质量排名
拿到清洗后的 DataFrame,可以先用 Spark SQL 计算各城市的平均 AQI 和平均 PM2.5:
from pyspark.sql.functions import avg, round city_stats = df.groupBy("city").agg( round(avg("AQI"), 2).alias("avg_AQI"), round(avg("PM25"), 2).alias("avg_PM25"), fast_count := None ) city_stats.orderBy(col("avg_PM25").desc()).show()在实际运行前,需要把上面的伪代码去掉,改成 Spark SQL 支持的写法。正确写法如下:
from pyspark.sql.functions import avg, round city_stats = df.groupBy("city").agg( round(avg("AQI"), 2).alias("avg_AQI"), round(avg("PM25"), 2).alias("avg_PM25"), round(avg("PM10"), 2).alias("avg_PM10") ) city_stats.orderBy(col("avg_PM25").desc()).show()这一段演示了 Spark 的聚合操作。结果可以写到 HDFS 输出目录:
city_stats.write.csv("/airquality/output/city_stats", header=True, mode="overwrite")5.2 时间维度统计
按月份统计全省平均 AQI,可以反映季节变化:
from pyspark.sql.functions import month, avg as _avg df_month = df.withColumn("month", month(col("date"))) month_stats = df_month.groupBy("month").agg( round(_avg("AQI"), 2).alias("avg_AQI"), round(_avg("PM25"), 2).alias("avg_PM25") ).orderBy("month") month_stats.show()这里使用month函数从日期列中提取月份。如果是跨年数据,最好再同时提取年份,例如withColumn("year", year(col("date"))),避免把不同年份的同一月份混在一起。
5.3 相关性分析
在做预测之前,先看一下哪些特征与 PM2.5 相关性更高。Spark 的 DataFrame 自带corr方法:
corr_value = df.stat.corr("PM25", "PM10") print("PM25 与 PM10 相关系数:", corr_value)这一步在毕业设计中很有展示价值,它可以作为后续选择特征列的依据之一。你也可以把多个相关性结果放到一张表格中写进论文。
6. 特征工程与 CatBoost 建模
6.1 特征工程思路
预测目标可以定义为:利用前一天的空气质量数据和滞后特征,预测当天的 PM2.5 浓度。
在实际数据处理时,我们常用“滞后特征”来表示历史信息。比如预测 1 月 3 日的 PM2.5,就使用 1 月 2 日、1 月 1 日的 PM2.5、AQI、SO2、NO2 等作为特征。
在 Spark 中,可以使用lag窗口函数构造滞后特征:
from pyspark.sql.window import Window from pyspark.sql.functions import lag window_spec = Window.partitionBy("city").orderBy("date") df_feat = df.withColumn("PM25_lag1", lag("PM25", 1).over(window_spec)) \ .withColumn("PM25_lag2", lag("PM25", 2).over(window_spec)) \ .withColumn("AQI_lag1", lag("AQI", 1).over(window_spec)) \ .withColumn("SO2_lag1", lag("SO2", 1).over(window_spec)) \ .withColumn("NO2_lag1", lag("NO2", 1).over(window_spec))这里需要注意的是,每个城市的第一个日期没有前一天的记录,所以lag会产生空值。后续需要对这些空值进行处理,最简单的办法是删除:
df_feat = df_feat.dropna()对于时间序列预测,更严谨的做法是使用训练集的均值填充空值,避免因为删除导致时间不连续。毕业设计阶段,如果日期数据足够长,直接删除前两行影响也不大,但要在论文里说明处理逻辑。
6.2 构造预测目标列
定义特征列feature_cols和预测目标列target_col。把 Spark DataFrame 转换为 Pandas DataFrame,方便 CatBoost 使用:
# 文件路径:prepare_train_data.py feature_cols = [ "AQI", "PM25", "PM10", "SO2", "NO2", "CO", "O3", "PM25_lag1", "PM25_lag2", "AQI_lag1", "SO2_lag1", "NO2_lag1" ] # 选择需要的列并转为 Pandas pandas_df = df_feat.select(feature_cols + ["date", "city"]).toPandas() # 目标列:当日 PM2.5 X = pandas_df[feature_cols] y = pandas_df["PM25"] print("样本数量:", len(X)) print("特征数量:", len(X.columns))这里有一个容易踩坑的地方:toPandas()会把所有数据拉到 Driver 节点内存,如果数据量超级大,会直接内存溢出(OOM)。因此建议先对数据做采样或先跑通本地小样本,再决定是否全量转换。也可以使用spark.sql("SET spark.sql.adaptive.enabled=true")优化执行计划。
6.3 训练 CatBoost 模型
划分训练集和测试集,然后构建 CatBoost 回归模型。为了体现模型调优能力,这里添加了评估集:
# 文件路径:train_catboost.py from catboost import CatBoostRegressor, Pool from sklearn.model_selection import train_test_split from sklearn.metrics import mean_absolute_error, mean_squared_error, r2_score import numpy as np X_train, X_test, y_train, y_test = train_test_split( X, y, test_size=0.2, random_state=42, shuffle=False ) model = CatBoostRegressor( iterations=1000, learning_rate=0.05, depth=6, loss_function='RMSE', eval_metric='RMSE', random_seed=42, od_type='Iter', od_wait=100, verbose=100 ) model.fit( X_train, y_train, eval_set=(X_test, y_test), use_best_model=True, plot=False ) # 模型评估 y_pred = model.predict(X_test) mae = mean_absolute_error(y_test, y_pred) rmse = np.sqrt(mean_squared_error(y_test, y_pred)) r2 = r2_score(y_test, y_pred) print("MAE:", mae) print("RMSE:", rmse) print("R2:", r2) # 保存模型 model.save_model("catboost_pm25_model.cbm")代码中shuffle=False很重要,因为时间序列数据不能像普通分类数据一样随机打乱,否则会造成数据泄漏,训练集包含未来的信息,测试结果会虚高。
这里再解释一个 CatBoost 特点:use_best_model=True表示在验证集上评估指标不再提升时,自动回退到历史上最优模型。配合od_type='Iter'和od_wait=100,可以提前停止训练,既节省时间,又能防止过拟合。
7. 完整实验流程与结果评估
7.1 流程串联
在实际运行项目时,建议把脚本拆成三个文件:
spark_preprocess.py:读取数据、清洗、滞后特征生成、导出特征数据。train_catboost.py:读取特征数据、训练模型、输出评估指标。predict_demo.py:加载模型,对新的特征数据做预测。
这种方式有利于后期写论文时解释每一个模块,也让代码更易维护。
7.2 模型评估指标说明
在回归预测任务中,最常看的三个指标:
| 指标 | 含义 | 越低/越高越好 |
|---|---|---|
| MAE | 平均绝对误差 | 越低越好 |
| RMSE | 均方根误差 | 越低越好,对大误差更敏感 |
| R2 | 决定系数 | 越接近 1 越好 |
如果发现 R2 很低,不要急着调模型参数,先检查特征列是否存在大量缺失,或者滞后特征数量是否足够。很多时候特征工程对 R2 的影响远大于模型调参。
7.3 示例运行结果
实际结果会随着数据集大小和特征丰富度变化,这里不贴具体跑分。你可以按以下逻辑整理输出:
样本数量: 12000 训练集: 9600 测试集: 2400 MAE: 12.xx RMSE: 16.xx R2: 0.8x答辩时如果能把预测值和真实值的折线图展示出来,效果会更好。可以简单使用 matplotlib 绘制:
import matplotlib.pyplot as plt plt.figure(figsize=(12, 5)) plt.plot(y_test.values[:200], label="真实值") plt.plot(y_pred[:200], label="预测值") plt.legend() plt.title("PM2.5 预测结果对比") plt.show()8. 常见问题与排查思路
在做这个项目的过程中,比较容易踩到下面这些问题,这里整理成一个排查表格:
| 问题现象 | 常见原因 | 解决思路 |
|---|---|---|
| SparkSession 启动失败 | JAVA_HOME 没有配置,或 JDK 版本过高 | 确认java -version正常,并 export JAVA_HOME |
| 读取 HDFS 文件报错 FileNotFound | 路径写错或 HDFS 尚未启动 | 先用hdfs dfs -ls /airquality/input检查路径 |
jar does not exist or is not a normal file | Hadoop classpath 异常或SPARK_DIST_CLASSPATH未设置 | 执行export SPARK_DIST_CLASSPATH=$(hadoop classpath) |
toPandas()导致 OOM | 数据量太大,Driver 节点内存不足 | 增加 Spark Driver 内存,或先采样小数据验证 |
| CatBoost 训练非常慢 | 迭代次数多、数据量大、CPU 负载高 | 降低iterations,使用od_type='Iter'提前停止 |
| 预测结果 R2 很低 | 特征缺失、滞后特征不足、数据未按时间排序 | 检查缺失值和shuffle参数 |
| HDFS DataNode 启动后自动退出 | 伪分布式配置副本数设太大 | 设置dfs.replication=1并格式化集群 |
其中最后一个问题非常常见。很多人启动 HDFS 后,DataNode 进程反复退出,在日志中看到java.io.IOException: Incompatible clusterIDs,这通常是因为多次格式化 NameNode 导致集群 ID 不一致。解决方法是删除 HDFS 临时目录和数据目录后重新格式化:
stop-dfs.sh rm -rf /usr/local/hadoop/tmp rm -rf /usr/local/hadoop/dfs/data hdfs namenode -format start-dfs.sh注意:这个操作会清空 HDFS 上的数据,所以只在测试环境执行,并且提前备份原始文件。
9. 最佳实践与工程建议
9.1 数据合法性与安全边界
本项目使用的空气质量数据应来自公开平台,仅用于学习与学术研究。在毕业设计报告中,建议写明数据来源和版权说明。不要使用未授权爬取的内部数据,也不要为了效果故意编造敏感字段。
9.2 先小数据跑通再上集群
很多同学一上手就搭三节点集群,结果环境问题占用了 70% 的时间。更合理的路线是:
- 在本地使用 Spark local 模式跑通全流程。
- 将数据量缩减到几百条,确认代码逻辑正确。
- 再扩展到大文件,测试 HDFS 读写。
- 最后才是多节点集群。
9.3 重视特征工程
写论文时,很多人会把重点放在调参上。实际上对于空气质量预测问题,滞后特征和气象特征是决定效果的核心。如果时间允许,可以再引入温度、湿度、风速等气象数据,利用 Spark 做多表 join,让模型效果明显提升。
9.4 模型持久化与展示
训练好的 CatBoost 模型可以保存为.cbm文件,也可以导出为.onnx。毕业设计答辩时,如果需要做在线演示,可以用简单 Flask 接口加载模型并返回预测结果,不需要写太复杂的前端页面。
9.5 日志与结果管理
建议为 Spark 任务开启日志记录,方便回溯。在spark-submit时加上--driver-memory 2g --executor-memory 2g等参数,避免运行到一半内存不足。中间结果也建议按日期分目录保存,不要全部堆在同一个输出路径。
10. 总结与后续扩展方向
到这里,一个完整的“基于 Spark 的河南省空气质量数据分析与预测系统”就搭建完成了。整个过程涵盖了 Hadoop HDFS 文件存储、Spark DataFrame 清洗聚合、窗口函数特征工程、CatBoost 回归预测和常见集群问题排查。毕业设计做到这一步,已经具备相当完整的技术链条。
如果想继续扩展,可以考虑三个方向:
第一,引入天气数据,把气象预报因子作为特征,提升预测精度;第二,使用 Spark MLlib 的模型作为对比基线,在论文中形成“机器学习基线 + CatBoost 精度”的对比实验;第三,用可视化大屏展示各地市空气质量排名和未来 24 小时趋势预测。
最后提醒三点:不要一上来就追求集群规模,先把单机流程跑通;特征工程优先级高于模型调参;所有实验过程保留截图和日志,方便写论文。希望这篇文章能帮你把毕业设计的核心流程顺利走通,也祝你在答辩时能把这套系统的价值清晰地讲出来。