ARTICLE DETAIL

资讯详情

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

Spark处理全国70年气象数据:Python分布式分析实战

Spark处理全国70年气象数据:Python分布式分析实战 简介这是一套面向计算机相关专业学生、教师及大数据初学者的Spark实战项目资料以Python为开发语言围绕全国历史气象数据展开分析可用于毕业设计、课程设计、作业或项目立项演示。资源包共75个文件约2.46MB包含Python源码、Markdown说明文档、PNG可视化图表、TXT数据文件、XML配置及答辩PDF等其中源码负责数据读取、清洗与统计图表直观呈现全国气象站气温、降水分布及历年变化趋势文档则记录环境配置与运行说明。项目代码完整、资料齐全已通过测试可正常运行读者可借鉴其Spark分布式计算思路与数据可视化实现也可在此基础上修改扩展功能。目前已有29人学习下载适合需要大数据分析案例参考或毕设选题落地的读者使用。1. 全国历史气象数据遇上 Spark一套能扛住 70 年日值记录的 Python 分析链路全国 2400 多个国家级气象站从 1951 年到现在日值数据累积下来是十亿级记录。单机 pandas 读一个省还能忍一旦要算「全国近 30 年逐年平均气温距平」或者「某月降水极值的空间分布」内存直接爆掉跑一晚上出不来结果。这个标题要解决的就是这件事用 Python 做上层分析、Spark 做分布式计算把全国历史气象数据从原始文本一路跑到可解释的统计结论。适合已经会写 Python、但没真正在集群上跑过气象数据的人也适合手里有 Spark 环境、却不知道拿它分析什么真实数据的人。下面按「数据长什么样 → 怎么清洗 → 怎么算 → 坑在哪」的顺序讲每一步都能照着复现。2. 全国历史气象数据的结构、来源与 Spark 选型理由2.1 气象日值数据的真实长相国内历史气象数据最常见的形态是中国气象数据网提供的日值数据集按站点打包每个站点一个文本文件文件名通常是SURF_CLI_CHN_MUL_DAY-TEM-YYYYMM.TXT这类格式。打开一个文件前几行是站号、经纬度、海拔、数据年份的说明之后才是逐日记录。真正要解析的是数据行字段用空格或制表符分隔典型字段包括站号、纬度、经度、观测日期、平均气温、最高气温、最低气温、降水量、平均风速、日照时数等。这里有个新手最容易翻车的地方不同要素的数据集字段数量不一样气温数据集和降水数据集的列定义完全不同不能拿同一套 schema 硬套。我一般会先抽一个文件用 Python 读前 20 行把列位置和单位确认清楚再写解析逻辑。单位也要注意气温常见是 0.1 摄氏度存储降水量是 0.1 毫米直接当摄氏度用会差 10 倍这种错误在结果里表现为「全国平均气温 150 度」一眼能看出来但如果只算距平就可能被掩盖。数据量级上单站 70 年日值大约 2.5 万行2400 站就是 6000 万行起步如果算上小时数据或者多个要素轻松过亿。这个量级用 pandas 单机处理读入就要几十 GB 内存所以必须上分布式。2.2 为什么是 Spark 而不是别的选 Spark 的核心理由有三个。第一气象数据是典型的「宽表 时间序列」Spark SQL 的 DataFrame API 对这类结构化数据支持成熟groupBy、窗口函数、join 都能直接写不用手写 MapReduce。第二气象分析里大量是聚合操作比如按站按年求平均、按经纬度网格求极值这些正好是 Spark 的强项shuffle 之后并行聚合比单机快一个数量级。第三Python 生态能和 Spark 无缝衔接清洗用 PySpark最后小结果拉回 pandas 画图不用在两套语言之间倒腾。和 Databricks 的区别这里也说一句很多人搜到 Databricks 就懵。Databricks 是托管版 Spark 的商业平台底层还是 Spark本地或自建集群用开源 Spark 就行分析逻辑完全一致不用为了跑这个项目去上云。环境上本地开发用pip install pyspark就能起一个 local 模式集群则常见 standalone 或 YARN。下面给一个最小可跑的本地环境确认命令。# 确认 Python 和 PySpark 版本匹配Python 3.8 对应 Spark 3.x python --version pip install pyspark3.5.0 # 启动 local 模式4 个核验证环境 pyspark --master local[4]参数说明local[4]表示用本机 4 个线程模拟分布式开发阶段够用上集群时把 master 换成spark://host:7077或yarn。版本上 PySpark 和 Spark 必须对应3.5.0 的 PySpark 配 3.5.x 的 Spark混用会报序列化错误。2.3 从原始文本到 Parquet 的清洗链路原始文本直接分析效率极低标准做法是先清洗成 Parquet 列式存储后续查询快很多。清洗要处理四件事跳过文件头、按固定宽度或分隔符切列、单位换算、缺失值标记。气象数据里缺失值常用999999或32766表示必须替换成 null否则参与平均会把结果拉飞。from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType spark SparkSession.builder \ .appName(WeatherClean) \ .config(spark.sql.parquet.compression.codec, snappy) \ .getOrCreate() # 定义 schema避免 inferSchema 全表扫描 schema StructType([ StructField(station_id, StringType(), True), StructField(lat, DoubleType(), True), StructField(lon, DoubleType(), True), StructField(obs_date, StringType(), True), StructField(avg_temp, IntegerType(), True), # 单位 0.1 摄氏度 StructField(max_temp, IntegerType(), True), StructField(min_temp, IntegerType(), True), StructField(precip, IntegerType(), True), # 单位 0.1 毫米 ]) raw spark.read.option(delimiter, \t).schema(schema).csv(hdfs:///weather/raw/) cleaned raw \ .withColumn(avg_temp, F.when(F.col(avg_temp) 999999, None) .otherwise(F.col(avg_temp) / 10.0)) \ .withColumn(precip, F.when(F.col(precip) 999999, None) .otherwise(F.col(precip) / 10.0)) \ .withColumn(obs_date, F.to_date(F.col(obs_date), yyyyMMdd)) \ .filter(F.col(obs_date).isNotNull()) cleaned.write.mode(overwrite).partitionBy(station_id).parquet(hdfs:///weather/clean/)逻辑说明schema 显式声明能省掉一次全表推断大文件上差别明显。when...otherwise把缺失值转 null 并同时做单位换算一步到位。partitionBy(station_id)让后续按站查询只扫对应分区是气象数据最自然的切分维度。参数上Parquet 压缩用 snappy 是速度和体积的平衡点追求更小体积可以换 gzip但读写会慢。3. 用 Spark SQL 算全国气温与降水的核心指标3.1 逐年平均气温与距平计算清洗完的数据要回答的第一个问题是全国逐年平均气温怎么变。距平是气象里的标准说法指某年值减去多年平均值正距平偏暖负距平偏冷。计算分两步先按年聚合再算基准期均值最后相减。基准期一般取 1981 到 2010 年这是气象业务里的常用气候态。clean spark.read.parquet(hdfs:///weather/clean/) # 按年求全国平均气温先按站按年再平均避免大站权重过高 yearly clean.filter(F.col(avg_temp).isNotNull()) \ .groupBy(station_id, F.year(obs_date).alias(year)) \ .agg(F.avg(avg_temp).alias(station_year_temp)) \ .groupBy(year) \ .agg(F.avg(station_year_temp).alias(national_temp)) # 基准期 1981-2010 的气候态 baseline yearly.filter((F.col(year) 1981) (F.col(year) 2010)) \ .agg(F.avg(national_temp).alias(base_temp)) \ .collect()[0][base_temp] anomaly yearly.withColumn(anomaly, F.col(national_temp) - baseline) \ .orderBy(year) anomaly.show(80)逻辑说明这里有个容易忽略的细节全国平均不能直接把所有站点记录丢进一个 avg那样站点多的年份权重会偏。正确做法是先算每站每年均值再对站求平均保证每个站权重一致。参数上基准期可以按研究需要调整比如算 1991-2020改 filter 条件即可。collect()[0]取单值基准是安全的因为聚合后只有一行。3.2 降水极值与空间分布降水分析关心的是极值和分布不是简单平均。比如要找出每个站历史上单日最大降水量以及全国某年降水最多的前 10 个站。窗口函数在这里很好用。from pyspark.sql.window import Window # 每站历史单日最大降水 w Window.partitionBy(station_id).orderBy(F.col(precip).desc()) top_precip clean.filter(F.col(precip).isNotNull()) \ .withColumn(rk, F.row_number().over(w)) \ .filter(F.col(rk) 1) \ .select(station_id, obs_date, precip) # 按经纬度 1 度网格统计年均降水用于画空间分布 grid clean.filter(F.col(precip).isNotNull()) \ .withColumn(grid_lat, F.floor(lat)) \ .withColumn(grid_lon, F.floor(lon)) \ .groupBy(grid_lat, grid_lon, F.year(obs_date).alias(year)) \ .agg(F.sum(precip).alias(grid_year_precip)) grid.write.mode(overwrite).parquet(hdfs:///weather/grid_precip/)逻辑说明窗口函数按站分区、按降水降序排名取第一名就是历史极值比 groupBy 加 max 再 join 回去更简洁。网格统计用floor把经纬度归到 1 度格子是气象空间分析里最省事的做法精度够用且结果规整。参数上网格大小按需求改0.5 度更细但数据量翻四倍。3.3 结果导出与可视化衔接Spark 算完的结果通常不大拉回 pandas 画图最方便。但要注意别把大表直接 toPandas会 OOM。正确做法是先聚合到小结果再拉。import matplotlib.pyplot as plt pdf anomaly.select(year, anomaly).toPandas() pdf pdf.sort_values(year) plt.figure(figsize(12, 4)) plt.bar(pdf[year], pdf[anomaly], color[#c0392b if v 0 else #2980b9 for v in pdf[anomaly]]) plt.axhline(0, colorblack, linewidth0.8) plt.xlabel(Year) plt.ylabel(Temperature Anomaly (C)) plt.tight_layout() plt.savefig(national_temp_anomaly.png, dpi150)逻辑说明toPandas只在小结果上调用这里 70 行数据毫无压力。颜色按正负距平分色是气象距平图的标准画法红暖蓝冷。参数上 dpi 150 够论文和报告用再高文件会大。4. 避坑与排查气象数据跑 Spark 最容易翻车的五件事4.1 缺失值没处理平均值直接失真现象算出来的全国平均气温比实际高或低十几度或者某年数据明显异常。原因气象原始数据用 999999、32766 这类哨兵值表示缺测直接参与 avg 会把它当成真实值。解决清洗阶段就用when...otherwise把所有哨兵值转 null聚合前统一filter(col.isNotNull())。我一般会把每个要素的缺失值编码写进配置不硬编码在代码里。4.2 单位没换算结果差 10 倍现象平均气温显示 150 而不是 15。原因气温和降水在原始数据里是 0.1 单位存储忘了除以 10。解决清洗时统一换算并在 schema 注释里写清单位。这个坑的隐蔽之处在于如果只算距平正负趋势可能还对但绝对值全错所以一定要抽查几个已知站点的值对不对。4.3 分区过多导致小文件灾难现象清洗任务跑完HDFS 上出现几十万个小文件后续查询启动慢。原因partitionBy(station_id)加上默认 200 个 shuffle 分区2400 个站乘上分区数就是海量小文件。解决清洗阶段先repartition到合理数量或者按年加站号组合分区。我一般会把分区数控制在几千以内写完用spark.sql.files.maxPartitionBytes控制读取时的合并。4.4 窗口函数没加分区全表 shuffle现象算每站极值的任务卡在 shuffle 阶段不动。原因窗口定义忘了partitionBySpark 把全表当一个分区排序。解决窗口必须partitionBy(station_id)让每个站独立排序。这个错误在数据量小时看不出来上亿行时直接拖垮集群。4.5 本地模式内存不够误以为是代码问题现象local 模式跑全国数据报java.lang.OutOfMemoryError。原因local 模式默认内存小且所有分区挤在一个 JVM。解决开发阶段先用单省或单年数据验证逻辑上集群再跑全量本地要跑就调大 driver 内存--driver-memory 8g并减少 shuffle 分区数。别在本地硬扛全量数据那是给自己找不痛快。5. 把分析做成可复用的气象数据管道调度、增量与验证前面几步跑通的是单次分析真正要长期用得把它做成能重复跑的管道。核心是三件事调度、增量更新、结果验证。调度上气象数据通常按月或按年更新用 Airflow 或 cron 触发 Spark 任务都行。关键是把清洗和分析拆成两个独立任务清洗产物是 Parquet分析只读 Parquet这样重跑分析不用重新解析原始文本。下面是一个 cron 风格的提交脚本。#!/bin/bash # 每月 5 号跑上月数据清洗再跑分析 spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --num-executors 20 \ --conf spark.sql.shuffle.partitions400 \ clean_weather.py --month $(date -d last month %Y%m) spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --num-executors 10 \ analyze_weather.py参数说明spark.sql.shuffle.partitions默认 200数据量大时调到 400 到 800 能减少单分区压力executor 内存按集群规格调气象聚合是内存密集型4g 起步比较稳。增量更新上清洗任务按mode(append)写入并按年月分区分析任务只读需要的分区避免全量重算。结果验证是很多人省掉但最不该省的一步。我的习惯是每次跑完抽三个检查点全国站点数是否在合理范围、逐年记录数是否连续无断档、几个已知站点的极值是否和历史记录吻合。这三个检查用 SQL 几行就能写。# 验证每年记录数和站点数是否合理 check clean.groupBy(F.year(obs_date).alias(year)) \ .agg(F.count(*).alias(rows), F.countDistinct(station_id).alias(stations)) \ .orderBy(year) check.show(80) # 正常情况站点数稳定在 2000 以上行数随年份平滑变化突降说明数据缺失逻辑说明站点数突然掉到几百说明某年数据没导全行数某年翻倍可能是重复导入。这两个信号比任何日志都直观。参数上不用调直接看输出即可。最后说个我自己的习惯。气象数据分析最怕的不是算不出来是算出来了但不知道对不对。所以我现在每写一个聚合逻辑都会先用一个省的数据在 pandas 里手算一遍拿结果去对 Spark 的输出对上了再上全量。这个笨办法帮我挡掉过好几次单位换算和缺失值处理的错误。全国历史气象数据是个值得长期做的方向数据公开、问题真实、技术栈通用跑通一次之后能延伸出很多分析。希望帮到你。本文还有配套的精品资源点击获取
返回列表