ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Python+Spark机器学习天气预测系统:从数据清洗到模型部署全解析

2026/9/11 8:43:44 拓冰建站 浏览量
Python+Spark机器学习天气预测系统:从数据清洗到模型部署全解析 简介大数据天气预测毕业设计项目资料包基于Python与Spark构建机器学习预测系统面向计算机相关专业学生、教师及企业开发人员适用于毕业设计、课程设计、项目演示也适合大数据方向初学者进阶。压缩包共15个文件整体约7.86MB包含Python爬虫脚本、Spark/Scala处理示例、Java工具类、Markdown文档、授权txt及9张png架构流程图覆盖天气数据采集、缓存存储、特征处理、模型训练与预测展示的完整链路。项目附带详细文档与全部资料代码经测试运行成功并获导师指导认可可直接作为毕设基准方案也可扩展算法或优化工程细节其中Jedis缓存工具与Spark排序归约示例可单独学习缓存读写、键值排序与归并聚合等常用大数据操作。目前已有161人学习适合需要完整参考实现、配套说明与可运行源码的毕业生或初学者。1. 天气预测系统PythonSpark 适合做毕设也值得做扎实如果毕业设计只跑通 sklearn 的线性回归答辩时很难让人相信你接触过大数据。天气预测是一个时间序列回归问题把它和 Python、Spark 机器学习组合起来恰好能覆盖数据采集、清洗、分布式计算、模型训练、结果评估的完整链路这也是“基于 PythonSpark 机器学习天气预测系统”这类题目常见的立意。这篇不打算逐行讲解一份现成源码而是把做这类系统真正绕不开的环节数据怎么来、Spark 怎么处理、模型怎么选、集群怎么跑、结果怎么验证一次讲清楚。准备大数据方向毕设的人以及想用 MLlib 处理真实数据的开发者都可以按这个顺序复现。2. 从气象数据到 Spark 可用的 DataFrame先搭对处理链路2.1 数据源选择公开数据与字段设计天气预测系统的第一步不是写代码而是选数据。公开气象数据源比较多国内可以用中国气象数据网国外可以用 NOAA 的 GSOD 或者 Kaggle 整理好的历史天气集。GSOD 按年份组织包含温度、气压、湿度、风向风速和站点信息日粒度适合做日级预测。常见做法是选一个站点或一个小区域拉取最近 3 到 5 年的数据而不是一开始就全量下载。原因是毕设要验证的是“大数据处理链路”数据量够用和复杂度正好即可数据一旦过杂排错成本会吃掉写文档的时间。字段建议控制在 8 到 12 个超过这个数量特征工程和文档会很难解释。下表是我常用的初始字段字段名类型说明station_idstring站点编号datedate观测日期tmax / tminfloat最高 / 最低温度prcpfloat降水量wdspfloat平均风速humfloat相对湿度slpfloat海平面气压如果数据源里没有 hum可以先用温度和露点计算相对湿度公式很多选一个简单公式并在文档里注明即可。字段类型尽量在源头定好避免后续 Spark 每次推断。2.2 Python 清洗脚本检查缺失、重复和异常值拿到 CSV 后不要立刻灌进 Spark先用 pandas 做一轮快速检查回答“数据能不能用”。常见问题是日期格式不统一、部分站点缺测多、风速出现负值、同一天有重复记录。下面这段脚本可以完成基础筛查import pandas as pd df pd.read_csv(weather_raw.csv, parse_dates[date], index_coldate) print(df.shape) print(df.isna().mean().sort_values(ascendingFalse)) # 查看缺失率 dup df[df.index.duplicated()] print(f重复日期: {len(dup)}) # 阈值过滤风速不能为负站点缺失率不超过 20% df df[(df[wdsp] 0)] df df[df.groupby(station_id)[tmax].transform(lambda x: x.isna().mean()) 0.2] # 统一日期格式 df df.reset_index() df[date] pd.to_datetime(df[date]).dt.strftime(%Y-%m-%d)代码逻辑是先看每条字段的缺失比例再按阈值过滤异常站点。transform(lambda ...)会把组内缺失率广播回每一行达到按站点过滤的目的。阈值 0.2 不是固定的如果站点数据本来就少可以放宽到 0.5但缺失太多会让模型学到错误模式宁可靠插值补一部分也不要留大量空洞。清洗后按站点和日期排序这一步直接影响下一章滞后特征的正确性。如果顺序没有排好lag取到的可能不是前一天而是上一条记录这样模型指标看着不错实际部署完全不可用。2.3 写入 SparkParquet 和分区策略Spark 读取清洗后的 CSV 很简单但我不推荐直接对 CSV 反复训练。CSV 没有 schema、没有压缩每次 read 都推断类型大文件上效率很低。常见做法是先用 Spark 读一次转成 Parquet 格式落盘并按年份分区from pyspark.sql import SparkSession from pyspark.sql import functions as F spark SparkSession.builder \ .appName(weather-etl) \ .master(local[*]) \ .getOrCreate() sdf spark.read.csv(weather_clean.csv, headerTrue, inferSchemaTrue) sdf sdf.withColumn(year, F.substring(sdf[date], 0, 4)) sdf.write.mode(overwrite).partitionBy(year).parquet(weather_parquet)上面代码中substring(date, 0, 4)取出年份字符串分区键选year的好处是单年数据查询时 Spark 只需要扫描对应目录。要注意withColumn生成的 year 是字符串后续与 date 过滤条件配合时最好把 date 转换成真正的日期类型。分区键不要选 station_id站点数少时分区数量小调度成本反而高于收益。partitionBy会改变文件的目录结构读取时不需要指定分区列Spark 会自动从路径恢复。写完以后用spark.read.parquet(weather_parquet).printSchema()验证 schema检查每个字段类型是否符合预期。如果发现 tmax 被读成 string通常是因为原始表里有空值或特殊字符回到清洗阶段处理不要在这个阶段硬转。3. Spark 机器学习管线特征工程和训练参数3.1 用 Spark MLlib 还是 Pandassklearn怎么选天气预测的数据量级通常不大用 Pandas sklearn 完全能跑那为什么还要用 Spark答案是题目要求体现出大数据处理而且 Spark 的管线可以扩展站点数和年份。实际做法是数据清洗、特征工程在 Spark 上完成模型训练按数据规模决定。如果单站点数据在几万行以内把 Spark 结果转成 pandas 用 sklearn 训练速度更快调试更方便如果想展示分布式机器学习就用 MLlib 的RandomForestRegressor或GBTRegressor。我会两个都做当时写在文档里的结论是sklearn 在单机上可能精度略高但 MLlib 的训练过程可以水平扩展这才符合大数据场景。MLlib 的回归器在分布式计算上做了优化但调参逻辑和 sklearn 类似。下面是VectorAssembler配合Pipeline的典型写法from pyspark.ml.feature import VectorAssembler from pyspark.ml.regression import RandomForestRegressor from pyspark.ml.evaluation import RegressionEvaluator from pyspark.ml import Pipeline feature_cols [lag1, lag2, rolling_mean7, month_sin, month_cos] assembler VectorAssembler(inputColsfeature_cols, outputColfeatures) rf RandomForestRegressor( featuresColfeatures, labelColtmax, numTrees100, maxDepth10, seed42 ) pipeline Pipeline(stages[assembler, rf]) model pipeline.fit(train_df)注意VectorAssembler的输入列必须全部是数值类型字符串列要先做StringIndexer或OneHotEncoder。maxDepth不是越大越好天气数据的有效深度通常在 8 到 12 之间超过 15 很容易把历史日期的偶然波动学进去。3.2 特征工程时序问题不能随机切分天气预测是时间序列问题特征设计比模型选择更重要。做这个系统时我主要构造三类特征。滞后特征是最直接的前 1 天、前 2 天和前 7 天的最高温能捕捉短期惯性和周周期。滚动统计包括近 7 天的最高温均值、最低温均值、降水总和这些特征相当于给模型一个“最近天气走势”的抽象表达。周期特征把日期转成一年中的第几天再计算 sin 和 cos避免 1 月和 12 月在欧氏距离上被强行拉开。周期特征对气候型预测很重要因为气温变化与季节强相关。构造滞后特征要用窗口函数Spark 不像 pandas 可以直接 shift。下面是按站点分区的窗口写法from pyspark.sql.window import Window from pyspark.sql import functions as F w Window.partitionBy(station_id).orderBy(date) train_df train_df.withColumn(lag1, F.lag(tmax, 1).over(w)) train_df train_df.withColumn(lag2, F.lag(tmax, 2).over(w)) train_df train_df.withColumn( rolling_mean7, F.avg(tmax).over(w.rowsBetween(-6, 0)) )参数说明lag(tmax, 1)在当前分区内按日期排序后取上一行最高温。如果该站点第一天没有前值结果是 null后续要过滤或填均值。rowsBetween(-6, 0)是把窗口从当前行向前扩展 6 行宽度为 7 行这正是“近 7 天滚动平均”的精确含义。如果不写partitionBy整个数据会被当作一个分区跨站点的日期顺序会串数据这是常见错误。切分训练集和测试集时严格执行“按时间切分”。毕设里直接randomSplit是错误的因为随机切分会让未来数据混进训练集模型相当于做过一遍“开卷考试”。常见做法是train_df.filter(F.col(date) 2023-01-01)作为训练集剩余作为测试集。3.3 网格搜索和参数调优TrainValidationSplit 的参数MLlib 可以用ParamGridBuilder配合TrainValidationSplit自动搜索参数。下面的网格设计用于对比numTrees和maxDepth两个最重要的参数from pyspark.ml.tuning import ParamGridBuilder, TrainValidationSplit grid ParamGridBuilder() \ .addGrid(rf.numTrees, [50, 100]) \ .addGrid(rf.maxDepth, [5, 10, 15]) \ .build() evaluator RegressionEvaluator(labelColtmax, metricNamermse) tvs TrainValidationSplit( estimatorpipeline, estimatorParamMapsgrid, evaluatorevaluator, trainRatio0.8 ) cv_model tvs.fit(train_df) best_model cv_model.bestModeltrainRatio0.8是把传入的数据按时间顺序再切出 20% 做验证集与外部测试集无关。TrainValidationSplit只做一次划分所以结果有一定随机性如果时间充足可以换成CrossValidator但会让训练时间成倍增加。毕设里用前者即可文档中写明“为了控制计算成本采用单次验证”。调参完成后保存PipelineModel因为里面包含特征处理逻辑。加载时PipelineModel.load(path)不要只保存随机森林模型而丢失特征列的组装方式。还要注意 MLlib 的模型序列化对 Spark 版本敏感训练和预测最好保持同一版本。4. 本地跑通和集群提交天气预测系统的部署与排错4.1 本地环境Python、Spark 的安装与配置问题做这类系统环境安装会占整体时间的一两成主要是 Java、Spark、Python 三个版本匹配问题。新版 Spark 依赖 Java 8 或 11Python 3.8 以上都可以用pip install pyspark能拿到 Spark 的 Python API但不包含 JVM所以要单独安装 JDK。测试环境是否正常的顺序是python --version java -version spark-submit --version三个命令都能输出版本后再跑数据。Windows 上首次运行spark-shell可能一闪而过通常是因为设置了错误的SPARK_HOME或缺少 Hadoop 的本地依赖。如果是本地纯local[*]模式可以先不配置HADOOP_HOME但读写文件时可能触发“Could not locate executable null\bin\winutils.exe”这样的提示。常见解决方法是下载与 Hadoop 版本匹配的winutils.exe放在没有空格的目录并对spark.local.dir做调整。我一般会在脚本开头打印运行模式避免集群和本地配置混淆from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(weather-predict) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 10) \ .getOrCreate() print(当前 master:, spark.sparkContext.master)参数spark.sql.shuffle.partitions10是给本地调试验证用的当数据量只有几千行时默认 200 个分区会调度很多空任务没必要。4.2 集群提交Spark on YARN 的参数怎么设毕设如果只用 local 模式演示效果有限。想展示大数据集群能力最少要有一个 3 节点的 YARN 集群然后把作业提交到 YARN 上。最常见的提交命令是spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 2g \ --num-executors 3 \ --executor-cores 2 \ --driver-memory 1g \ weather_predict.py参数含义executor-memory是每个容器内存上限num-executors是容器数量executor-cores是每个容器占用的虚拟核数driver-memory给提交方和规划任务用。总内存大约是 3×2GB 1GB如果集群节点是 8GB 配置这个默认可以跑如果不够就降成--executor-memory 1g --num-executors 2。注意executor-cores不是设得越大越好太多核会互相抢 CPU通常 1 到 2 核效果最好。提交到 YARN 后作业运行日志不在本地终端里要去yarn logs -applicationId application_xxxx或者集群资源管理页看。常见失败是容器超过内存限制被杀英文日志里会有“Container killed”关键字。这时不要盲目调大内存先把代码里的宽依赖处理掉比如 groupBy 后加repartition(10)。4.3 三个高频坑时区、数据倾斜、特征泄漏第一个坑是时区。CSV 里的日期是字符串Spark 读成 timestamp 或 date 时会用服务器默认时区。如果集群时间和本地不一致日期过滤结果会少一天。建议统一用F.to_date(date, yyyy-MM-dd)并给 SparkSession 设置spark.sql.session.timeZoneUTC让所有计算有明确的时区依据。第二个坑是数据倾斜。按站点做groupBy时大城市站点数据量明显多于其他站点部分 Executor 会成为长尾。天气场景通常把数据压缩成一天一行让每个站点记录数相近如果还有特别大的 key再考虑加 salt key 做两阶段聚合。毕设数据量小一般不会遇到严重倾斜但文档里提到这个优化点会给答辩加分。第三个坑是特征泄漏。泄漏主要发生在滞后特征构造之后没有排序就切分了训练集导致测试集里包含未来的滚动统计值。排查方法是打印特征列的最小日期和最大日期确保测试集最小日期等于切分点。另一点是rolling_mean7里如果包含当天的 tmax并且在训练模型时同时把 tmax 和 rolling_mean7 放进特征列相当于用答案预测答案这个错误很隐蔽。构造滚动特征时要把当天的值排除即窗口为rowsBetween(-7, -1)。5. 用回归评估指标和残差图反推模型问题最后一步不是保存模型而是验证。天气预测是回归问题评估指标建议写三个MAE、RMSE、R²。MAE 解释性最好单位与温度一致RMSE 对大误差敏感R² 表示模型解释的方差比例。三者同时出现在文档里比只写“准确率 95%”专业得多。把预测结果落地后用下面的方式一次性计算from pyspark.ml.evaluation import RegressionEvaluator pred best_model.transform(test_df) pred.select(tmax, prediction).show(10) for metric in [mae, rmse, r2]: evaluator RegressionEvaluator(labelColtmax, predictionColprediction, metricNamemetric) print(metric, evaluator.evaluate(pred))MLlib 的 RegressionEvaluator 支持rmse、mae、r2三种指标直接在循环里调用比自己写公式省事。注意只传入 labelCol、不传 predictionCol 会默认取prediction列如果你的特征列里恰好有这个名称指标看起来会特别小因为模型“作弊”了。指标算完要看残差。把结果收集到本地画真实值-预测值散点图和残差直方图import matplotlib.pyplot as plt pdf pred.select(tmax, prediction).toPandas() pdf[residual] pdf[prediction] - pdf[tmax] fig, ax plt.subplots(1, 2, figsize(12, 4)) ax[0].scatter(pdf[tmax], pdf[prediction], alpha0.3) ax[0].plot([-10, 40], [-10, 40], r--) ax[0].set_xlabel(实际最高温) ax[0].set_ylabel(预测最高温) ax[1].hist(pdf[residual], bins40) ax[1].set_xlabel(残差) plt.tight_layout() plt.savefig(residual.png, dpi150)散点图如果在低温端明显偏离 45 度线说明模型对低温天气的学习不足可以从数据里检查冬季样本量。残差直方图如果出现双峰通常不是模型问题而是训练集和测试集季节不平衡回到第 3 章检查切分时间。进阶做法是按月份统计 MAEfrom pyspark.sql import functions as F pred.withColumn(month, F.month(date)).groupBy(month) \ .agg(F.avg(F.abs(F.col(prediction) - F.col(tmax))).alias(mae_by_month)) \ .orderBy(month) \ .show()输出后如果发现冬季月份 MAE 明显更高就在特征里增加月份相关特征或者对冬季样本做重采样。这套“指标 残差 月度误差”的验证流程比单纯保存模型更有说服力把这些现象和改进写进毕设文档答辩时比自己说“我调过参”有用得多。本文还有配套的精品资源点击获取