ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Python+Spark全国历史气象数据分析实战:从环境搭建到清洗聚合

Python+Spark全国历史气象数据分析实战:从环境搭建到清洗聚合 简介这是一个基于Spark的全国历史气象数据分析项目面向计算机相关专业学生、教师及科研工作者完整覆盖气象数据清洗、站点统计、MySQL存储、地图可视化等开发环节从原始数据导入到结果展示形成闭环可直接作为毕业设计、课程设计或大数据入门实战参考。资源共75个文件压缩包约2.46MB结构上以5个Python源码为核心配合35个Markdown说明文档、7个PNG成果图表、7个TXT数据文件以及若干XML配置与备份文件既提供可运行的脚本也有分步讲解和图示结果便于对照理解。项目内置全国2018年最高/最低/平均气温与降水量分布图、历年平均气温与降水量变化曲线等成型输出数据源和可视化结果一应俱全并附赠额外资料包已有29人学习下载反馈显示环境配置与运行效果良好。对于希望快速上手Spark处理真实气象数据的开发者可基于现有代码二次开发或直接用于课题演示与课程报告。1. PythonSpark全国历史气象数据分析这个毕设选题值不值得做“PythonSpark做全国历史气象数据分析”几乎是毕业设计和课程设计里最稳的一个方向数据好找、量级够大、能展示的东西也多。我拆过的这套项目核心不是写几百行复杂算法而是把一件事跑通——用Spark把全国多年份、多站点的气象CSV数据清洗干净再按省份和年份算出年均温、年降水量、极端温度这类指标最后落盘成可汇报的结果。适合两类人一类是拿到几GB气象数据、Pandas一读就卡死的新手另一类是想要一份完整Spark分析案例做参照的从业者。下面按我实际拆项目的顺序把数据格式、环境搭建、清洗聚合和坑位一次讲透。2. 项目地基从气象数据格式到Spark开发环境的落地选型2.1 全国历史气象数据长什么样字段、量级与脏数据要动手之前得先知道这个项目到底处理什么数据。全国历史气象数据最常见的形态是全国地面气象站日值或逐小时数据集一般按站点分文件或者按年份打包成一个大的CSV目录。字段基本固定区站号五位数字比如54511是北京南郊站、年、月、日、平均气温、最高气温、最低气温、降水量、平均风速。区站号前两位隐含了地理区域不过做省份维度分析时通常还得配一张站点元数据表里面有省份名称、纬度、经度和海拔。数据量级方面以国家级气象站日值数据为例全国约有两千多个站点取二十年的记录就是两千万行上下。如果换成逐小时数据行数会涨到几亿单个CSV文件几个GB很常见。这也是这个题目为什么不建议用Pandas硬扛的根本原因——数据规模决定了技术选型。脏数据是这题真正花时间的地方。常见的有三类缺失值用-9999、32766、99999这类占位符填充气温出现物理上不可能的异常值比如夏天某站跑出-80度同一个站点不同年份编码重排或者重复记录。这类数据没有标准答案项目里的常规策略是先看字段统计分布再按阈值过滤规则宁可宽松一点也不能把有效记录误删。2.2 为什么选Spark而不硬刚Pandas先给一个反直觉的结论如果数据只有几百MBSpark单机跑并不比Pandas快甚至更慢。那为什么这个项目选Spark因为气象数据的特点是跨年份、跨站点、文件数量多量级不稳定。Pandas读取时需要全量载入内存千万行级别的CSV在8G内存笔记本上很快就开始swap任务卡死还不知道原因。Spark把数据按分区切分每个executor只处理一部分数据内存压力被摊开多文件场景下直接用通配符读取不用手工合并。另一个理由是Spark DataFrame自带的Catalyst优化器对这类宽表查询有谓词下推和列裁剪实际跑聚合时比手写Pandas循环干净得多。下面这张对比表可以快速说明差距对比点PandasSpark DataFrame内存占用全量载入数据一大概率swap分区处理内存压力分散千万行聚合容易内存溢出通常几十秒到几分钟多文件读取需手工合并或循环读取通配符和目录直接读学习成本低中等需要理解分区与Shuffle这不是说Pandas没用而是这个场景里Spark的容错和扩展性更适合。项目后面把聚合结果转回Pandas做可视化正好各取所长。2.3 开发环境搭建Python、Spark与JDK的版本搭配环境搭错是新手第一个翻车现场。我拆这套项目时用的是一套比较稳的组合JDK 1.8 Spark 3.2.0 Python 3.8pyspark版本和Spark保持一致。不少教程里的版本乱搭最容易踩的是JDK版本过高导致Spark起不来或者Python 3.10以上和旧版pyspark之间出现兼容报错。安装顺序我一般这样走先装JDK 1.8并配置JAVA_HOME再下载spark-3.2.0-bin-hadoop3.2解压到固定目录并设置SPARK_HOME然后执行pip install pyspark3.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, truncateFalse)通配符路径是最省事的做法不用写循环合并。每个文件列顺序不一致时显式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()把清洗结果缓存住后续多次聚合就不用重新跑一遍清洗流程。数据里还有一层关键信息在站点表里需要单独读入准备joinstation_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个分区对单机处理两千万行数据是合理的。这里用的是最简单直接的ProvinceYear分组。如果想看城市粒度把分组键换成station再join站点表把省份和城市名一起带出来即可。四个聚合指标分别对应毕业设计里最常被问到的年均温、年降水量、极端最高最低温一张表全齐了。3.3 落盘分区Parquet与CSV格式的选择聚合结果要落盘才能交给下一步可视化或者答辩展示。Parquet是首选格式列式存储压缩率高查询时只读需要的列速度明显优于CSVoutput_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.partitions200对这个场景不一定合适倾斜的分区就爆了。解决先把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任务都强制先跑完这一步再放开全量这个习惯在气象数据这类脏数据多的场景里救了我不止一次。希望帮到你。本文还有配套的精品资源点击获取
返回列表