简介:这是一个基于Spark的全国历史气象数据分析项目,面向计算机相关专业学生、教师及科研工作者,完整覆盖气象数据清洗、站点统计、MySQL存储、地图可视化等开发环节,从原始数据导入到结果展示形成闭环,可直接作为毕业设计、课程设计或大数据入门实战参考。资源共75个文件,压缩包约2.46MB,结构上以5个Python源码为核心,配合35个Markdown说明文档、7个PNG成果图表、7个TXT数据文件以及若干XML配置与备份文件;既提供可运行的脚本,也有分步讲解和图示结果,便于对照理解。项目内置全国2018年最高/最低/平均气温与降水量分布图、历年平均气温与降水量变化曲线等成型输出,数据源和可视化结果一应俱全,并附赠额外资料包;已有29人学习下载,反馈显示环境配置与运行效果良好。对于希望快速上手Spark处理真实气象数据的开发者,可基于现有代码二次开发,或直接用于课题演示与课程报告。
1. Python+Spark全国历史气象数据分析:这个毕设选题值不值得做
“Python+Spark做全国历史气象数据分析”几乎是毕业设计和课程设计里最稳的一个方向:数据好找、量级够大、能展示的东西也多。我拆过的这套项目,核心不是写几百行复杂算法,而是把一件事跑通——用Spark把全国多年份、多站点的气象CSV数据清洗干净,再按省份和年份算出年均温、年降水量、极端温度这类指标,最后落盘成可汇报的结果。适合两类人:一类是拿到几GB气象数据、Pandas一读就卡死的新手;另一类是想要一份完整Spark分析案例做参照的从业者。下面按我实际拆项目的顺序,把数据格式、环境搭建、清洗聚合和坑位一次讲透。
2. 项目地基:从气象数据格式到Spark开发环境的落地选型
2.1 全国历史气象数据长什么样:字段、量级与脏数据
要动手之前,得先知道这个项目到底处理什么数据。全国历史气象数据最常见的形态是全国地面气象站日值或逐小时数据集,一般按站点分文件,或者按年份打包成一个大的CSV目录。字段基本固定:区站号(五位数字,比如54511是北京南郊站)、年、月、日、平均气温、最高气温、最低气温、降水量、平均风速。区站号前两位隐含了地理区域,不过做省份维度分析时,通常还得配一张站点元数据表,里面有省份名称、纬度、经度和海拔。
数据量级方面,以国家级气象站日值数据为例,全国约有两千多个站点,取二十年的记录就是两千万行上下。如果换成逐小时数据,行数会涨到几亿,单个CSV文件几个GB很常见。这也是这个题目为什么不建议用Pandas硬扛的根本原因——数据规模决定了技术选型。
脏数据是这题真正花时间的地方。常见的有三类:缺失值用-9999、32766、99999这类占位符填充;气温出现物理上不可能的异常值(比如夏天某站跑出-80度);同一个站点不同年份编码重排或者重复记录。这类数据没有标准答案,项目里的常规策略是先看字段统计分布,再按阈值过滤,规则宁可宽松一点,也不能把有效记录误删。
2.2 为什么选Spark而不硬刚Pandas
先给一个反直觉的结论:如果数据只有几百MB,Spark单机跑并不比Pandas快,甚至更慢。那为什么这个项目选Spark?因为气象数据的特点是跨年份、跨站点、文件数量多,量级不稳定。Pandas读取时需要全量载入内存,千万行级别的CSV在8G内存笔记本上很快就开始swap,任务卡死还不知道原因。Spark把数据按分区切分,每个executor只处理一部分数据,内存压力被摊开,多文件场景下直接用通配符读取,不用手工合并。
另一个理由是Spark DataFrame自带的Catalyst优化器,对这类宽表查询有谓词下推和列裁剪,实际跑聚合时比手写Pandas循环干净得多。下面这张对比表可以快速说明差距:
| 对比点 | Pandas | Spark DataFrame |
|---|---|---|
| 内存占用 | 全量载入,数据一大概率swap | 分区处理,内存压力分散 |
| 千万行聚合 | 容易内存溢出 | 通常几十秒到几分钟 |
| 多文件读取 | 需手工合并或循环读取 | 通配符和目录直接读 |
| 学习成本 | 低 | 中等,需要理解分区与Shuffle |
这不是说Pandas没用,而是这个场景里Spark的容错和扩展性更适合。项目后面把聚合结果转回Pandas做可视化,正好各取所长。
2.3 开发环境搭建:Python、Spark与JDK的版本搭配
环境搭错是新手第一个翻车现场。我拆这套项目时用的是一套比较稳的组合:JDK 1.8 + Spark 3.2.0 + Python 3.8,pyspark版本和Spark保持一致。不少教程里的版本乱搭,最容易踩的是JDK版本过高导致Spark起不来,或者Python 3.10以上和旧版pyspark之间出现兼容报错。
安装顺序我一般这样走:先装JDK 1.8并配置JAVA_HOME,再下载spark-3.2.0-bin-hadoop3.2解压到固定目录并设置SPARK_HOME,然后执行pip install pyspark==3.2.0,最后在命令行敲pyspark确认能进入交互式界面。这一步跑通,后面就顺了。
2.4 从CSV到DataFrame:显式Schema、编码与通配符读取
环境准备好之后,第一步是初始化SparkSession。这里有几个参数值得说清楚:
from pyspark.sql import SparkSession spark = (SparkSession.builder .appName("ChinaWeatherAnalysis") .master("local[*]") # 本地开发用所有CPU核;部署集群改成yarn .config("spark.sql.shuffle.partitions", "48") .config("spark.driver.memory", "4g") .config("spark.executor.memory", "6g") .getOrCreate())master设置成local[*]是开发阶段最省事的写法,表示用本机所有可用核跑;正式部署到集群时改成yarn。spark.sql.shuffle.partitions是Shuffle后的分区数,默认200,在单机开发时往往偏大,这里按核数调整为48。driver和executor内存按机器配置给,8G笔记本上driver给4g已经够用。
数据读取这一步,我强烈建议不要用inferSchema自动推断类型。气象CSV里常有不规则字符串,自动推断既慢又容易在后续join时暴露类型不匹配问题。显式定义Schema更可控:
from pyspark.sql.types import (StructType, StructField, StringType, IntegerType, DoubleType) weather_schema = StructType([ StructField("station", StringType(), True), StructField("year", IntegerType(), True), StructField("month", IntegerType(), True), StructField("day", IntegerType(), True), StructField("temp_avg", DoubleType(), True), StructField("temp_max", DoubleType(), True), StructField("temp_min", DoubleType(), True), StructField("precip", DoubleType(), True), StructField("wind", DoubleType(), True) ]) df = (spark.read .option("header", True) .option("encoding", "UTF-8") .schema(weather_schema) .csv("input/*.csv")) # 通配符一次读入所有年份文件 df.printSchema() df.show(5, truncate=False)通配符路径是最省事的做法,不用写循环合并。每个文件列顺序不一致时,显式Schema能保证字段按名字对应,而不是按位置错位。encoding这个参数特别容易被忽略,如果源文件是GBK编码,这里就要改成GBK,否则中文列名和部分值会变成乱码。
3. 核心实战:全国站点数据的清洗、聚合与分区落盘
3.1 清洗:缺失值、异常值、重复记录一网打尽
Spark开发里有个共识:写分析逻辑之前,先把数据洗成可信任的状态。气象数据的清洗核心就三件事:过滤占位符、过滤物理异常值、去重。
from pyspark.sql.functions import col, concat_ws, to_date df_clean = (df .filter(col("year") >= 1951) # 丢掉早期完整性差的记录 .filter((col("temp_avg") >= -60) & (col("temp_avg") <= 60)) .filter(col("precip") >= 0) # 降水量不可能是负值 .withColumn("date", to_date( concat_ws("-", col("year").cast("string"), col("month").cast("string"), col("day").cast("string")), "yyyy-M-d")) .dropDuplicates(["station", "year", "month", "day"]) ) print("清洗前行数:", df.count()) df_clean.cache() print("清洗后行数:", df_clean.count())温度阈值用的是全国历史极端气温参照:最冷在漠河附近约-52度,最热在吐鲁番约49度,所以取-60到60是安全的物理范围,既能过滤异常,又不会误杀有效记录。precip >= 0这行代码顺带把-9999这类占位符过滤掉了,因为缺失降水在源文件里通常填负数。
to_date的格式坑特别多。如果月份和日期是1位数,用"yyyy-MM-dd"解析会得到null,因为格式串里的MM期望两位数。这里用"yyyy-M-d"对单数字串天然兼容。最后按站点和日期去重,保证同一站点同一天只保留一条记录。
清洗完调用count()是因为Spark是惰性计算的,不触发action之前,前面的filter根本不会真的执行。cache()把清洗结果缓存住,后续多次聚合就不用重新跑一遍清洗流程。
数据里还有一层关键信息在站点表里,需要单独读入准备join:
station_df = (spark.read .option("header", True) .option("encoding", "UTF-8") .csv("input/stations.csv")) station_df = station_df.select( col("station").cast("string"), col("province").cast("string"), col("lat").cast("double"), col("lon").cast("double") ) station_df.show(5)站点表字段比较杂,这里只挑后面用到的四列,顺便把类型强制固定。province字段在后续省份聚合里是分组键,经纬度则留给最后的可视化。
3.2 聚合:按省份和年份算年均温、降水量与极端温度
清洗完的数据还需要join站点表才能做省份维度分析。join时用左连接,因为理论上可能存在清洗后保留下来的站点不在站点元数据表里的情况:
from pyspark.sql.functions import avg, sum, min, max df_joined = df_clean.join( station_df.select("station", "province", "lat", "lon"), "station", "left" ).filter(col("province").isNotNull()) # 丢不掉无省份信息的记录 annual_stats = (df_joined .groupBy("province", "year") .agg( avg("temp_avg").alias("annual_avg_temp"), sum("precip").alias("annual_sum_precip"), max("temp_max").alias("annual_max_temp"), min("temp_min").alias("annual_min_temp") ) .orderBy("province", "year") ) annual_stats.show(10)groupBy是Spark作业里最耗时的一环,它会触发一次全量Shuffle,把所有相同省份和年份的数据汇聚到同一个分区里计算。前面设定的spark.sql.shuffle.partitions在这里生效,48个分区对单机处理两千万行数据是合理的。
这里用的是最简单直接的Province+Year分组。如果想看城市粒度,把分组键换成station,再join站点表把省份和城市名一起带出来即可。四个聚合指标分别对应毕业设计里最常被问到的年均温、年降水量、极端最高最低温,一张表全齐了。
3.3 落盘:分区Parquet与CSV格式的选择
聚合结果要落盘才能交给下一步可视化或者答辩展示。Parquet是首选格式,列式存储压缩率高,查询时只读需要的列,速度明显优于CSV:
output_path = "output/annual_stats" annual_stats.write \ .mode("overwrite") \ .partitionBy("year") \ .parquet(output_path) print("已写入:", output_path)partitionBy("year")的好处是后续按年份过滤时直接跳过无关目录。但要注意一个副作用:年份多时每个年份目录下文件数量可能膨胀,产生大量小文件。这个坑第四节专门展开。如果只是想给老师一份能直接用Excel打开的CSV,就把Parquet读回来再单独写一次:
spark.read.parquet(output_path) \ .coalesce(1) \ .write.mode("overwrite") \ .option("encoding", "UTF-8") \ .csv("output/annual_stats_csv")coalesce(1)把数据合并到单分区再写,保证只产出一个CSV文件,而不是一堆碎片。这里注意coalesce是窄依赖操作,比repartition便宜,但会降低并行度,只适合在结果集已经很小的时候用。
4. 避坑:气象数据跑Spark最容易翻车的四个现场
4.1 中文列名乱码与编码推断错误
现象:CSV读进来后,中文字段名变成乱码,字符串列的值全部显示为null,但printSchema看到的类型又是对的。
原因:源文件是GBK或ANSI编码,Spark默认按UTF-8解析,中文多字节字符被拆成了无法识别的序列,值自然null。
解决:读取时显式指定编码。做法是在读CSV的option里加encoding参数,值改成GBK或GB18030;也可以先把源文件用文本编辑器批量转成UTF-8。我一般直接把编码参数跟着文件名一起写进配置文件,避免每次手改。
4.2 groupBy之后Stage一直重试,Executor报OOM
现象:跑省份聚合时,任务卡在Shuffle阶段,Stage反复重试,最后某几个Executor直接OOM退出。
原因:气象站点分布极不均匀,东部省份站点密集,西部省份稀疏。groupBy province时,站点多的省份数据量集中到同一个分区,单分区数据量过大,加上默认的spark.sql.shuffle.partitions=200对这个场景不一定合适,倾斜的分区就爆了。
解决:先把shuffle分区数调到48或96,和集群并行度匹配;如果仍然倾斜,可以对groupBy键做预分区,比如在groupBy之前repartition(col("province")),让相同省份的数据提前分布到多个分区。极端情况下还可以给倾斜key加盐,但气象数据场景一般用不到那么重的手段。
4.3 日期解析全null:格式没对齐
现象:to_date之后date列全为null,按时间筛选结果为空,但原始字符串看起来没问题。
原因:日期字段的实际格式和格式串不匹配。最常见两种:源文件里是20240101这种八位整数或字符串,却用了"yyyy-M-d"解析;另一种是月份、日不补零,用"yyyy-MM-dd"解析单数字段直接失败。
解决:先统一转成字符串,再按实际格式解析。八位数字用to_date(col("date_str"), "yyyyMMdd");单数字月日就用不带前导零的格式串。稳妥的做法是在清洗阶段直接对原始字段做一次格式探测,打印几条样本确认后再写死格式串。
4.4 输出目录几百上千个小文件,读起来和看起来都难受
现象:Parquet写出后,output目录下堆了上百个小文件,每个几十KB,后续读入时任务数量暴涨,处理反而变慢。
原因:Shuffle结束后分区数决定了文件数,默认200个分区就会产生200个文件;再加上partitionBy("year")按年份拆目录,年份越多,每个年份目录下的文件碎片越多。
解决:写出前先repartition或coalesce控制分区数。小文件问题的药方是写出前重新分区:
annual_stats.repartition(4) \ .write.mode("overwrite") \ .partitionBy("year") \ .parquet("output/annual_stats")repartition(4)让每个年份目录下最多4个文件,目录数量由年份数决定。实际项目的取舍是按字段重要程度来定:省份多就按省份分,年份多就按年份分,别两个一起细化,否则小文件灾难一定会回来。
5. 验证手段:先跑通一个站点再放开全量
整套项目跑完后,最怕的是结果看着有数,实际经不起追问。我自己习惯在放开全量数据之前,先挑一个熟悉的气象站做全流程验证。这个方法成本低、定位快,能挡住绝大多数翻车。
验证方式很直接:取北京南郊站(站号54511)的数据作为样本,跑同一套清洗和聚合逻辑:
sample_df = df_clean.filter(col("station") == "54511").cache() sample_stats = (sample_df .groupBy("station", "year") .agg( avg("temp_avg").alias("annual_avg_temp"), sum("precip").alias("annual_sum_precip") ) .orderBy("year")) sample_stats.show(20)然后对答案。日值数据一年约365条记录,20年就有7000多行;按站和年份聚合后,应该有20行左右。如果行数和预期对不上,说明清洗阈值过狠误杀了记录,或者日期解析出问题导致部分年份数据消失。如果行数对了但某个年份温度明显离谱,就单独把那一年的原始记录拉出来看。
这一步跑通后,再把filter条件去掉,对全量数据执行同一套代码。另一个便宜好用的基线校验是直接检查落盘结果能否被重新读回,并核对聚合后的记录数:
check_df = spark.read.parquet("output/annual_stats") print("省份数量:", check_df.select("province").distinct().count()) print("总记录数:", check_df.count())省份数量应该和站点表里的省份数一致,总记录数应该约等于省份数乘以年份跨度。这两条对不上,说明join或分组有问题。从那以后我每次跑Spark任务,都强制先跑完这一步再放开全量,这个习惯在气象数据这类脏数据多的场景里救了我不止一次。希望帮到你。
本文还有配套的精品资源,点击获取