
这次我们来看一个典型的计算机毕业设计选题基于 Spark 的河南省空气质量数据分析与预测系统。它把 Hadoop、Spark、Hive 这些大数据组件串起来再用 CatBoost 做 AQI 预测是一个覆盖“数据采集 → 清洗 → 分析 → 建模 → 可视化”全流程的大数据项目。如果你正在准备大数据方向的毕业设计或者想找一个能同时展示分布式计算和机器学习能力的课题这个方向值得重点考虑。这个课题最核心的看点是数据量够大技术栈够全业务逻辑清晰。空气质量监测数据天然带有时间序列特征和空间分布特征适合用 Spark 做分布式清洗和特征统计也适合用 CatBoost 这类梯度提升模型做数值预测。整个系统从底层存储到上层应用都有明确的分工不是那种看起来高大上、实际上只有几个接口的演示型项目。本文会围绕这套系统拆解它的技术栈定位、环境搭建、数据预处理、特征工程、CatBoost 建模、Spark 任务调度、可视化展示和常见问题排查。无论你是打算自己从零开发还是准备找团队定制这篇文章都可以作为技术选型和方案设计的参考。1. 核心能力速览先把这套系统的整体规格做一个汇总方便快速判断它是否符合你的毕设需求。能力项说明项目类型大数据分析 机器学习预测系统适合计算机科学与技术、大数据、数据科学方向毕业设计核心技术栈HadoopHDFS、Spark、Hive、MySQL、CatBoost、Python、可视化框架数据来源河南省各地市空气质量监测站点公开数据包含 AQI、PM2.5、PM10、SO2、NO2、CO、O3 等指标主要功能空气质量数据清洗、多维度统计分析、时序趋势分析、城市间对比、AQI 预测预测算法CatBoost支持特征重要性分析和模型可解释性支持平台Linux 服务器或本地虚拟机配合 PyCharm 或 Jupyter 开发部署方式Hadoop 集群 / 伪分布式 Spark 任务提交 可视化 Web 服务是否支持批量任务支持可通过 Spark 批量处理多年份、多城市监测数据数据存储HDFS 作为原始数据存储Hive 做数据仓库建模MySQL 存放统计结果和预测结果适合读者大数据方向毕业生、需要快速搭建完整数据项目的开发者、准备答辩展示系统的学生从材料看这个课题的组合方式是Hadoop 生态做数据底座Spark 做分布式计算CatBoost 做预测模型相比只做 ETL 或者只跑单个模型的课题它能覆盖的面试考点和答辩知识点更多。实际部署时集群规模可以根据机器配置调整不需要一上来就搭多节点集群先用伪分布式跑通流程是完全可行的。2. 适用场景与使用边界2.1 这个课题适合谁如果你符合以下任何一种情况这个选题方向是比较合适的需要完成大数据方向毕业设计但不想做纯理论分析希望有一个能运行、能演示、能讲清楚技术细节的系统。学习过 Hadoop 和 Spark 的基础知识但缺少一个完整的业务场景来串联这些组件。想在简历上补充一个包含分布式存储、分布式计算、机器学习的项目经历。需要演示数据量大了怎么办的场景空气质量监测数据的累积量可以很好地体现 HDFS 和 Spark 的处理优势。2.2 能解决什么问题从业务角度来说这套系统把原始监测数据变成了可查询、可分析、可预测的信息对河南省各地市的空气质量状况进行统一汇总按城市、按月份、按季度统计 AQI 均值、优良天数比例、首要污染物分布。分析不同污染物之间的相关性比如 PM2.5 和 PM10 的关系、CO 与 O3 在不同季节的变化规律。基于历史数据训练 CatBoost 模型预测未来 24 小时或未来几天的 AQI 数值和空气质量等级。将分析结果和预测结果通过可视化的方式呈现支撑环境变化趋势分析。2.3 不适合什么场景这套系统的定位是教学和毕业设计不建议用于真实的生产环境影响评估。原因是监测数据的完整性和实时性要求很高生产环境需要接入实时数据流和更严格的模型验证流程。此外如果只是单纯做论文理论研究不需要开发完整系统这个课题的工程量可能会显得偏重。最后如果完全没有大数据基础需要先补齐 Hadoop 和 Spark 的基本概念否则开发过程中遇到的报错会比较难定位。2.4 合规与边界提醒空气质量数据本身属于公开环境数据但需要注意几点数据采集必须来自官方公开渠道或者已获授权的数据源不能爬取非公开接口涉及具体城市、具体站点的数据在论文和演示中应做必要的脱敏处理预测结果只能作为学术研究和趋势参考不能用于正式的环境评估报告也不能替代官方空气质量发布信息。3. 系统总体设计3.1 分层架构这套系统的整体架构可以拆成五层每一层的职责都比较清晰层级组件职责数据采集层Python 爬虫 / 数据导入脚本从公开渠道抓取或下载河南省各地市空气质量历史数据存储层HDFS、Hive、MySQLHDFS 存原始 CSV/JSON 文件Hive 建外部表做数据仓库MySQL 存统计结果和预测结果计算层Spark Core、Spark SQL做数据清洗、聚合统计、特征计算生成模型训练所需的宽表算法层CatBoost、Scikit-learn训练 AQI 预测模型评估特征重要性输出模型文件应用层Flask / FastAPI、ECharts提供查询接口和可视化页面展示分析结果和预测曲线3.2 数据流向数据从采集到最终展示的完整流程大概是这样的原始监测数据以 CSV 或 JSON 格式上传到 HDFS。Spark 读取原始数据完成缺失值处理、异常值过滤、时间字段标准化。清洗后的数据写入 Hive 表按天分区存储。Spark SQL 在 Hive 表上执行多维聚合统计结果写入 MySQL。从 Hive 中抽取一段时间窗口的特征构造训练集和测试集。CatBoost 模型完成训练和评估模型序列化保存。Flask 服务读取模型和 MySQL 中的统计结果通过 ECharts 渲染前端页面。这个流程覆盖了一个大数据项目从数据到决策的完整链路。答辩的时候可以按这个数据流逐层讲解系统设计思路清晰也能展示你理解整个数据管道。4. 环境准备与前置条件4.1 硬件要求这个课题对硬件的要求不算苛刻可以做两种选择。第一种是伪分布式部署单台机器模拟 Hadoop 集群。推荐配置是 8 核 CPU、16GB 内存以上、200GB 可用磁盘。伪分布式模式适合开发和调试大部分功能都能正常跑通。第二种是集群部署使用 3 台或更多机器。每台机器推荐 4 核 CPU、8GB 内存以上其中 NameNode 所在节点内存建议至少 16GB。如果云服务器磁盘和内存有限可以先用伪分布式完成开发答辩前再在服务器上完整部署。4.2 软件版本清单以下是通用的版本选型参考具体版本以实际环境为准软件版本建议用途Ubuntu / CentOSUbuntu 18.04 及以上 / CentOS 7.x操作系统JDKJDK 8 或 JDK 11Hadoop、Spark 运行依赖HadoopHadoop 3.xHDFS 分布式存储SparkSpark 3.x分布式计算HiveHive 3.x数据仓库MySQLMySQL 5.7 / 8.0结果存储PythonPython 3.8 / 3.9数据分析、模型训练CatBoostcatboost 1.0梯度提升模型PySpark与 Spark 版本匹配Python 调用 Spark4.3 安装前检查在开始安装之前先用下面这组命令确认系统基础环境# 检查 Java 版本 java -version # 检查 Python 版本 python3 --version # 检查 SSH 是否可用Hadoop 伪分布式需要 SSH localhost 免密 ssh localhost # 检查系统内存 free -h这里有几个容易踩坑的地方Java 版本不兼容会导致 Hadoop 启动失败建议优先使用 Java 8SSH 如果没有配置免密Hadoop 启动时会要求反复输入密码磁盘空间建议预留 100GB 以上因为 HDFS 默认会保存多份副本数据量会翻倍。5. Hadoop、Spark 与 Hive 部署步骤5.1 Hadoop 伪分布式部署Hadoop 的安装配置是很多同学遇到的第一个门槛。这里给出一套伪分布式部署的通用流程核心配置在core-site.xml、hdfs-site.xml和yarn-site.xml三个文件中。先解压 Hadoop 并配置环境变量# 解压 Hadoop这里以 hadoop-3.3.x 为例 tar -zxvf hadoop-3.3.4.tar.gz -C /usr/local/ cd /usr/local/hadoop-3.3.4 # 编辑 /etc/profile vim /etc/profile # 在文件末尾追加以下内容 export HADOOP_HOME/usr/local/hadoop-3.3.4 export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64编辑core-site.xmlconfiguration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/usr/local/hadoop-3.3.4/tmp/value /property /configuration编辑hdfs-site.xmlconfiguration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name valuefile:///usr/local/hadoop-3.3.4/namenode/value /property property namedfs.datanode.data.dir/name valuefile:///usr/local/hadoop-3.3.4/datanode/value /property /configuration配置完成后格式化 NameNode 并启动服务cd /usr/local/hadoop-3.3.4 # 格式化 NameNode只有第一次需要执行 bin/hdfs namenode -format # 启动 HDFS sbin/start-dfs.sh # 启动 YARN sbin/start-yarn.sh # 验证进程是否启动对应进程为 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager jps启动后浏览器访问http://localhost:9870可以看到 NameNode 的管理页面这也说明 HDFS 已经正常工作了。5.2 Spark 安装配置Spark 的安装相对简单只要与 Hadoop 版本兼容即可。下载 Spark 安装包并解压然后配置环境变量。# 解压 Spark tar -zxvf spark-3.3.2-bin-hadoop3.tgz -C /usr/local/ cd /usr/local/spark-3.3.2-bin-hadoop3 # 编辑 spark-env.sh cp conf/spark-env.sh.template conf/spark-env.sh vim conf/spark-env.sh # 追加以下内容 export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 export HADOOP_HOME/usr/local/hadoop-3.3.4 export SPARK_MASTER_HOSTlocalhost启动 Spark 独立集群模式cd /usr/local/spark-3.3.2-bin-hadoop3 sbin/start-master.sh sbin/start-worker.sh spark://localhost:7077启动后访问http://localhost:8080可以看到 Spark Master 的 Web 页面Worker 状态正常即可。5.3 Hive 安装与 MySQL 配置Hive 的作用是把 HDFS 上的结构化数据映射成数据表方便用 SQL 做统计分析。Hive 的元数据默认存在 derby 中建议改成 MySQL这样可以避免多人访问时锁表的问题。-- 在 MySQL 中创建 Hive 元数据库 CREATE DATABASE hive_metastore CHARACTER SET utf8mb4 COLLATE utf8mb4_bin; -- 创建 Hive 用户并授权 CREATE USER hive% IDENTIFIED BY hive_password; GRANT ALL PRIVILEGES ON hive_metastore.* TO hive%; FLUSH PRIVILEGES;配置hive-site.xml核心内容如下configuration property namejavax.jdo.option.ConnectionURL/name valuejdbc:mysql://localhost:3306/hive_metastore?useSSLfalse/value /property property namejavax.jdo.option.ConnectionDriverName/name valuecom.mysql.cj.jdbc.Driver/value /property property namejavax.jdo.option.ConnectionUserName/name valuehive/value /property property namejavax.jdo.option.ConnectionPassword/name valuehive_password/value /property /configuration这里需要注意的是Hive 3.x 需要额外的 Guava 版本适配问题。如果启动时提示 Guava 版本冲突需要把 Hadoop 目录下的 guava jar 包替换到 Hive 的 lib 目录下具体版本以实际报错为准。6. 数据采集与预处理6.1 数据采集方案空气质量数据可以从公开环境监测平台获取也可以使用公开数据集。采集脚本可以用 Python 编写将抓到的数据保存为 CSV 文件。为了演示这里给出一份模拟数据结构date,city,station,aqi,pm25,pm10,so2,no2,co,o3 2024-01-01,郑州,郑东新区监测站,85,62,90,8,38,0.9,65 2024-01-01,洛阳,涧西区监测站,92,68,105,10,42,1.1,58 2024-01-01,开封,龙亭区监测站,76,55,78,7,35,0.8,70采集脚本的通用模板如下import csv import random import datetime import os def generate_demo_data(start_date, end_date, cities, output_dir): 生成模拟空气质量数据格式与公开数据保持一致。 用于开发和演示真实数据请替换为授权数据源。 os.makedirs(output_dir, exist_okTrue) current datetime.datetime.strptime(start_date, %Y-%m-%d) end datetime.datetime.strptime(end_date, %Y-%m-%d) with open(os.path.join(output_dir, air_quality.csv), w, newline, encodingutf-8) as f: writer csv.writer(f) writer.writerow([date, city, station, aqi, pm25, pm10, so2, no2, co, o3]) while current end: for city in cities: # 随机生成各监测指标模拟数据仅用于开发测试 pm25 random.randint(20, 150) pm10 int(pm25 * random.uniform(1.2, 1.5)) so2 random.randint(5, 30) no2 random.randint(20, 60) co round(random.uniform(0.5, 1.8), 1) o3 random.randint(40, 90) station f{city}监测站 writer.writerow([ current.strftime(%Y-%m-%d), city, station, int(pm25 * 0.5 pm10 * 0.3 random.uniform(10, 30)), pm25, pm10, so2, no2, co, o3 ]) current datetime.timedelta(days1) print(f模拟数据已生成到 {output_dir}/air_quality.csv) if __name__ __main__: cities [郑州, 洛阳, 开封, 安阳, 南阳, 商丘, 新乡, 平顶山] generate_demo_data(2022-01-01, 2024-12-31, cities, ./data)实际开发中建议先确认数据源的字段格式然后调整清洗逻辑。模拟数据的好处是可以快速跑通全流程等系统稳定后再替换成真实数据。6.2 Spark 数据清洗数据清洗是 Spark 处理的第一步。读取 CSV 后主要完成以下操作去除日期、城市、监测站字段为空的行。过滤 AQI 数值明显异常的数据例如 AQI 小于 0 或者超过 500。统一时间格式将日期列转换为标准格式并提取年、月、日、季度等维度字段。对缺失的污染物指标进行填充通用做法是针对同一城市前一日的数值做线性插值。下面是一段 PySpark 清洗代码示例from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, year, month, quarter, round as spark_round spark SparkSession.builder \ .appName(AirQualityClean) \ .master(local[*]) \ .enableHiveSupport() \ .getOrCreate() # 读取 HDFS 上的原始数据 df spark.read.csv( hdfs://localhost:9000/user/airquality/raw/air_quality.csv, headerTrue, inferSchemaTrue ) # 清洗逻辑 df_cleaned df \ .dropDuplicates([date, city, station]) \ .na.drop(subset[date, city]) \ .filter(col(aqi).isNotNull()) \ .filter((col(aqi) 0) (col(aqi) 500)) \ .withColumn(year, year(col(date))) \ .withColumn(month, month(col(date))) \ .withColumn(quarter, quarter(col(date))) # 输出到 Hive 表 df_cleaned.write.mode(overwrite) \ .partitionBy(year, month) \ .format(parquet) \ .saveAsTable(airquality.tb_air_quality_clean) spark.stop() print(数据清洗完成)清洗后的结果按年和月分区存入 Hive后续统计分析只需要扫描对应分区效率会明显提高。7. 特征工程与统计分析7.1 统计分析维度在 Hive 表基础上可以用 Spark SQL 快速完成一些常规聚合分析。比如按城市统计年度 AQI 均值SELECT city, year, ROUND(AVG(aqi), 2) AS avg_aqi, ROUND(AVG(pm25), 2) AS avg_pm25, COUNT(*) AS record_count FROM airquality.tb_air_quality_clean GROUP BY city, year ORDER BY city, year;也可以分析污染天数占比SELECT city, SUM(CASE WHEN aqi 100 THEN 1 ELSE 0 END) AS prev_good_days, COUNT(*) AS total_days, ROUND(SUM(CASE WHEN aqi 100 THEN 1 ELSE 0 END) / COUNT(*) * 100, 2) AS good_ratio FROM airquality.tb_air_quality_clean WHERE year 2023 GROUP BY city ORDER BY good_ratio;这些统计结果可以写到 MySQL用于前端可视化展示。7.2 预测特征构造用 CatBoost 做 AQI 预测之前需要构造特征集。常用的特征思路如下时间特征年、月、日、星期几、是否周末、是否节假日、所在季度。历史滑动窗口特征前一日的 AQI、PM2.5、PM10 等前 3 日 AQI 均值前 7 日 AQI 均值。滞后污染物特征PM2.5 滞后 1 日、SO2 滞后 1 日等。气象相关特征如果数据中有温度、湿度、风速等字段也可以加入。特征宽表可以在 Spark 中构造输出为 Parquet 文件from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import lag, avg as spark_avg spark SparkSession.builder \ .appName(FeatureEngineer) \ .enableHiveSupport() \ .getOrCreate() df spark.sql(SELECT * FROM airquality.tb_air_quality_clean WHERE year 2022) windowSpec Window.partitionBy(city).orderBy(date) df_feature df \ .withColumn(aqi_lag1, lag(aqi, 1).over(windowSpec)) \ .withColumn(pm25_lag1, lag(pm25, 1).over(windowSpec)) \ .withColumn(pm10_lag1, lag(pm10, 1).over(windowSpec)) \ .withColumn(aqi_lag3_avg, spark_avg(aqi).over(windowSpec.rowsBetween(-3, -1))) # 过滤前几天的空值 df_feature df_feature.na.drop(subset[aqi_lag1, pm25_lag1, pm10_lag1]) df_feature.write.mode(overwrite) \ .format(parquet) \ .save(hdfs://localhost:9000/user/airquality/feature/feature_table)特征构造完成后统计特征数量和记录数量确认没有明显异常再进入模型训练。8. CatBoost 模型训练与评估8.1 模型选型依据CatBoost 是梯度提升决策树模型对表格数据效果稳定不需要做复杂的特征标准化支持类别特征也自带特征重要性输出。对于空气质量预测这种表格型回归任务是一个比较稳妥的选择。相比 XGBoost 和 LightGBMCatBoost 在默认参数下对类别特征的处理更友好训练过程也不容易过拟合。8.2 训练代码从 HDFS 读取特征表划分训练集和测试集然后训练 CatBoost 回归模型import pandas as pd from catboost import CatBoostRegressor, Pool from sklearn.model_selection import train_test_split from sklearn.metrics import mean_squared_error, mean_absolute_error, r2_score # 从 HDFS 或本地读取特征数据 # 示例从本地 Parquet 读取实际项目中改为 Spark 导出的文件路径 df pd.read_parquet(feature_table.parquet) feature_cols [ year, month, day, weekday, is_weekend, pm25, pm10, so2, no2, co, o3, aqi_lag1, pm25_lag1, pm10_lag1, aqi_lag3_avg ] target_col aqi X df[feature_cols] y df[target_col] X_train, X_test, y_train, y_test train_test_split( X, y, test_size0.2, random_state42, shuffleFalse ) # 指定类别特征例如 is_weekend 虽然是数值也可以作为类别处理 cat_features [is_weekend] model CatBoostRegressor( iterations1000, learning_rate0.05, depth6, loss_functionRMSE, eval_metricRMSE, random_seed42, verbose100 ) train_pool Pool(X_train, y_train, cat_featurescat_features) test_pool Pool(X_test, y_test, cat_featurescat_features) model.fit(train_pool, eval_settest_pool, use_best_modelTrue) # 评估 pred model.predict(X_test) mse mean_squared_error(y_test, pred) mae mean_absolute_error(y_test, pred) r2 r2_score(y_test, pred) print(fMSE: {mse:.2f}) print(fMAE: {mae:.2f}) print(fR2: {r2:.4f}) # 特征重要性 feature_importance model.get_feature_importance() for name, importance in zip(feature_cols, feature_importance): print(f{name}: {importance:.4f}) # 保存模型 model.save_model(catboost_aqi_model.cbm) print(模型已保存)8.3 模型评估要点判断模型效果时优先看 R2 和 MAE 两个指标。R2 越接近 1 说明模型对训练集和测试集解释力越强MAE 表示平均绝对误差。从经验上看AQI 预测任务 MAE 控制在 10 到 20 之间、R2 在 0.85 以上已经是不错的效果。但具体指标和数据集质量、特征丰富度直接相关不能作为绝对标准。实际开发中还要注意时间序列数据的特殊性不能随机打乱数据后划分训练集和测试集必须按时间先后切分否则会引入未来信息模型评估结果虚高。上面的代码中shuffleFalse就是按原始时间顺序做的切分。9. 可视化展示与 Web 服务9.1 技术方案可视化层可以采用 Flask 或 FastAPI 提供后端接口配合 ECharts 做前端展示。统计结果存放在 MySQL预测结果可以由模型实时计算后返回。推荐展示的页面包括首页总览全省 AQI 均值、优良天数比例、主要污染物分布。城市对比页柱状图展示各地市 AQI 排名。时间趋势页折线图展示某城市 AQI 随时间变化的趋势。特征相关性页热力图展示各污染物之间的相关系数。预测页展示未来 7 天 AQI 预测曲线以及置信区间。9.2 后端接口示例用 Flask 写一个简单的 AQI 预测接口from flask import Flask, request, jsonify from catboost import CatBoostRegressor import pandas as pd import datetime app Flask(__name__) model CatBoostRegressor() model.load_model(models/catboost_aqi_model.cbm) feature_cols [ year, month, day, weekday, is_weekend, pm25, pm10, so2, no2, co, o3, aqi_lag1, pm25_lag1, pm10_lag1, aqi_lag3_avg ] app.route(/predict, methods[POST]) def predict(): data request.get_json() # 根据前端传入的日期拼装特征 target_date datetime.datetime.strptime(data[date], %Y-%m-%d) features { year: target_date.year, month: target_date.month, day: target_date.day, weekday: target_date.weekday(), is_weekend: 1 if target_date.weekday() 5 else 0, pm25: data[pm25], pm10: data[pm10], so2: data[so2], no2: data[no2], co: data[co], o3: data[o3], aqi_lag1: data[aqi_lag1], pm25_lag1: data[pm25_lag1], pm10_lag1: data[pm10_lag1], aqi_lag3_avg: data[aqi_lag3_avg] } df pd.DataFrame([features]) pred float(model.predict(df)[0]) return jsonify({date: data[date], predict_aqi: round(pred, 2)}) if __name__ __main__: app.run(host0.0.0.0, port5000)启动服务后用 curl 可以快速验证接口curl -X POST http://127.0.0.1:5000/predict \ -H Content-Type: application/json \ -d { date: 2024-06-01, pm25: 65, pm10: 90, so2: 8, no2: 38, co: 0.9, o3: 70, aqi_lag1: 88, pm25_lag1: 60, pm10_lag1: 85, aqi_lag3_avg: 82 }接口返回 JSON 数据后前端可以直接用 ECharts 渲染预测曲线。10. 批量任务与调度设计10.1 Spark 批量统计批量任务是这套系统非常实用的一部分。空气质量数据是逐日累积的每天都会有新数据流入可以设计定时任务完成以下批量操作每天凌晨将前一天的监测数据上传到 HDFS。Spark 任务读取新增数据追加到 Hive 表。Spark SQL 更新 MySQL 中的统计结果表。保存最新的模型输入特征文件用于后续预测。用户如果希望一键运行整套流程可以写一个 Shell 脚本依次执行数据上传、清洗、统计分析、模型预测等步骤#!/bin/bash # 设置环境变量 export HADOOP_HOME/usr/local/hadoop-3.3.4 export SPARK_HOME/usr/local/spark-3.3.2-bin-hadoop3 export HIVE_HOME/usr/local/hive-3.1.3 # 1. 上传当天数据到 HDFS hadoop fs -put /data/airquality/20240601.csv /user/airquality/raw/ # 2. 提交 Spark 清洗任务 spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 2g \ --driver-memory 2g \ scripts/air_quality_clean.py # 3. 提交统计分析任务 spark-submit \ --master yarn \ --deploy-mode cluster \ scripts/air_quality_stats.py # 4. 触发模型训练或预测 python scripts/train_or_predict.py echo 批量任务执行完成10.2 任务监控Spark 任务提交后可以通过 YARN 的 Web 页面查看任务状态、日志和执行进度。如果任务卡住优先检查资源是否充足、数据路径是否正确、是否有数据倾斜的情况。批量任务建议在脚本中增加日志输出和失败重试逻辑避免因为一个文件格式问题导致整批任务中断。11. 资源占用与性能观察11.1 观察方法运行 Spark 任务时可以通过几个命令实时观察资源占用# 查看 HDFS 空间使用 hadoop dfs -df -h # 查看 YARN 任务列表和资源分配 yarn application -list yarn application -status application_xxxxx # 查看系统 CPU 和内存 top -u hadoop free -g在 Spark Master 的 Web 页面中可以看到每个 Executor 的核数和内存分配。如果发现某个 Executor 内存溢出需要调整spark.executor.memory参数或者减少单次处理的分区数。11.2 性能优化思路数据倾斜如果某个城市的数据量明显大于其他城市可以在groupBy之前加盐打散或者使用repartition重新分区。小文件过多Hive 表如果按天分区但每天数据量很小会产生大量小文件。建议使用spark.sql.shuffle.partitions控制分区数量并定期合并小文件。模型训练资源CatBoost 训练建议放在 Spark 任务之后单独执行不要和 Spark 同时抢占资源避免内存不足。结果写入 MySQL批量写入时使用批量插入避免逐条写入减少数据库压力。12. 常见问题与排查方法问题现象可能原因排查方式解决方案Hadoop 启动失败jps 看不到 NameNode未格式化 NameNode或 core-site.xml 配置错误查看 Hadoop 日志检查 tmp 目录停掉所有进程删除 tmp 目录重新格式化Spark 任务提交后一直等待YARN 资源不足或 Worker 没有启动查看 YARN 页面检查 Worker 状态关闭其他占用内存的任务增加 Executor 内存Hive 启动报 Guava 版本冲突Hadoop 与 Hive 自带 Guava 版本不一致查看日志中的版本信息将 Hadoop lib 中的 guava jar 替换到 Hive libCatalyst 内存溢出 OOM单次处理数据量过大或分区数过少查看 Executor 日志增加分区数调大 executor-memoryMySQL 写入中文乱码表字符集不是 utf8mb4检查表结构和 JDBC URL建表时指定 utf8mb4URL 加characterEncodingutf8CatBoost 训练结果偏差大特征泄漏或数据切分方式错误检查是否使用未来数据将数据按时间排序切分禁止随机切分接口请求超时模型文件过大或服务器配置低打印接口耗时日志减小模型规模或使用批量预测接口前端图表不显示后端返回字段格式不一致使用 curl 测试接口统一 JSON 字段命名检查 ECharts 数据格式13. 最佳实践与使用建议13.1 开发顺序建议不要试图一次性把全链路做完。建议先按以下顺序推进先搭 Hadoop 伪分布式环境确保 HDFS 能正常读写。用少量模拟数据跑通 Spark 清洗流程。建立 Hive 表完成最简单的 SQL 统计。做特征表训练一个基础版 CatBoost 模型。封装 Flask 接口做可视化页面。最后扩展批量调度和性能优化。13.2 工程化管理数据文件、脚本、模型、日志分目录存放避免所有文件堆在同一个目录下。每一个处理步骤都输出日志记录输入数据量、处理耗时、异常记录数。模型文件版本化训练完成后保存到独立目录同时记录训练参数和评估指标。接口服务上线前用测试数据跑一遍完整流程确认输入输出格式稳定。13.3 答辩注意事项答辩时讲解重点放在三块分布式存储解决什么问题、Spark 如何处理海量监测数据、CatBoost 模型如何预测 AQI。建议准备好一份数据流图每层使用的组件、数据格式、接口方式都标清楚。演示的时候优先展示批量任务的日志输出和可视化图表让评审看到系统可以真实运行不只是一个概念框架。14. 总结与下一步这套基于 Spark 的河南省空气质量数据分析与预测系统最大的价值是串联了完整的大数据技术链路。从 HDFS 存储、Hive 数仓、Spark SQL 统计到 CatBoost 预测、Flask 接口和 ECharts 可视化每一层都能找到对应的工作模块适合作为毕业设计深入了解大数据工程的全貌。最先建议验证的功能是数据清洗和基础统计先把数据链路跑通确认 HDFS 和 Spark 环境稳定再进入模型训练阶段。最容易踩的坑是环境配置尤其是 Hadoop 和 Hive 的版本兼容问题建议在正式开发前先安装一套最小可运行环境。后续如果想继续扩展可以接入实时数据流处理把系统升级为近实时监测平台也可以引入更多气象数据和卫星遥感数据提高模型预测精度还可以把预测目标从城市级细化到站点级分析单个监测站点的污染变化规律。对于毕业设计来说先把基础链路做到稳定再在这个基础上做一两个亮点功能整个系统的完成度会高很多。