尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Spark处理两百GB全国气象数据:从清洗到性能调优的完整实战

Spark处理两百GB全国气象数据:从清洗到性能调优的完整实战 简介大数据处理中传统单机Pandas面对数十GB乃至上百GB的结构化数据时往往因内存瓶颈而无法胜任。Spark作为分布式计算引擎通过内存计算、分区读取与Catalyst优化器为海量表格数据提供了高效的分析方案。在实际工程中数据清洗是保证分析质量的关键包括缺测值映射、异常值过滤及质量控制并将清洗结果转为Parquet列式存储以提升查询性能。同时合理设置分区数、使用广播变量优化JOIN、调整序列化方式等性能调优手段能显著提升Spark作业的稳定性与效率。本文以全国两千多个气象观测站、七十余年的气候记录为实例完整展示了从技术选型、数据预处理、核心统计到可视化交付的实战路径为处理大规模气象数据或类似多维表格数据的开发者提供一套可复用的参考方案。 去年接了个数据处理项目甲方扔过来好几块硬盘打开一看是二十多个省份、两千多个气象观测站、七十多年的逐日气候观测记录。压缩包全部解压完大概两百多GB我第一反应是这玩意儿用单机Pandas想都不用想光读一遍就能把内存干到爆。所以我直接用Spark搭了一套分析流程配合Python做图形化输出前后一个月把整个流程从数据清洗、核心统计、可视化到结果交付完整跑通中间踩了不少坑。这篇文章我想把整个过程从技术选型、数据清洗、核心分析代码到性能调优完整记录下来给准备用Spark处理大规模表格数据、或者正打算做气象/地理类数据分析的同学一份可以直接抄作业的参考。文章里涉及的都是我在实际项目里验证过的东西没有什么“最佳实践”式的空话基本都是线性的真实操作步骤。1. 为什么用Spark做气象数据分析——技术选型实录1.1 先算算全国历史气象数据到底有多大很多人对“大数据”这个词没有概念觉得几个GB的CSV也敢叫大数据实际上拿到全国历史气象数据之后你会发现这个量级远超单机内存能承受的上限。先做一个简单的估算全国大约有2400多个国家级气象观测站如果使用逐日数据从1951年算到2023年一共73年每年365天左右那数据量大约就是2400 × 73 × 365 ≈ 6400万行。但实际项目里拿到的数据往往不止国家级站点还有区域级自动站站点数量能到几万个时间粒度也可能细化到逐小时。如果是逐小时数据那行数直接乘以24大约15亿行。再加上温度、降水、风速、湿度、气压、日照等十几个要素列每条记录按100字节算逐日数据大约6~7GB逐小时数据超过150GB。这还只是原始文本。如果换成未经压缩的文本格式加上各种分隔符、表头、换行符整体体积还要膨胀30%~50%。我自己拿到的那批数据解压后是两百多GB里面覆盖了逐小时和逐日两种粒度加了质量控制标记字段所以比估算值还大一些。这种量级的数据单机Pandas不是“慢”的问题而是读不读得进内存的问题。64GB内存的机器一次性加载两百多GB原始文件直接OOM没有任何悬念。所以第一步决策就是把处理引擎换成Spark。1.2 Pandas、Dask、Spark我为什么选了Spark在处理大数据表格时通常有三条路Pandas、Dask、Spark。Pandas单表哪怕只有一亿行在普通服务器上也已经很难玩得转了groupby、merge这类操作一旦触发中间结果膨胀内存会瞬间翻倍。Dask可以在单机上做并行计算也能处理超过内存的数据但它的生态和Spark相比还是薄弱一些尤其是SQL优化、JOIN优化、窗口函数方面性能差距明显。Spark的DataFrame API本身是模仿Pandas的Python端的程序员上手几乎没有学习成本同时它背后有Catalyst查询优化器能对逻辑计划做谓词下推、列剪枝、分区裁剪这部分优化是Dask目前比不了的。另外一个关键原因是项目后期需要做全国维度的聚合分析比如按年份聚合全国平均气温、按站点统计极端高温天数、按月份算降水量分布这类操作本质上都是分组聚合和窗口计算正是Spark最擅长的场景。还有一点纯粹从数据分析角度来说Spark的SQL能力非常成熟。我可以直接把DataFrame注册成临时表用SQL写聚合逻辑这比用Pandas链式写法直观很多尤其在处理多重group by和窗口函数时SQL的可读性优势非常明显。1.3 Spark在这个项目里具体解决了什么问题我总结一下Spark在这个气象项目里实际发挥的作用主要有四个方面第一分布式存储和读取。数据源放在HDFS或者本地文件系统Spark按分区读取不会把所有数据加载到一个节点上。第二内存计算与迭代优化。气象分析经常会做多阶段聚合比如先按站算月均值再按月算年均值最后在全国层面做趋势拟合。Spark的RDD和DataFrame会把中间结果缓存在内存里避免每次聚合都重新读一遍磁盘。第三容错机制。跑一个几十亿行的任务中间任何一个节点挂了在MapReduce里可能整个Job重跑Spark基于Lineage血缘关系只重算丢失的分区这在长时间运行的任务里太重要了。第四生态衔接。Spark计算结果可以非常方便地转成Pandas DataFrame也可以直接写Parquet、CSV、数据库。我最后就是用toPandas把聚合完成的小结果集拿回本地画图完全够用。2. 数据清洗与预处理把一堆原始文件变成可用数据2.1 原始气象数据长什么样气象数据的原始格式没有统一标准常见的有两种来源。一种是国内气象数据网下载的地面气候资料日值数据集一般是一个站点一个文件文件名带区站号内容是固定的列格式每天一行包含日期、平均气温、最高最低气温、降水量、风速、日照时数等。另一种是NOAA的GSOD/ISD数据按年打包全球站点都在一起带有经纬度和质量控制字段。两种格式都得先做解析。国内数据集最坑的地方在于缺测值往往用特殊数字表示不同版本的数据集缺测值还不一样。我常见到的就有32700、32744、32766、99999等这些都是历史数据里约定俗成的缺测编码没有统一的规范文档只能对着数据说明手工核对。还有一个常见问题是站点迁站。同一个区站号对应的经纬度在不同年份可能不一样如果不处理后续做地理可视化时会出现站点漂移的诡异现象。所以我把清洗流程设计成三个阶段解析阶段把不同来源的文件统一读成DataFrame列名规范成英文小写加下划线日期解析成标准timestamp类型。过滤阶段缺测编码全部映射成NULL质量控制码非0的记录直接过滤异常物理值剔除。落地阶段清洗好的数据重新分区写成Parquet列式存储后续所有分析都从Parquet读取。2.2 清洗逻辑和关键编码规则清洗过程看起来简单但细节决定成败我把几个关键点单独拿出来说。第一个是缺测值映射。国内日值数据集的缺测值通常有多个例如32744表示缺报或未观测32766表示降水微量或者没有降水在降水要素里还分0和32700两种语义。这些值如果直接参与聚合运算会把平均气温拉低十几度结果完全失真。我的处理办法是写一个UDF把这些特殊值统一替换成NULL之后再聚合时用avg函数自然跳过NULL。第二个是异常值过滤。气温数据的物理合理范围大致在-55℃到50℃之间降水的物理合理范围是0到400毫米日降水量风速一般不超过75米/秒。超过这些范围的值要么是仪器故障要么是录入错误直接删掉即可。第三个是质量控制码。部分数据集每行都有QC字段0表示数据正确1表示可疑2表示错误3表示已修改。我保留了0和3过滤掉1和2这个策略可以最大程度保留有效数据的同时把噪声降到最低。第四个是站点信息表。我会单独建一张站点维度表包含站点ID、站名、省份、经度、纬度、海拔并处理迁站情况——如果一个站点的经纬度在不同时间段变化很大我会保留最新的经纬度并在分析维度上把它视为同一个站点。2.3 从原始CSV到列式存储Parquet清洗完的数据不会直接在原始格式上反复读取而是统一转换成Parquet格式落地。为什么选Parquet三个原因列式存储、高压缩比、内置统计信息。气象分析经常会用时间和站点做过滤比如只看2020年以后的数据或者只看某一省份的数据。Parquet的谓词下推再加上分区裁剪可以让Spark在读取阶段就跳过大量无关文件实际测试下来查询速度比直接读CSV快三到五倍。当时我用Spark读取清洗后的数据按“年份”字段做了动态分区写入每个年份一个分区目录。这样后续分析例如算年度趋势时Spark只需要扫描对应年份的分区不需要全量扫描整个表。写到Parquet时的编码也很关键因为原始数据里如果有中文列名Parquet schema的兼容性会有问题所以我在清洗阶段统一把列名改成了英文。时间字段建议存储成timestamp类型不要存字符串这样后续过滤效率高很多。3. 核心分析代码与实现细节3.1 气温趋势分析年度聚合与线性拟合分析任务里最核心的一个指标是“全国年平均气温变化趋势”。实现思路不复杂从日值数据里按年份聚合出年平均气温再用线性回归看斜率。第一步是读取Parquet数据from pyspark.sql import SparkSession from pyspark.sql.functions import col, year, avg, count spark SparkSession.builder \ .appName(weather_analysis) \ .config(spark.sql.shuffle.partitions, 200) \ .config(spark.sql.parquet.enableVectorizedReader, true) \ .getOrCreate() df spark.read.parquet(hdfs:///user/weather/clean_data/)第二步按年份聚合全国平均气温。这里有个细节不同站点的日值缺失情况不一样如果直接对所有站点做简单平均某些年份缺测站数特别多平均结果的代表性会变差。所以我一般会要求有效站点数必须超过全国站点数的三分之二才认为该年均值有效。yearly_temp df.filter(col(avg_temp).isNotNull()) \ .groupBy(year(date).alias(year)) \ .agg(avg(avg_temp).alias(nation_avg_temp), countDistinct(station_id).alias(station_count)) yearly_temp yearly_temp.filter(col(station_count) 1600) \ .orderBy(year)第三步把聚合结果拿到本地做线性拟合。因为到这里数据量已经非常小了只有几十行用toPandas拉回本地完全没问题。import pandas as pd import numpy as np pdf yearly_temp.toPandas() x pdf[year].values y pdf[nation_avg_temp].values slope, intercept np.polyfit(x, y, 1)算出来的斜率就是每年平均升温多少摄氏度。这个数字如果要用于正式报告最好再算一下置信区间我用的是scipy的linregress函数可以同时给出p值和标准误。3.2 极端天气事件统计高温日数和暴雨日数除了平均气温极端天气的频次变化也是气象分析里非常关心的指标。这里我做了两个典型的统计高温日数和暴雨日数。高温日数的业务定义是日最高气温≥35℃记为一个高温日≥40℃记为强高温日。我按年、按站点把高温日数统计出来再汇总成全国或者分省的总量。hot_days df.filter(col(max_temp) 35) \ .groupBy(year(date).alias(year), col(province), col(station_id)) \ .agg(count(*).alias(hot_day_count))暴雨日数的定义是日降水量≥50mm记为一个暴雨日≥100mm记为大暴雨日。rain_storm df.filter(col(precipitation) 50) \ .groupBy(year(date).alias(year), col(station_id)) \ .agg(count(*).alias(storm_day_count))这类统计的坑在于降水要素和气温要素的缺测判断方式不一样。降水如果为0表示无降水是正常观测值而气温的0度是正常值不能过滤。所以清洗时不能一刀切把所有0值都删掉必须按要素区别处理。再往深一层做可以统计高温日数随年份的变化趋势。做法和前面气温趋势类似把每年全国高温总日数或者平均单站高温日数做线性回归看斜率为正还是负。这部分结果通常比单一平均气温更有说服力因为极端天气的频次变化对公众感知更直观。3.3 站点维度的精细化分析窗口函数与滞后计算如果要看更精细的规律比如单站点的气温变化率、季节转换时间的变化就需要用到窗口函数。比如我想计算每个站点的年平均气温在时间上的变化可以用lag函数计算相邻年份之间的温差from pyspark.sql.window import Window w Window.partitionBy(station_id).orderBy(year) station_yearly df.groupBy(station_id, year(date).alias(year)) \ .agg(avg(avg_temp).alias(year_temp)) station_yearly station_yearly.withColumn( prev_year_temp, lag(year_temp).over(w) ).withColumn( temp_change, col(year_temp) - col(prev_year_temp) )这样每个站点每年相对于上一年的温差就出来了。汇总统计时如果把temp_change按年做平均可以看出整体升温在哪些年份最明显哪几年有降温波动。窗口函数的坑是如果不指定partitionBy默认按全表排序那会把所有数据放到一个分区里算压力极大。气象数据分析时分区键一定是站点或者年份这种分布均匀的列而不是温度这种连续值。再举个例子30年气候平均值的计算也适合用窗口函数实现。所谓气候平均一般取最近30年如1991-2020年平均。用窗口函数按站点和月份分组滑动30年窗口算平均比先过滤再聚合高效很多因为不需要反复扫描整个数据表。4. 可视化落地与业务洞察4.1 什么时候把Spark数据拉回本地很多初学者拿到Spark就喜欢把数据collect回Pandas再处理这个思路在数据量大时是致命的。Spark计算结果分为两种一种是聚合后的结果比如年均温度每年全国就一行这种数据量很小用toPandas拉回本地画图完全没问题另一种是全量明细数据比如所有站点每天的原始记录几百GB这绝不能拉回本地。我的原则是所有分析都在Spark里完成聚合最终拉回本地的必须已经是“结果表”。比如画全国气温趋势图我只需要年份和对应平均温度两列最多几十行画站点分布图我需要每个站点的经纬度和某种度量值大约几千行。几千行的DataFrame用toPandas一点压力都没有。画图部分我用的是Matplotlib加Seaborn地理分布图用Cartopy画中国地图边界然后再叠加散点图。4.2 趋势图、热力图、分布图怎么画整个项目里我生成了三类核心图表。趋势图是最直观的。把年份作为x轴平均气温作为y轴画折线图再叠加一条线性回归直线能一眼看出升温趋势。这里有一个视觉优化细节折线图的线宽调到2.5左右回归线用虚线加粗颜色错开不然两种线容易混在一起。热力图适合看月份-年份维度的气温变化。把数据透视成年份×月份的矩阵用Seaborn的heatmap画暖色表示高温冷色表示低温。不过从视觉可读性出发我建议还是按年月折叠成趋势曲线更直观热力图只适合做展示大图不适合做核心结论。分布图用于看地理特征。把每个站点的经纬度和指标值比如年均温变化斜率组合起来在地图上画散点点的大小和颜色映射指标数值。这里必须注意地图投影坐标系要和散点经纬度坐标系一致我使用的是EPSG:4326即WGS84经纬度坐标系Cartopy的ccrs.PlateCarree就是匹配这个的。4.3 从数据里看到的几个有意思的规律分析结束之后我总结了一些有价值的业务发现这里挑三个有代表性的简单说一下。第一个是增温速率的空间差异。从单个站点的气温变化斜率来看我国北方大部分站点年均温升高速率高于南方具体数字这里不展开但从业务角度看这意味着不同区域的气候适应策略应该差异化而不是一刀切。第二个是极端高温日数的年际波动明显增大。高温日数不是线性增长的而是呈现所谓“跳跃式增长”某些年份突然暴增之后维持在较高水平。用滑动平均处理后会看得更明显。第三个是暴雨事件的不均匀性。暴雨日数在南方站点显著高于北方但近年来的趋势是北方部分站点的暴雨日数也在增加这种区域性差异从聚合统计里很容易发现但从单一站点看很难看清全局。这些规律是整个项目的核心交付物也是Spark这种全量分析工具的价值体现——不用抽样直接跑全量数据结论的置信度完全不一样。5. 性能调优与集群部署踩坑记录5.1 分区数和并行度怎么设Spark的性能调优第一刀永远是看分区数和并行度。默认的spark.sql.shuffle.partitions是200这个值对几十亿行的数据明显不够。我一般会根据executor数量和CPU核心数来估算分区数设为executor总数乘以核心数再乘以2到3倍。举个例子我用的集群是8台机器每台8核那总核心数是64。目标并行度我设置在128到192之间。设置方式spark SparkSession.builder \ .config(spark.sql.shuffle.partitions, 160) \ .config(spark.default.parallelism, 64) \ .getOrCreate()还有一个细节读取Parquet时每个文件最好在128MB左右太小的文件会有大量task空转太大的文件则容易让单个task处理时间过长。写Parquet时如果发现小文件过多可以先repartition或coalesce再写。5.2 广播变量用得好JOIN速度翻倍气象分析经常要把大表和站点信息表做关联。站点信息表全国也就两千多行这种“大表join小表”的场景直接使用broadcast join会产生极大性能收益。默认情况下Spark的join如果遇到一张小表会自动广播但保险起见我会显式指定from pyspark.sql.functions import broadcast df df.join(broadcast(station_info), station_id, left)广播之后小表会被复制到每个executor的内存里整个join过程不需要shuffle性能提升立竿见影。这里需要注意广播变量的内存限制。默认的广播阈值是10MB如果小表超过这个值要调大参数.config(spark.sql.autoBroadcastJoinThreshold, 104857600)5.3 Shuffle与序列化配置经验Shuffle是Spark作业里最昂贵也最容易出问题的环节尤其是groupBy和join操作。我调优时主要关注几点。第一序列化方式。Java自带的序列化器性能比较差我换成Kryo.config(spark.serializer, org.apache.spark.serializer.KryoSerializer)Kryo序列化后的数据更紧凑网络传输和磁盘溢写的压力都会小很多实测性能提升20%以上。第二Shuffle落盘。如果executor内存不足Spark会把shuffle数据溢写到磁盘。这时要关注spark.shuffle.spill.compress配置默认是true一般保持默认即可。如果发现溢写严重说明executor内存不够或者分区数太少优化方向是增加shuffle.partitions。第三Executor的内存分配。每个executor分配的内存不宜过高太高会导致GC压力增大。我试过4核8G的executor配置比较稳定。堆内内存和堆外内存的比例用默认值就行不需要特别调。第四缓存策略。同一个DataFrame如果会被多次使用用cache或persist缓存在内存里。气象数据聚合之后经常会被多个分析模块复用比如“日值数据”这个表我先做了一次清洗缓存再分别计算趋势、极端天气、月度统计节省了大量重复IO。5.4 本地调试与集群运行的差异本地跑Spark和集群跑Spark体验差异非常大很多坑都是本地不出问题一上集群就挂。本地模式下Spark默认把driver放在当前进程所有任务都在本地线程池执行executor内存默认是1GB。数据量大一点就会报OOM。所以本地调试一定要限制数据量我通常是先filter某一年或者某个省份的小样例数据跑通逻辑再在集群上全量铺开。还有一个常见问题是Python版本和依赖不一致。Spark的PySpark执行Python代码时会通过pyspark.daemon在executor上启Python进程如果每台worker的Python环境缺少依赖包任务会在运行中途失败。解决办法是在每台机器上装好相同的依赖或者用spark.pyspark.python指定一个统一的Python解释器。6. 常见问题速查与排障实录6.1 高频问题排查表我把这个项目里遇到的高频问题整理成一张表后面做类似项目可以直接对照排查。问题现象常见原因解决方案Executor OOMexecutor内存过小或shuffle数据量过大减少executor核数增加executor内存调大shuffle.partitions任务卡在某个stage长时间不动数据倾斜某个key数据量特别大给key加盐打散或者改用广播join读取CSV出现乱码源文件是GBK编码未指定读取时加option(encoding, gbk)时间字段解析失败字符串格式不一致用to_timestamp(col, yyyy-MM-dd)指定格式写Parquet后产生海量小文件shuffle.partitions设置过大写之前coalesce或者降低shuffle.partitionsjoin速度极慢大表join大表没有优化检查是否可以先过滤再join或者使用bucket join本地跑得通集群上挂Python环境不一致统一每台worker的Python环境和依赖6.2 一次典型的OOM排查过程我印象最深的一次OOM排查发生在算全国极端降水日数的时候。那次任务跑了一个多小时在某个stage反复失败错误信息显示executor lost每个executor都报OutOfMemoryError。我的排查步骤是先看Spark UI里的executor内存使用曲线发现shuffle read数据量异常大说明shuffle.partitions设置不够导致每个task拉取的数据过多。我先把shuffle.partitions从200调到500问题缓解了一些但还有两个executor在崩溃。再进一步看数据发现降水数据的groupBy key是“省份年份”但部分省份的站点数量远多于其他省份比如站点特别多的省份那个task要处理的数据量是其他省份的几倍。这属于典型的数据倾斜。最终解决方案是用两阶段聚合第一阶段先按“站点年份”聚合出每个站点的暴雨日数第二阶段再按省份汇总。这样第一阶段的shuffle key是站点ID分布均匀第二阶段的中间结果已经大幅减小不会再产生倾斜。这个案例的启示是遇到倾斜问题优先从业务逻辑层面拆解先把细粒度的中间结果聚合出来再用上层维度汇总往往比直接在最终维度上做groupBy要快一个数量级。6.3 运行环境与资源规划避坑建议最后补充几个关于运行环境的建议。我这次用的是8台机器的小集群每台机器64GB内存、16核CPU。Spark部署用的是Standalone模式没有上YARN主要是因为集群私有不需要多租户资源隔离。如果你手头没有集群也可以在本机用Spark的local模式跑这个项目不过需要把数据量缩小到一个合理的范围比如先只处理某几年或者某几个省的数据。想跑全量的话建议把clean数据做成分区Parquet然后至少准备一台32GB内存以上的机器local[*]模式也值得一试。另外一个建议是在大规模作业跑之前先用小的数据集跑一遍逻辑验证再上全量。这里面最基本的做法是用limit(1000)或者year过滤来做开发测试不要一上来就全表扫描。开发环境和生产环境的分离我吃过不少亏这一条值很多时间成本。写在最后的一点体会这个项目跑完之后我最大的感受是Spark真正解决的是“单机装不下、跑不动”的硬性诉求而不是为了用框架而用框架。实际上最后交付的分析报告核心图表也就十几张但每一次聚合计算都建立在几千万甚至几十亿行的全量数据之上这种结论的说服力是完全不一样的。如果你手头正好有一批类似的多维表格数据我的建议是先在本地用一个小样例跑通逻辑再放到Spark上铺开。清洗环节宁可多花点时间把缺测值和质量控制搞清楚也不要急着做分析因为气象分析这种领域数据本身的准确度比分析技巧重要得多。踩过几次坑之后你就会明白数据清洗阶段每多花一小时后面分析阶段就能省下十几个小时。本文还有配套的精品资源点击获取
返回列表