拓冰建站拓冰建站
首页 / 资讯中心 / 正文

基于Spark与CatBoost的河南省空气质量数据分析与预测

上周有个读者给我看了他的毕业设计选题基于Spark的河南省空气质量数据分析与预测系统技术栈写着 Hadoop、Spark、CatBoost。他原话是“模型我都调过一轮了CatBoost 在 Jupyter 里跑得挺准但放到 Spark 上就各种报错。”我问了三个问题数据从哪来数据放在 HDFS 还是本地你是用 Spark 训练还是只用 Spark 做预处理他沉默了。这不是他一个人的问题很多选择大数据分析类选题的人第一反应都是“先调模型”但这类项目真正的隐藏考点从来不是准确率而是能不能把一条从原始数据到最终预测结果的完整链路跑通。如果你也选了类似题目或者正在犹豫要不要选我的建议是先别急着定算法先把数据处理链路画出来。链路通模型准不准反而只是调参问题链路不通模型再厉害都只是 PPT 上的演示。1. 这个选题真正考的不是模型而是数据链路1.1 为什么很多类似项目最后都卡在数据接入先说一个很容易被忽略的事实河南省空气质量数据分析与预测听起来是一个机器学习任务但第一步并不是模型训练而是把数据拿到手并且搞成干净的、可计算的表格。很多学生选题时会把重点放在“我要用 CatBoost 预测 PM2.5”然后花大量时间调参等到中期检查才发现数据源还没搞定或者数据拿到了但 HDFS 里根本没法直接分析。这种场景在毕业设计里非常常见。空气质量的原始数据通常不是一份完美的 CSV。公开数据平台可能按小时、按站点返回数据字段包括时间、城市、站点名称、AQI、PM2.5、PM10、SO2、NO2、O3、CO但不同站点的缺失情况差别很大。有的站点凌晨数据经常断有的站点因为设备维护连续缺一周。更麻烦的是不同平台的字段命名、时间格式、数值精度都不同这些数据要进入 Spark 分析必须统一成一套规范的 Schema。所以真正卡住进度的往往不是算法而是数据怎么定期获取或导入缺失值怎么处理时间字段怎么统一站点名称怎么编码不同数据源怎么对齐。这些工作占整个项目的工作量通常超过一半。如果你只把它当成“模型训练前的数据清洗”就会一直觉得烦但如果你把它理解成“数据链路设计”就会知道这是整个项目的地基。1.2 河南省空气质量数据的典型形态以河南省为例空气质量数据一般可以按照“站点-时间-污染物浓度”来组织。每个站点每小时会有一条记录字段大概包含站点编码、城市、测点名称、监测时间、PM2.5、PM10、SO2、NO2、O3、CO、AQI 等。这里有几个实际处理时要注意的点AQI 不是直接读出来就能用的字段很多平台会给出 IAQI 或者污染物浓度AQI 需要按国家标准计算而计算逻辑在不同污染物的分级限值下并不一样。时间粒度一般是小时级但有些平台可能提供日平均数据。如果你的预测目标是对未来 24 小时的 PM2.5 做预测小时级数据是基础但如果只想分析季节趋势日平均也够用。站点编码是类别特征不是数值特征。建模时不能直接把“站点”当整数输入最好用 CatBoost 内置的类别特征处理或者做有序编码。气象数据比如温度、湿度、风速、风向、气压通常需要从气象数据平台单独获取。和空气质量数据合并时要按站点和小时做关联而两套数据的坐标系、时间标准、站点位置不一定完全一致。这些看起来都是细节但每一条都会决定你的特征工程能不能做下去。如果输入层就乱了后面的 Spark 算得再快也没有意义。1.3 判断要不要用 Spark数据规模才是第一标准很多人在选题时会把 Spark 当作一个“必须用的技术”理由是毕业设计需要体现大数据分析能力。这个动机可以理解但 Spark 并不只代表“先进”它更代表“成本和复杂度”。你需要先做一个简单的判断如果只有一个站点的几年小时级数据行数大概在几十万量级用 Pandas 处理完全足够如果有全省几十个站点的多年小时级数据行数可能达到几千万甚至上亿单机内存会非常紧张用 Spark 做分布式处理才合理如果你还要做多表关联、时间窗口聚合、滞后特征批量生成单机脚本的运行时间会非常不可控Spark 的计算模型更适合这种重活。所以我不是无条件推荐上 Spark。如果你的数据量不够大强行引入 Hadoop、Spark 只会让你陷入环境调试的泥潭。但如果你的目标是做一个能体现大数据处理能力的系统并且数据确实到了千万行以上那么 Hadoop Spark 就是合理选择。还有一个判断标准是“将来会不会重复跑”。毕业设计可能只需要跑一次结果但如果这个系统要能接收新数据、定期更新预测那么 Spark 任务的重复执行能力、日志输出和异常处理都会比单机脚本更合适。2. 技术选型Hadoop、Spark、CatBoost 各承担什么角色2.1 一个清晰的职责划分很多初学同学会把 Hadoop、Spark、CatBoost 混在一起以为它们是同一个层面的东西。实际上它们解决的是不同阶段的问题。可以用一条流水线来理解Hadoop HDFS 负责存储原始数据和中间结果。空气质量数据是典型的文本格式文件非常适合放在 HDFS 里统一管理。Hadoop YARN 负责资源调度但初学者可以先不深入 YARN 细节用 Spark 的 local 模式跑通流程后再提交到集群。Spark 负责分布式数据处理。读文件、清洗、过滤、时间窗口聚合、特征构造这些都可以用 Spark DataFrame / Spark SQL 完成。CatBoost 负责模型训练和预测。它和 Spark 的关系没有那么大既可以直接读 Spark 导出的特征文件也可以在 Spark 任务里调用 Python 训练脚本。这个职责划分很重要。不要在“用 Spark 训练 CatBoost”这个问题上纠结太久。CatBoost 本质上是一个单机机器学习库虽然支持一定的多线程训练但它不是为分布式训练设计的。常见的做法是Spark 负责把数据变成特征宽表然后导出成 Parquet 或 CSV再交给 CatBoost 训练。2.2 CatBoost 的优势与边界在空气质量预测场景里CatBoost 有几个很实在的优势支持类别特征。站点编号、城市、星期、小时都可以直接作为类别特征传入不用做繁琐的 One-Hot 编码。对缺失值有内置处理逻辑。空气质量数据经常有缺失CatBoost 在训练时会根据特征分布处理缺失值减少手工插值带来的偏差。在中小规模数据上训练速度通常可控。相比深度神经网络表格数据用梯度提升树更容易获得稳定效果。但要注意边界。CatBoost 也不是万能的。它适合处理“特征维度不是特别高、数据之间有大量表格型关系”的问题。如果你的目标是把预测结果延展到空间分布加入经纬度和地理距离或者要捕捉很复杂的时间序列依赖那么你可能还需要 LSTM、Prophet 或者图神经网络这些就不是 CatBoost 能简单覆盖的了。对毕业设计而言我更建议先做 CatBoost把基线模型跑通。如果时间充裕再做一个 LSTM 对比突出“不同模型在同一任务上的差异”。这样既体现了分析深度也不会因为一个模型调不出来而卡住整体进度。2.3 组件取舍能省则省但 HDFS 不能省Hadoop 生态非常大Hive、HBase、Kafka、Flume、Redis 各有各的用途但并不是所有组件都在同一个项目里必须出现。一个典型的毕业设计技术栈可以收敛为Hadoop HDFS存储原始文件和模型预测结果Spark完成数据处理、特征工程和部分统计分析CatBoost完成回归预测MySQL 或 SQLite存储最终结果方便可视化系统查询Superset / ECharts / Grafana做图表展示。至于很多人问的 Spark 读取 Redis通常用于缓存结果或实时查询。在毕业设计里如果老师没有特别要求实时性Redis 不是必须。你可以把它写进“系统优化方向”作为未来改进点而不是一开始就引入否则会让环境依赖变多、调试难度变大。从实际答辩角度看你能把 HDFS 的目录设计和 Spark 的计算流程讲清楚就已经能证明大数据分析能力了。堆组件只会让问题变得复杂不会让分数变高。3. 从零跑通一个最小可用的空气质量分析预测流程3.1 先把环境跑在 local 模式如果你还没有搭过集群不要一上来就部署三台机器的伪分布式集群。那样很可能浪费大量时间在配置、网络和权限问题上而不是真正做数据分析。我更建议的环境路线是在自己的笔记本上安装 Java、Hadoop先跑通 HDFS 的启动和文件上传下载安装 Spark先用 local 模式跑spark-shell或spark-submit用一个小样本 CSV 文件验证 Spark 能正常读取并执行 SQL把这个流程稳定之后再考虑要不要做集群部署。在常见实践中spark.local 模式不需要 YARN也不需要配置很多参数只要 Java 版本和 Spark 版本兼容启动后就可以在本地进程里运行。等你的数据链路和代码逻辑都验证过了再提交到集群可以避免把环境问题和业务问题搅在一起。如果你已经有一个 Hadoop 集群那么要注意 Spark 和 Hadoop 的版本兼容性。不同发行版之间的配置差别很大不要盲目复制网上的安装命令。至少要先看官方文档确认 Spark 版本支持的 Hadoop 版本范围。3.2 数据清洗统一时间、处理缺失值、规范化站点假设你已经把河南省各站点的空气质量数据导入到 HDFS 的某个目录下比如/data/airquality/raw接下来要做的是 Spark 数据清洗。这里有一个通用处理思路核心是三个动作第一步统一 Schema。Spark 读 CSV 时最好手动指定schema否则时间字段会被当成字符串数值字段可能被自动推断为 Double 或 Long后续容易出现类型不匹配。from pyspark.sql.types import StructType, StructField, StringType, DoubleType, IntegerType, TimestampType schema StructType([ StructField(station_code, StringType(), True), StructField(city, StringType(), True), StructField(monitor_time, TimestampType(), True), StructField(pm25, DoubleType(), True), StructField(pm10, DoubleType(), True), StructField(so2, DoubleType(), True), StructField(no2, DoubleType(), True), StructField(o3, DoubleType(), True), StructField(co, DoubleType(), True), StructField(aqi, IntegerType(), True) ]) df spark.read.option(header, True).schema(schema).csv(/data/airquality/raw)第二步处理缺失值。最简单的策略是先按站点分组查看每一列的缺失率。如果缺失比例很高比如超过 30%那这一列作为特征时就要谨慎。如果只是偶发缺失可以用前后均值插值或者直接保留让 CatBoost 自己处理。第三步统一站点编码。站点名称可能有中文别名、旧编码、新编码建议在 Spark 里维护一张站点维表用join把原始表关联到标准编码。station_dim spark.read.csv(/data/airquality/station_dim.csv, headerTrue) df df.join(station_dim, onstation_code, howleft)清洗完成后把结果写成 Parquet 格式方便后续快速读取df.write.mode(overwrite).parquet(/data/airquality/clean)关键是每一步都要打印处理前后的行数和关键字段的数值范围。不要等建模时才发现数据有问题。3.3 特征工程滞后特征、滑动窗口和时间周期空气质量预测不能只用当前时刻的污染物浓度还要用过去一段时间的变化趋势。所以特征工程非常重要。常见的特征包括滞后特征当前时刻前 1 小时、3 小时、6 小时、12 小时、24 小时的 PM2.5 浓度滑动窗口特征过去 6 小时、24 小时的平均值、最大值、最小值、标准差时间周期特征小时、星期、月份、是否节假日、季节气象特征温度、湿度、风速、风向、气压如果能获取的话。在 Spark 里生成滞后特征可以用窗口函数from pyspark.sql.window import Window from pyspark.sql.functions import lag, col windowSpec Window.partitionBy(station_code).orderBy(monitor_time) df_feature df.withColumn(pm25_lag1, lag(pm25, 1).over(windowSpec)) \ .withColumn(pm25_lag3, lag(pm25, 3).over(windowSpec)) \ .withColumn(pm25_lag24, lag(pm25, 24).over(windowSpec))滑动窗口特征可以用avg、max、min配合rangeBetween来写但要注意窗口大小和计算成本。对几千万行数据来说窗口计算在 Spark 里是能接受的但不要一次生成太多极端复杂的窗口否则 Shuffle 开销会变大。做完特征后先看一眼describe()或者采样几条数据确认特征值不是全为空也没有离谱的极值。3.4 CatBoost 模型训练参数不是越多越好特征宽表准备好了接下来就是模型训练。CatBoost 可以直接读取 Pandas DataFrame所以你可以先在 Spark 里把特征表转成 Pandas需要保证数据量能在单机内存中承载或者直接用toPandas()导出抽样数据。如果特征表太大单机放不下可以先按时间切片用最近两年的数据训练前面的数据做验证而不是一次性全量导出。一个通用训练示例是import pandas as pd from catboost import CatBoostRegressor, Pool feature_cols [station_code, hour, weekday, month, pm25_lag1, pm25_lag3, pm25_lag24, temp, humidity, wind_speed] X df_pandas[feature_cols] y df_pandas[pm25_future] train_pool Pool(X_train, y_train, cat_features[station_code, hour, weekday, month]) model CatBoostRegressor( iterations1000, learning_rate0.05, depth6, eval_metricRMSE, random_seed42, verbose100 ) model.fit(train_pool, eval_setvalid_pool)要注意几个默认参数的理解iterations是树的数量。不是越大越好要看验证集是否继续下降可以用 Early Stopping。learning_rate和学习步长有关。学习率越小需要越多树训练时间越长。depth控制树的复杂度。深度过大容易过拟合空气质量数据里 6 到 8 通常够用。cat_features必须明确指定类别特征否则像“站点编码”这样的字段可能被当作数值连续特征模型解释性会变差。模型评估不要只看 RMSE还要看验证集在不同污染等级下的表现。尤其要关注高污染日因为这才是预警系统真正关心的场景。如果一个模型在普通天预测很准但在重污染天误差很大那实际价值会打折扣。3.5 结果落库和可视化让结论看得见模型输出的是预测值。为了让这个系统“像一个大数据分析系统”还需要把预测结果和实际值整理成可视化图表。建议把结果写回 MySQL 或者直接保存为 Parquet然后用 Superset 或 ECharts 做图表。展示内容至少包含各站点 PM2.5 实际值与预测值的时序对比曲线河南省各城市近 7 天 AQI 均值排名模型特征重要性排序预测误差在不同时间段的分布。这里有一个很实在的经验可视化不是锦上添花而是让评委快速理解你做了什么的最短路径。如果只有一堆训练代码和 RMSE别人是看不出你做了什么的。但如果你能展示出“周一早高峰郑州市 PM2.5 预测偏高”这种具体观察整个项目的完成度会明显提升。4. 从单机到集群迁移 Spark 项目时最容易踩的四个坑4.1 路径和依赖导致的启动失败第一个大坑是路径问题。很多人会在提交 Spark 任务时看到类似这样的报错jar does not exist or is not a normal file: /usr/local/hadoop/share/hadoop/m...这个报错通常不是代码逻辑问题而是 Spark 找不到 Hadoop 依赖包或者提交命令里的--jar路径写错了。排查思路很简单先确认 Hadoop 安装路径是不是你引用的那个路径再确认spark-submit参数里有没有多余或错误的--jar最后确认 classpath 是否正确是否有版本冲突。如果你用的是 Python API还要注意PYSPARK_PYTHON和PYSPARK_DRIVER_PYTHON的环境变量是否指向正确的 Python 解释器。很多时候任务一提交就报错不是代码错而是 driver 端和 executor 端的 Python 环境不一致。4.2 Spark OOM先查分区和资源不要盲目加内存当数据量上来以后最容易遇到的问题就是 OOM内存溢出。但一看到 OOM 就去调spark.executor.memory通常并不能根治问题。你需要先判断 OOM 发生在哪个阶段如果是读取阶段 OOM说明单个分区加载的数据量太大需要增加分区数如果是 Shuffle 阶段 OOM说明spark.sql.shuffle.partitions太小或者某个 key 的数据过于集中如果是 driver OOM说明collect()到本地时数据量超出了 driver 内存这种情况要避免把全量数据toPandas改为抽样或分批处理。一个更稳妥的做法是先用sample(0.1)跑通全流程观察每个 stage 的时间和输入数据量再加入更多数据。不要一开始就用全量数据调参否则你很难分辨是资源不足还是代码有问题。4.3 数据倾斜和小文件按站点聚合时尤其明显空气质量数据按站点聚合时不同城市的数据量可能差别很大。比如郑州可能有多个国控站点数据量是其他城市的好几倍这样在做groupBy(city)或窗口函数时就会出现数据倾斜一个 executor 忙死其他 executor 空闲。解决数据倾斜比较有效的方法有增加rand()扰动先打散 key再做二次聚合使用repartition(partitionExprs[city])手动调整分区数把倾斜的站点单独拆出来处理最后再合并结果。另一个容易被忽略的是小文件问题。如果清洗完数据后写入 HDFS 时用了太小的分区数会产生大量小文件。小文件不会让任务直接报错但会让下一次读取和 NameNode 管理变慢。建议写入前用coalesce()或repartition()控制文件数量尽量生成 128MB 左右的文件块。4.4 验证链路先小样本再单机再集群最后全量这一步是我最想强调的。如果你一开始就在集群上跑全量数据出了问题以后很难定位。验证顺序应该是一层一层放大的先用 1000 条样本在 local 模式跑通确认 Schema、清洗逻辑和模型输入输出再用单机模式跑一个站点的完整数据确认特征工程时间和模型训练能完成然后提交到集群用limit(100000)或sample(0.01)验证分布式执行确认没有路径错误、OOM、数据倾斜后再放开全量数据。每跑一个阶段都要把中间结果的行数和关键字段统计值打印到日志里。这样一旦后面出问题你能快速判断是输入变了、清洗逻辑错了还是资源不够。5. 这个项目适合哪些人不适合哪些人5.1 适合的场景和人群如果你满足下面这几个条件这个选题是合适的你至少有一定 Python 基础会写简单的 Pandas 处理代码你对 Hadoop/Spark 没有抵触情绪愿意花时间安装和学习你不追求把模型的准确率做到极致而是更想体验一个“从数据到系统”的完整流程你的毕业设计时间分配比较充足前一个月可以接受在环境配置上花时间你想在简历上写“熟悉 Spark 数据处理掌握 CatBoost 建模”这个项目比玩具项目更有支撑力。这类项目的答辩亮点通常在“工程性”而不是“算法创新”。你能把 HDFS 的存储设计、Spark 的数据清洗逻辑、特征构造的窗口计算、CatBoost 的结果解释串起来已经很完整了。5.2 不适合的场景反过来如果你有下面任何一种情况建议谨慎选择你只有两周时间之前从没接触过 Spark那环境搭建和踩坑就会耗尽你的时间你只拿到了很小的一份数据比如几百条那用 Hadoop Spark 会很刻意答辩时容易被问“为什么不用 Excel”你想做一个实时系统对延迟要求很高那 Spark 批处理加 CatBoost 并不是最优方案你完全没有数据源。如果拿不到至少连续一年以上的河南省空气质量数据分析结果很难有说服力。还有一个边界要说明这个项目不是图像识别、不是自然语言处理、也不是实时推荐系统它更适合作为“大数据分析 机器学习在环境科学中的应用”类选题而不是万能选题。5.3 毕业设计之外的长期价值从长期角度看这个项目真正值得沉淀的不是某一门技术而是“如何把数据存储、数据清洗、特征工程、模型训练、结果可视化串成一条可复用流水线”的思维方式。你做完之后会发现很多模块是可以抽出来复用的HDFS 目录设计可以复用到其他分析项目Spark 数据清洗模板可以复用到交通、气象、电商日志数据CatBoost 特征工程和参数理解可以复用到其他表格型预测任务构建特征宽表和评估验证集的思路远比一个调好的模型有价值。甚至可以说这个项目真正让你学到的是“工程化思考”如何先定边界、再定方案、最后落地。这种能力放到真实工作中比多记住一个算法 API 有用得多。如果你以后想深入还可以把 Spark 读取 Redis 这类能力加进去用来做实时预测缓存也可以把 HDFS 扩展成多节点集群学习 NameNode 高可用和数据均衡但那是后话了。对现在来说先踏踏实实把链路跑通是唯一值得做的事。选这个题之前不要先问“CatBoost 准不准”先问自己一个问题如果我拿到的是 1000 万行空气质量数据我能不能在两天内把它变成一份可训练的特征宽表如果能这个题你就选对了如果不能那么你现在需要做的不是继续选型而是先去把一条 1 万行的小样本链路跑通。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门