尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

基于Spark的河南省空气质量数据分析与预测系统实战

基于Spark的河南省空气质量数据分析与预测系统实战 每年到毕业设计选题季很多同学都会纠结一个问题想做大数据方向又担心纯理论不好做想带上机器学习预测又怕自己基础不够。如果你也处于这个阶段可以考虑“基于 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 等模块。MapReduceHadoop 中的批处理计算模型编程复杂度较高适合离线清洗但写起来偏繁琐。CatBoostYandex 开源的梯度提升决策树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 / CentOS20.04 / 7.x长期稳定版JDK1.8 或 11Hadoop 与 Spark 都需要Hadoop3.x相比 2.x 更适合学习Spark3.x建议带 PySparkPython3.8配合 PySpark 与机器学习库CatBoost1.x通过 pip 安装pandas1.x/2.x配合数据处理scikit-learn1.x计算评估指标PyCharm / Jupyter任意用于开发和调试需要注意版本不是越新越好。Hadoop 3.x 与 Spark 3.x 的组合在网络上资料最多遇到问题容易查到。3. 环境准备与版本说明3.1 环境清单在动手之前先把环境分成两部分分布式存储计算环境和 Python 机器学习环境。分布式存储计算环境需要安装JDKHadoop可先做伪分布式SparkPython 机器学习环境需要安装Python 3.8pysparkpandascatboostscikit-learnmatplotlib如果你暂时不具备搭建集群的条件可以先让 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 namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/usr/local/hadoop/tmp/value /property /configurationhdfs-site.xml中建议把副本数设置为 1避免伪分布式环境下报副本不足的警告configuration property namedfs.replication/name value1/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 python3.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-01city城市郑州AQI空气质量指数98PM25PM2.5 浓度μg/m³72PM10PM10 浓度μg/m³105SO2二氧化硫浓度μg/m³12NO2二氧化氮浓度μg/m³45CO一氧化碳浓度mg/m³0.9O3臭氧浓度μg/m³68grade空气质量等级良为了简化可以准备一份henan_airquality.csv按城市和日期存放逐日监测结果。字段中不要使用PM2.5这种带点号的列名否则 Spark SQL 处理时需要用反引号包裹略显麻烦统一改为PM25更省事。4.2 将数据上传到 HDFS启动 HDFS 后先创建项目目录再把数据文件上传到 HDFShdfs 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/, headerTrue, inferSchemaTrue ) 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 字段往往是字符串如果不转换后续按时间排序和求滞后特征时会出现逻辑错误。inferSchemaTrue虽然能自动推断字段类型但日期类型不一定能自动识别所以最好手动指定。清洗后的数据还可以继续做去重df df.dropDuplicates([date, city])同一城市同一天如果存在重复监测记录只保留一条避免污染统计结果和预测标签。5. Spark 核心数据分析5.1 城市级空气质量排名拿到清洗后的 DataFrame可以先用 Spark SQL 计算各城市的平均 AQI 和平均 PM2.5from 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, headerTrue, modeoverwrite)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.enabledtrue)优化执行计划。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_size0.2, random_state42, shuffleFalse ) model CatBoostRegressor( iterations1000, learning_rate0.05, depth6, loss_functionRMSE, eval_metricRMSE, random_seed42, od_typeIter, od_wait100, verbose100 ) model.fit( X_train, y_train, eval_set(X_test, y_test), use_best_modelTrue, plotFalse ) # 模型评估 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)代码中shuffleFalse很重要因为时间序列数据不能像普通分类数据一样随机打乱否则会造成数据泄漏训练集包含未来的信息测试结果会虚高。这里再解释一个 CatBoost 特点use_best_modelTrue表示在验证集上评估指标不再提升时自动回退到历史上最优模型。配合od_typeIter和od_wait100可以提前停止训练既节省时间又能防止过拟合。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 fileHadoop classpath 异常或SPARK_DIST_CLASSPATH未设置执行export SPARK_DIST_CLASSPATH$(hadoop classpath)toPandas()导致 OOM数据量太大Driver 节点内存不足增加 Spark Driver 内存或先采样小数据验证CatBoost 训练非常慢迭代次数多、数据量大、CPU 负载高降低iterations使用od_typeIter提前停止预测结果 R2 很低特征缺失、滞后特征不足、数据未按时间排序检查缺失值和shuffle参数HDFS DataNode 启动后自动退出伪分布式配置副本数设太大设置dfs.replication1并格式化集群其中最后一个问题非常常见。很多人启动 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 小时趋势预测。最后提醒三点不要一上来就追求集群规模先把单机流程跑通特征工程优先级高于模型调参所有实验过程保留截图和日志方便写论文。希望这篇文章能帮你把毕业设计的核心流程顺利走通也祝你在答辩时能把这套系统的价值清晰地讲出来。
返回列表