
刚接到“Spark第一次作业”的时候我心里其实没有底。课上讲了RDD、DataFrame、Spark SQL这些概念听着像那么回事可真要上手写一个能跑的作业还是另一码事。我们这次作业的任务是围绕一份农产品价格数据做清洗和分析要求用Spark完成数据的读取、过滤、聚合和结果输出并且要提交脚本和运行记录。整个流程走下来我最大的感受是Spark的学习门槛不在于API本身而在于环境、数据、资源三者之间的协调。这篇文章就从头到尾记录我当时从拆解需求、搭建环境、写代码到排查问题、验证结果的完整过程希望能给正卡在第一次Spark作业上的同学一些直接能用的参考。1. 拿到作业后的三个决定1.1 这个作业到底要做什么先把需求彻底搞清楚。作业给的数据是一批农产品市场报价记录大概长这样农产品的名称、所属地区、上报日期、最低价、最高价、平均价偶尔还有单位。任务要求分三步走一是把原始数据读进来并且清洗干净包括去掉空值、去掉重复记录、把价格字段统一成数值类型二是做若干统计分析比如按月计算各品种的平均价格、找出每个地区价格波动最大的农产品三是把结果导出成可读的文件并给出关键结论。我一开始犯了一个典型错误就是拿到需求就想着赶紧写代码。后来发现如果不先想清楚“输入长什么样、中间要做什么变换、输出给谁看”这三个问题代码写一半就得推倒重来。所以建议所有第一次做Spark作业的人先花30分钟把数据文件打开看几行把字段名、格式、脏数据长什么样都摸清楚再开始写第一行代码。这一步省下来的时间远比你想象的多。1.2 为什么选Spark而不是其他工具其实一开始我也有过疑问数据量也不算特别大用Pandas处理不就行了为什么非要用Spark后来我理解了一个关键点Spark的核心价值不只是“处理大数据”而是在于它有一套统一的分布式计算模型和内存计算引擎。哪怕数据量还没到TB级别Spark的DataFrame API配合Catalyst优化器也能自动帮我们做谓词下推、列剪枝这些优化写出来的代码结构比Pandas脚本更清晰也更适合迁移到生产环境。用MapReduce做对比就更明显了。MapReduce每一步都得落盘反复读写磁盘在一个多次清洗和多阶段聚合的作业里I/O开销非常大。Spark把中间结果尽量放在内存里迭代计算效率高一个数量级。当然选型不是越重越好如果只是单机小数据处理用Pandas确实更快也更省事。但既然作业要求是“用Spark完成”那就该把Spark的思维方式学透先定义DataFrame再通过变换得到结果而不是用循环逐行处理。1.3 学习路径怎么排不贪多先跑通我给自己定了一个原则先把Local模式跑通再考虑集群先把DataFrame API跑通再琢磨RDD级别的东西先能出结果再谈性能优化。第一次做作业目标不是成为Spark专家而是把“数据进、结果出”这条链路完整走一遍。具体步骤我安排成四个阶段环境搭建安装Spark并跑通一个最简单的测试、API学习理解RDD、DataFrame、Spark SQL的区别和使用场景、作业实现完成清洗和分析、排查验证跑通后核对结果并处理各类报错。这样安排的好处是每阶段都有明确产出不会一头扎进源码或底层原理里出不来。那些源码分析、内存模型调优等作业写完、踩过坑之后再去补反而理解得更快。2. 环境准备与集群搭建跑通第一个Spark程序2.1 本机实验还是搭集群作业数据量不大完全没必要一开始就搭三台机器的集群本地模式完全够用。我在本地开发机上用Spark的Local模式跑Master地址直接写local[4]意思是本地用4个线程来模拟并行执行。这样既能验证代码逻辑又不用考虑机器之间如何通信问题排查也简单得多。网上有很多关于Spark集群搭建的教程热门搜索词里也有“spark集群搭建”但我建议第一次做作业的同学冷静一点集群搭建是一个独立的学习项目它涉及节点规划、SSH免密、启动Master和Worker等一整套流程和你写数据分析作业其实是两码事。先把作业本身跑通有精力了再回头研究集群这样两条线都不会太受影响。如果作业确实要求集群环境那就老老实实开三台虚拟机按官方文档一步步来别跳步。2.2 Spark版本与JDK兼容性安装Spark之前必须先确认JDK版本。我用的是Spark 3.5.x官方要求JDK 8/11/17都可以但实测下来JDK 11最稳。如果你装了JDK 21某些老版本Spark反而会报模块访问相关的错误排查起来很头疼。Python这边PySpark 3.5版本支持Python 3.8到3.11用3.9或3.10比较保险别上来就上Python 3.12配套的依赖包可能还没有跟上。安装过程其实不复杂核心几步是这样# 1. 下载Spark二进制包选带hadoop的版本 wget https://dlcdn.apache.org/spark/spark-3.5.4/spark-3.5.4-bin-hadoop3.tgz # 2. 解压并放到统一目录 tar -zxvf spark-3.5.4-bin-hadoop3.tgz sudo mv spark-3.5.4-bin-hadoop3 /opt/spark # 3. 配置环境变量 export SPARK_HOME/opt/spark export PATH$PATH:$SPARK_HOME/bin:$SPARK_HOME/sbin # 4. 验证安装 spark-shell --version2.3 内存参数设置的基本逻辑跑作业之前需要了解一下Driver和Executor的关系。简单说Driver是“总指挥”负责接收代码、分解任务、汇总结果Executor是“干活的工人”负责真正执行计算和数据存储。所以Driver的内存不能太小否则代码一复杂、结果一多就容易OOMExecutor的内存决定了每个计算节点能处理多少数据。我当时用本地模式跑作业配置是这样spark-submit \ --master local[4] \ --driver-memory 2g \ --executor-memory 4g \ --num-executors 2 \ spark_homework.py注意在Local模式下--executor-memory和--num-executors不会真正启动两个独立Executor进程但参数逻辑还是要理解。为什么Driver内存给2g因为作业里会反复缓存DataFrame中间结果都压在Driver的内存数据结构里为什么本地线程数用4因为本机CPU是8核Spark建议Local模式下线程数不超过CPU核数-2留出系统余量。2.4 第一个Spark程序的启动与验证环境装好后我写了一个极简的PySpark脚本来做冒烟测试确保环境真的没问题from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(FirstSparkJob) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate() df spark.createDataFrame([(apple, 5.0), (banana, 3.0)], [item, price]) df.show() spark.stop()用spark-submit执行这个脚本如果控制台能正常打印出一个表格就说明Spark环境基本没问题。这里有一个小细节spark.sql.shuffle.partitions默认值是200这是为了适配集群环境下的大规模shuffle设计的。但本地小数据作业如果保持200会导致每个分区里几乎没有数据白白浪费调度开销。我把这个参数调成4让每个分区里有实际内容作业跑起来明显更快。3. 作业核心RDD、DataFrame与Spark SQL怎么选3.1 三个API的定位与区别有一次面试官问过我一个问题RDD和DataFrame到底什么区别我当时答得稀碎。后来自己想通了RDD是把数据当成一个个对象你告诉Spark“怎么算”比如map一个函数处理每个元素DataFrame是把数据当成一张表你告诉Spark“算什么”比如筛选price 10的记录至于底层怎么扫描、怎么优化Catalyst优化器自动帮你做。这里可以用一个生活化的类比。RDD就像你逛菜市场每个摊位都要自己去问价、自己挑货、自己装袋DataFrame就像是写了一张购物清单“买3斤西红柿单价不超过5块的”店员拿着清单自己规划去哪里拿货最划算。RDD给足了你自由度但一切优化都要自己管DataFrame牺牲了一点灵活性换来的是自动优化的执行计划。在我这个作业里数据处理涉及大量结构化的字段过滤、聚合、连接用DataFrame和Spark SQL明显比RDD更高效。一方面代码更短逻辑一目了然另一方面Spark能自动做的优化比如只读需要的列、尽早过滤数据这些靠RDD手写的话得写很多样板代码。3.2 为什么最终选择DataFrame加Spark SQL最初我有三个备选方案纯RDD算子、DataFrame API、Spark SQL写HQL风格的语句。最终的选择是“DataFrame API为主Spark SQL为辅”。原因是作业里清洗逻辑复杂DataFrame的withColumn、filter、dropDuplicates这些方法链式调用写起来非常直观。某些统计查询比如按月分组求平均值、计算环比变化用Spark SQL写更简练代码量可以少一半。两种方式底层共用同一套逻辑计划也就是说你不用担心性能有什么差异Spark会把两者翻译成同一种执行计划。一个小技巧是把DataFrame注册成临时视图这样就能用SQL语句直接查询它。我经常先用DataFrame做清洗然后df.createOrReplaceTempView(price_daily)再写SQL做分析。这种混合写法非常适合“清洗用API、分析用SQL”的作业场景。3.3 读取JSON数据时最容易忽略的细节热搜关键词里有一条是“spark中读取json”这确实是作业里的高频操作。Spark读取JSON很方便但有几个坑值得提前知道第一个坑是schema推断。如果数据里有缺失值Spark推断出来的字段类型经常和预期不一致。比如某一行价格字段是字符串“12.5元”其他行是数字Spark可能把整列推断成StringType后面做聚合时就傻眼了。解决办法是手动指定schema不要再让Spark猜。from pyspark.sql.types import (StructType, StructField, StringType, DoubleType, DateType) schema StructType([ StructField(product, StringType(), True), StructField(region, StringType(), True), StructField(date, DateType(), True), StructField(min_price, DoubleType(), True), StructField(max_price, DoubleType(), True), StructField(avg_price, DoubleType(), True), StructField(unit, StringType(), True), ]) df_raw spark.read.schema(schema).json(data/agriculture_prices.json) df_raw.printSchema()第二个坑是嵌套JSON。如果原始数据是类似{info: {price: 12.5, region: 北方}}的结构直接读进来后整列都是StructType。这时候要用col(info.price)的方式取内层字段或者用from_json函数把字符串字段先解析成结构化字段。第三个坑是多行JSON的格式。Spark的read.json支持的标准JSON格式是“每行一个JSON对象”也就是JSON Lines格式。如果你拿到的文件是每行一个JSON但中间带换行缩进的“漂亮格式”读取时很容易解析失败。我当时就遇到这个问题后来用了一个小脚本把多行JSON转成JSON Lines格式再用spark.read.json读取就正常了。4. 第一次作业的完整实现从清洗到分析4.1 数据准备与字段设计我们的数据源是从文件系统读入的农产品价格记录约几万行字段有产品名称product、地区region、日期date、最低价min_price、最高价max_price、平均价avg_price、单位unit。业务目标有三个求各月份全国平均价、找出每个地区价格波动最大的产品、按产品统计月均价变化趋势。先把数据加载进来做初步探查。我用df_raw.show(10, truncateFalse)看一眼原始数据再用df_raw.groupBy(region).count().show()看每个地区有多少条记录df_raw.describe(min_price, max_price, avg_price).show()看数值分布。这一步很重要你能直观看到哪些字段有大量空值、哪些价格明显不合理、哪些字段格式需要调整。数据探查做完我心里基本有数了。4.2 清洗逻辑处理空值、去重、格式统一清洗步骤我拆成了四步每步都用DataFrame的API实现并中间打印行数来验证清洗效果。from pyspark.sql.functions import col, regexp_replace, trim, lower, when # 第1步去掉完全重复的记录 df_dedup df_raw.dropDuplicates([product, region, date]) # 第2步统一文本格式 df_clean df_dedup.withColumn(product, lower(trim(col(product)))) \ .withColumn(region, lower(trim(col(region)))) # 第3步价格字段清洗去掉单位后缀 df_clean df_clean.withColumn( avg_price_numeric, regexp_replace(col(avg_price), [^0-9.], ).cast(double) ) # 第4步过滤异常值 df_clean df_clean.filter( (col(min_price).isNotNull()) (col(max_price).isNotNull()) (col(avg_price_numeric) 0) (col(max_price) col(min_price)) ) print(原始行数:, df_raw.count()) print(去重后行数:, df_dedup.count()) print(清洗后行数:, df_clean.count())这里每一步的逻辑都不是随便写的。去掉重复记录是因为同一品种同一天在不同渠道上报了多次业务上只需要一条价格字段用正则把非数字字符去掉是因为原始数据里价格带“元/斤”这种单位后缀异常值过滤是为了把max_price min_price这样的逻辑错误数据排除掉。有一个细节很多教程不会提文本统一大小写这一步极其重要。原始数据里“苹果”“Apple”“苹果 ”三种写法都有如果不提前统一分组聚合时会各自成组统计结果就会碎掉。我做完这一步分组数从几十个降到了正常水平的个位数可见数据品质对后续分析的影响有多大。4.3 分析逻辑月均价格与波动排行清洗完数据后我先把清洗后的DataFrame注册成临时视图然后用SQL做分析。第一个分析是各月份各产品的平均价格df_clean.createOrReplaceTempView(price_clean) df_monthly spark.sql( SELECT product, substr(date, 1, 7) AS month, ROUND(AVG(avg_price_numeric), 2) AS avg_price FROM price_clean GROUP BY product, substr(date, 1, 7) ORDER BY month, avg_price DESC ) df_monthly.show(20)这段SQL看着简单但有一个SQL基础不牢容易犯的错GROUP BY里用了表达式substr(date, 1, 7)那么SELECT里也必须用同一个表达式不能直接写month。我一开始想当然地写了SELECT product, month, ... GROUP BY product, month结果Spark直接报错提示month不在GROUP BY里。后来改成在SELECT里也用substr(date, 1, 7) AS month并且GROUP BY使用同一个表达式代码才跑通。第二个分析是找出每个地区价格波动最大的农产品。这里的“波动”我定义为某产品在该地区月均价格的标准差标准差越大说明价格起伏越明显。用窗口函数可以优雅地完成from pyspark.sql.window import Window from pyspark.sql.functions import stddev, row_number df_volatility spark.sql( SELECT region, product, ROUND(STDDEV(avg_price_numeric), 3) AS price_volatility FROM price_clean GROUP BY region, product ) windowSpec Window.partitionBy(region).orderBy(col(price_volatility).desc()) df_top_volatility df_volatility.withColumn(rank, row_number().over(windowSpec)) \ .filter(col(rank) 1) \ .select(region, product, price_volatility) df_top_volatility.show(20, truncateFalse)窗口函数的思路是先在region维度上分组然后组内按波动幅度降序排列用row_number()赋一个序号最后只保留每组序号为1的那一条。这里有一个效率上的小经验把STDDEV这种聚合在SQL里先做再在外面接窗口函数不要让窗口函数直接作用在原始数据上否则shuffle的数据量会大很多。4.4 结果导出与验证Spark作业一般要把结果写到文件里通常用CSV或Parquet格式。我第一次写的时候直接指定了输出路径结果跑完发现目录下有很多part-0000x文件。这是分布式计算的正常现象每个executor各写各的分片。实际操作中如果需要汇总成一个文件可以把输出重新分区df_result spark.read.csv(output, headerTrue).repartition(1) df_result.write.mode(overwrite).option(header, True).csv(output_single)结果是output_single目录下只有一个主要数据文件其他的是_SUCCESS标记文件和临时文件这在提交作业时看着就干净多了。验证结果这一步我做得比较仔细。我用两组数据对比一组是Spark作业输出一组是直接用Python脚本统计的原始数据。两条链路各自独立计算如果两个数字完全对得上基本上可以确认作业没有逻辑硬伤。当时月均苹果价格的计算结果和Python脚本自测结果相差0.01元原因是一边用了四舍五入、一边用了截断取整调整统一后结果完全吻合。5. 踩坑记录与问题排查5.1 Executor内存溢出的处理思路第一次跑作业我最怕的就是java.lang.OutOfMemoryError但偏偏就是遇上了。排查的过程很有代表性。当时我做了一个大数据量的JOIN和GROUP BY结果Executor直接OOM。排查顺序是先看日志确认是Driver内存问题还是Executor内存问题。如果是Executor内存溢出优先做这几件事一是检查代码里有没有不必要的collect()把全部数据拉到Driver二是看shuffle分区数是不是太多导致每个分区太小、排序列过大三是检查有没有及时释放缓存df.cache()之后必须用df.unpersist()释放。还有一个容易被忽略的点是堆外内存。Spark的spark.executor.memoryOverhead参数默认是executor内存的10%但这部分内存是给存储、连接缓冲等用的。如果数据里有大量递归解析或正则处理堆外内存可能先爆掉。我当时把spark.executor.memoryOverhead从默认值调到了512m问题就缓解了。5.2 TaskNotSerializable——闭包陷阱这个报错是新手非常容易踩的。现象是程序运行到某个时刻Spark突然报org.apache.spark.SparkException: Task not serializable。原因在于Spark的计算任务会被序列化后发送到不同的Executor如果我们在RDD的map或DataFrame的UDF里引用了当前Driver上的某个对象而这个对象没有序列化就会触发这个错误。我遇到的具体场景是在withColumn里用了一个自定义函数函数内部引用了ConfigParser初始化出来的配置对象。解决方案很简单把需要的字段作为参数传进函数不要在函数体内引用外部对象。如果确实需要共享配置可以用broadcast变量广播一份只读数据到所有Executor。现在用DataFrame API加内置函数基本不用写UDF这类问题也少了很多但一旦你开始写udf或rdd.map就一定要警惕闭包对象能不能序列化。5.3 数据倾斜部分任务特别慢数据倾斜的表现是整个作业卡在某一个stage绝大部分任务几十秒就完成了但有一两个任务要跑十几分钟甚至直接OOM。这是分布式计算的经典问题原因通常是某个key的数据量太大导致处理这个分区的Executor严重超载。我在分析“哪些地区价格波动最大”时遇到“山东”这个地区的记录特别多单独一个分区就占了整体数据的40%。处理方法有两种一种是把分区数调大让大key的数据散到多个分区另一种是加盐也就是先给大key加上随机后缀拆散聚合完再合并。对于作业而言最简单有效的做法是先查看spark.sql.shuffle.partitions当前值把200调成400或500让shuffle后的数据分得更细。如果调完还不行再考虑加盐方案。还有一个验证的小工具df.groupBy(region).count().orderBy(desc(count)).show()先看看哪些key的数据量是异常的。不先做这个检查你根本不知道是哪里倾斜。5.4 格式和依赖相关的隐蔽坑还有一些看起来很小但足以卡住半天的坑比如JSON文件里日期是yyyy/MM/dd格式Spark的DateType默认只能解析yyyy-MM-dd这时候要么用to_date函数传yyyy/MM/dd格式要么提前把日期字符串转换成标准格式。依赖问题也很常见。我在跑PySpark作业时遇到过Python环境里装了多个版本的numpyPySpark内部某些模块加载了不兼容的native库直接报Illegal instruction错误。解决方法是建一个干净的虚拟环境把PySpark和必要的依赖重新安装一遍。这里特别提醒如果本地Python版本和服务器不一致作业里用到了第三方库一定要把环境打包成.zip或使用--py-files参数一起提交否则在代码少一个依赖根本跑不起来。6. 关于第一次Spark作业的一些体会做完第一次作业如果要总结最核心的体会我会说三件事。第一不要跳过环境验证直接写大作业代码环境问题往往比代码问题更隐蔽、更难排查。第二不要等到最后才做结果验证每个清洗步骤都该打印一下行数数据从几万行变成几千行时如果没有合理说明心里要立刻打一个问号。第三不要害怕报错Spark的日志其实非常结构化先看是哪一段代码出错再往上游找数据问题多数报错都能在半小时内定位。最后分享一个我后来一直在用的小习惯写完每段代码后先用df.explain()看一眼执行计划。这个操作不会跑数据只是把Spark打算怎么执行打印出来能直观看到谓词下推有没有生效、哪些步骤是先shuffle再过滤。第一次作业时我完全看不懂这个输出但跑过几次、对照过几个问题之后它就成了我调优的起点。第一次做Spark作业跑通本身就是最有价值的收获。