
第一次做 Spark 作业的人多少都有点被“大数据”这三个字唬住。我最初也是这样以为要搭一个多壮观的集群处理多吓人的数据量结果真正动手发现第一次作业的核心根本不是“大”而是把 Spark 的基础链路跑通环境能起来、数据能读进来、逻辑能算出来、结果能看明白。这篇文章就围绕这份“第一次作业”的完整经历来写从环境搭建到代码实现再到调试排错把我踩过的坑、验证过的方法、总结出来的经验都放进来给同样在入门阶段的朋友一条可以直接照着走的路。1. 第一次Spark作业到底在考什么1.1 作业背后真正要掌握的能力很多第一次布置 Spark 作业的老师或团队领导本质上不是要你搞出多复杂的算法而是要确认三件事你有没有把 Spark 环境跑起来知不知道程序提交后发生了什么能不能用 DataFrame 或 SQL 完成基本的数据处理。说白了这是一次“全链路体检”。我拿到的作业要求也很典型读取一份 JSON 数据做清洗再完成几个统计指标。听起来简单但它刚好覆盖了 Spark 最核心的几个环节输入、转换、输出。这里我特别想提醒第一次接触 Spark 的同学不要一上来就研究 RDD 的复杂算子也不要沉迷于调优参数。第一次作业最重要的是“闭环”。从启动 Spark 到写完代码到看到输出结果这个完整链路只要走通一次后面所有的进阶学习才有附着点。我自己第一次做作业的时候在 spark-shell 里连 sc 和 spark 的区别都没搞明白就开始写代码结果浪费了大量时间在莫名其妙的问题上。1.2 环境选型先选对武器再动手环境怎么搭是第一次作业里最劝退人的环节也是搜索引擎里 spark 相关热词里搜索量最大的部分。我当时纠结了很久是本地装一个单机版还是用虚拟机搭集群我的结论是第一次作业用本地单机版就够了而且优先选 Standalone 模式不要碰 YARN更不要尝试三台机器以上的集群。理由很简单。作业的数据量基本都是 MB 级别单机完全跑得动本地模式不需要处理跨节点通信、用户权限、资源调度这些额外问题可以把注意力全放在 Spark 本身的语法和逻辑上。我见过有人第一次作业就非要在三台服务器上搭集群最后折腾了两天网络配置连第一个 wordcount 都没跑出来这完全违背了作业的初衷。我推荐的方式是Mac 或 Linux 笔记本安装 JDK 8下载 Spark 预编译包解压即用。Windows 也行但需要额外装 Hadoop 的 winutils比较麻烦如果条件允许借一台 Linux 机器或虚拟机能省掉很多坑。语言上我选了 Python PySpark后面会细说为什么。2. 从零搭建Spark运行环境2.1 JDK、Hadoop与Spark的版本匹配问题版本匹配是我踩的第一个大坑。很多新手直接下载了最新版 JDK 17 和 Spark 4.x结果 Spark 启动的时候报了一堆模块相关的错误一脸懵。这里要记住一个原则Spark 是一个“吃版本”的框架它不是越新越好而是要匹配稳定组合。实际验证下来最稳的组合是 JDK 8 配 Spark 3.3.x 或 3.2.x。如果你用 Python还需要确认 PySpark 的版本和 Spark 主版本保持一致。我第一次装的时候Spyder 里早就装了一个老版本的 pyspark和后来单独下载的 Spark 版本对不上导致明明正确写的代码一直报类找不到的错误折腾了整整一个下午才反应过来是两个版本混了。还有一点要注意Spark 它自己带了一个 HDFS 客户端依赖但本地跑作业时不需要搭建完整的 Hadoop 集群只需要把 Hadoop 相关的环境变量指向 Spark 内置的依赖就行。很多人不知道这一点非要去装一个完整的 Hadoop 发行版属于纯粹的浪费时间。注意常见的报错java.lang.NoClassDefFoundError: org/apache/hadoop/...多半是环境变量或版本冲突导致的优先检查 spark-env.sh 和 ~/.bashrc 里是否指向了同一个 Hadoop 路径而不是急着重装系统。2.2 Spark目录结构与核心配置文件Spark 解压完以后你会看到一堆目录第一次看很容易头晕。但我用下来发现真正需要关心的只有三个地方conf/目录里的模板配置文件bin/目录里的启动脚本以及logs/目录里的运行日志第一次启动后会自动生成。先说配置文件。进入conf/目录把spark-env.sh.template复制一份成spark-env.sh然后在里面设置 Java 路径和 Spark 的工作目录。我最初没有设置SPARK_LOCAL_IP结果启动历史服务的时候一直报地址绑定错误加上这一行就解决了。还有一个值得关注的变量是SPARK_WORKER_MEMORY它决定了每个 Worker 能使用的默认内存大小如果你只是跑作业不启动 Worker 进程这个其实不太重要但很多教程都让你设设上也无妨。接下来是启动脚本。本地跑作业有两种方式一种是直接在终端里执行spark-submit提交脚本另一种是在bin/pyspark或spark-shell里交互式写代码。我第一次做作业的时候先用交互式调试因为可以一行一行看到结果等调试完再整成脚本用 spark-submit 提交。这个思路很推荐给所有人先用交互模式快速验证语法和逻辑再用脚本模式做最终提交效率高很多。2.3 启动验证跑通第一个Spark程序环境配好以后一定要先跑一个最小程序验证环境不要直接上作业内容。我习惯用的一个“环境体检”代码非常简单from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(EnvironmentCheck) \ .master(local[2]) \ .getOrCreate() data [(a, 1), (b, 2), (c, 3)] df spark.createDataFrame(data, [letter, number]) df.show() print(Total count:, df.count()) print(Spark version:, spark.version) spark.stop()这段代码如果能在终端里正常输出三行数据、一个计数和版本号说明 Spark 环境基本没问题。注意.master(local[2])里的[2]表示使用 2 个本地线程这个参数在第一次作业里随便设就行只要机器内存不是太小设成local[*]自动使用所有可用 CPU 线程也是 OK 的。我第一次跑这段代码时卡在了一个很奇怪的地方pyspark 能 import但是SparkSession.builder一执行就开始报错后来发现是 Python 默认版本是 3.11而 PySpark 3.3 对这个版本支持不全于是降到 Python 3.8 就好了。这种环境问题几乎每个人都会遇到排查思路就是一点点减少变量不要同时改多个配置。3. 作业实操基于订单数据的清洗与统计3.1 作业需求拆解与数据准备我拿到的作业数据是一个“模拟订单流水”一份 JSON 文件大概几千条记录字段包括 order_id、product_name、category、price、quantity、order_time、customer_region。作业要求做三件事清洗明显异常的数据统计每个品类的销售总额和订单量找出下单量最高的前五个商品。这个需求非常典型几乎覆盖了 DataFrame 的日常操作比单纯跑一个 wordcount 有营养得多。分析任何 Spark 作业的第一步都不是写代码而是先看数据。我强烈建议在写任何处理逻辑前先用任意工具打开原始 JSON 文件人工浏览前几十行搞清楚四个问题字段名是什么、嵌套结构是什么样的、哪些字段可能为空、数值字段的格式是否统一。我第一次就没看数据直接写代码结果读出 DataFrame 以后发现 price 字段里有“12.5元”这种脏数据如果提前看过原始数据很多清洗逻辑可以一次性写对。作业数据我放在了和 Spark 脚本同一个目录下的data/文件夹里目录结构如下spark_homework/ ├── data/ │ └── orders.json ├── job/ │ └── analysis.py └── output/这样组织的最大好处是提交命令里用相对路径就行不用写一长串绝对路径也方便反复调试时看清楚输入输出在哪里。3.2 读取JSON数据并转为DataFrameSpark 读取 JSON 比想象中简单它自带spark.read.json()方法一个调用就能解析整个文件。但这里有个很长一段时间我一知半解的细节这个方法会自动根据文件内容推断 schema包括嵌套字段和数组类型省去了手动定义表的步骤。对于第一次作业这种半结构化数据来说体验非常好。我实际用的读取代码是from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, DoubleType, IntegerType spark SparkSession.builder \ .appName(FirstHomework) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 4) \ .getOrCreate() df spark.read \ .option(multiline, true) \ .json(data/orders.json)有一点特别值得说明如果你的 JSON 文件是每行一个对象的标准格式multiline选项设不设都无所谓如果整个文件是一个大 JSON 数组或者格式化过的多行 JSON就必须设置multilinetrue否则会解析失败。我那次作业的数据恰好是多行格式化形式第一次没加这个参数读取结果只有一行数据半天没反应过来。读取完成以后一定要执行df.printSchema()和df.show(5)这两行输出就是你的“数据体检报告”能确认字段名有没有被自动改成小写、类型判断准不准、嵌套结构有没有被摊平。我做完这两步才发现 price 字段被推断成了 StringType因为里面有脏数据后续就必须先清洗类型。3.3 数据清洗空值、类型、去重一条龙数据清洗在作业里占了最重的篇幅也是面试官或老师最愿意追问的部分。我把清洗拆成三个步骤来做先处理空值再修正类型最后去重。空值处理我用的是分列判断而不是一刀切dropna()。产品名称或件数如果为空这行数据没有分析价值直接删除价格异常大的也用过滤条件排除。但下单时间为空的记录有时候可以通过其他字段推断所以我保留了抽查后手动处理的选项。真正执行的代码大致如下# 1) 删除产品名为空、价格为空的记录 clean_df df.dropna(subset[product_name, price]) # 2) 过滤异常价格数据低于0或高于100000的视为脏数据 clean_df clean_df.filter( (clean_df[price] 0) (clean_df[price] 100000) ) # 3) 清洗价格字段去掉元等后缀并转成DoubleType from pyspark.sql.functions import col, regexp_replace clean_df clean_df.withColumn( price_clean, regexp_replace(col(price), 元, ).cast(double) ) # 4) 基于订单ID去重 clean_df clean_df.dropDuplicates([order_id])这一步里有几个细节非常容易被忽略。比如regexp_replace的返回值是字符串类型必须再跟一个.cast(double)才能变成数值类型否则后面求和会变成字符串拼接。再比如dropDuplicates时如果指定多个字段作为去重键那任何组合一样的行都会被删除只留第一条这个语义和 SQL 里的DISTINCT不完全一样。清洗完以后我建议立刻把结果clean_df.count()和清洗前df.count()对比一下并且打印出清洗掉的记录数。这个数字写进作业报告里特别加分因为它说明你能意识到数据质量对结果的影响而不只是机械地调用 API。3.4 Spark SQL统计分析完整脚本清洗完成后统计部分我用了两种方式实现一种是用 DataFrame API另一种是注册成临时视图后用 Spark SQL。推荐第一次作业的同学至少尝试一次 Spark SQL 的方式因为无论是面试还是后续做项目SQL 的表达方式都更通用而且createOrReplaceTempView注册临时表这个操作是通往底层原理理解的一扇门。完整脚本的核心部分是这样的# 注册临时视图 clean_df.createOrReplaceTempView(orders_clean) # 统计每个品类销售额、订单量 category_stats spark.sql( SELECT category, ROUND(SUM(price_clean * quantity), 2) AS total_sales, COUNT(DISTINCT order_id) AS order_count FROM orders_clean GROUP BY category ORDER BY total_sales DESC ) # 找出下单量最高的前5个商品 top_products spark.sql( SELECT product_name, SUM(quantity) AS total_qty FROM orders_clean GROUP BY product_name ORDER BY total_qty DESC LIMIT 5 ) category_stats.show() top_products.show() # 写结果到CSV文件并标注为覆盖写入 category_stats.write.mode(overwrite).csv(output/category_stats) top_products.write.mode(overwrite).csv(output/top_products)写完这段以后我还做了两件额外的事。一件是执行clean_df.cache()把清洗后的数据缓存起来因为后面要跑两个不同的统计任务缓存可以避免重复读取源文件、重复执行清洗逻辑。另一件是把两个统计结果都写到了 output/ 目录下因为 Spark 默认是生成一个文件夹而不是一个文件文件夹下面会有part-00000-xxx.csv这样的真正数据文件很多人第一次看到会以为自己写失败了这是正常现象拼接一下文件名就是完整结果。4. 作业调试与常见问题排查4.1 日志里最常出现的几个报错第一次写 Spark 作业报错几乎避不开关键是别慌。我把自己和周围人踩过的最常见错误整理成了速查表遇到问题可以先对号入座。报错信息片段大概率原因解决思路Py4JJavaErrorPython 与 JVM 通信时底层 Java 代码抛了异常往前翻Caused by:部分那才是真正的原因AnalysisException: Path does not exist输入文件路径不对或相对路径基准错误改成绝对路径或确认执行 spark-submit 时的工作目录java.lang.OutOfMemoryErrorExecutor 内存不够调大spark.executor.memory或减少 shuffle 分区数Column xxx does not existschema 中列名大小写不匹配执行printSchema()确认列名Spark 默认不严格区分大小写TextFileFormat相关的写入错误输出目录已存在且未设置覆盖写入前加mode(overwrite)特别说一下Py4JJavaError这个报错对 Python 用户来说最具迷惑性因为一大堆堆栈信息看得人头大真正的异常提示被淹没在中间的 Java 堆栈里。我调试时的习惯是在终端里搜关键词Caused by:它后面跟的内容基本就是根本原因比如列名错误、类型不匹配、路径异常都能在几行内找到。4.2 性能与内存问题Driver和Executor怎么调第一次作业的数据量不大正常来说根本不会触发性能问题但有一点我却花了很长时间才真正理解为什么会 OOM明明数据才几 MB关键原因是 Spark 的内存不只看数据量还看任务执行时的中间结构尤其是 shuffle 操作。GROUP BY这种操作会触发全量数据重分区也就是 shuffle这个过程会在内存和磁盘之间产生大量临时数据。如果这个时候spark.sql.shuffle.partitions设置的默认值 200 个分区被触发即使数据很小也可能每个分区都分配了独立的缓冲积少成多就把内存吃满了。所以我给自己设了一个基准在本地跑小型作业时显式加上spark.sql.shuffle.partitions4只保留必要的并行度。这个参数我第一次完全没概念是看到报错里一直有ExternalSorter的日志才去查的详细说明。另一个值得调的参数是spark.executor.memory。在本地模式下这个参数基本等同于驱动进程可用的堆内存默认 1G 并不大。如果你的机器有 8G 内存把它调到2g是安全且舒适的。但注意不要贪心比如4g在只有 8G 内存的机器上会拖垮系统spark 进程是会真实占用物理内存的不是你设了它就一定分配那么多。4.3 结果验证技巧确认答案真的正确作业的统计结果出来以后不要直接交给老师或直接写进报告。第一次做的时候我不放心把结果导出来以后用 Excel 做了一遍同样的统计结果发现对不上后来检查才发现是清洗规则里价格字段的脏数据没过滤干净导致求和比预期大。这个教训非常深刻。这里分享一套我自己后来固定的“过河验证法”先数一下原始记录总数再数清洗后记录总数检查差值是否符合预期再找一个最小的业务口径比如某一个已知分类的销售额单独用filter查出来算一遍和统计数据里的对应数字对一下最后用describe()看关键数值列的 min、max、mean 是否在合理范围内。三步都过了统计结果基本就靠谱了。# 验证例单独计算某一类别的销售额 pe_orders clean_df.filter(clean_df[category] 数码产品) pe_total pe_orders.selectExpr(SUM(price_clean * quantity) AS total).collect()[0][0] print(单独计算值:, pe_total)虽然只是作业但养成“结果可解释”的习惯特别值钱。真正的业务环境里没有人会验收你的代码但所有人都会质疑你的数据。如果一个数据结果能被人问一句“你为什么算出这个数字”而不被问倒那这份作业和很多工作里最核心的能力就是一致的。5. 第一次作业结束后的进阶路线5.1 把作业升级成小型项目交完作业只是开始。我强烈建议在第一次作业的基础上花两个晚上把它扩展成一个小型项目扩展方向有三个任选一个都能让简历或总结里多一个可以讲的故事。方向一是加数据源原来只读 JSON可以加一份 CSV 或 Parquet 文件做 join 操作方向二是加输出方式原来是写文件可以尝试写入数据库比如 MySQL体会JDBC写入的配置流程方向三是加调度用脚本循环处理多个日期的数据文件模拟一次定时任务。我自己选的是方向一因为 join 操作是 Spark 里最常用、面试也最爱问的点。我人为制造了两份数据一份订单记录、一份商品分类字典然后通过product_name做 inner join 和 left join 对比还顺手体验了broadcast join的加速效果。这个扩展过程让我真正理解了为什么 Spark 的 join 那么快以及什么情况下需要广播变量。结构化作业文档也要同步准备。我后来的标准格式是需求背景、数据描述、环境信息、处理流程配流程图或步骤列表、核心代码、结果截图、遇到的问题及解决、改进方向。这份文档的完整度很多时候比代码本身更能说明你的能力尤其对刚入门的人而言。5.2 学习顺序建议与避坑清单如果让我给第一次作业的后来人一条学习顺序我会说先跑通 wordcount再做一次含 join 的数据统计再接触流处理或者结构化 Streaming。不要在第一次就研究调优也不要在第一次就陷入 Spark 内存模型的细节里。同时我总结了一份避坑清单都是从自己和周围人真实踩过的坑里提炼的版本组合固定先用 JDK 8 Spark 3.3 Python 3.8或 3.10不要盲目追新。代码里的路径一律优先用绝对路径或者运行时在项目根目录执行命令别指望相对路径永远正确。每次修改代码后先跑一个count()验证读取正常再跑完整逻辑。使用mode(overwrite)时小心误删之前的结果建议每天结果输出到带日期的新目录。日志一定要保留遇到报错先看最后几行和Caused by不要执着于从头读完整日志。写 SQL 统计前先打印两行清洗后的样例数据确认转换逻辑生效再继续。这份清单看起来短但每一条背后都是真实教训。比如日期目录这件事我一开始总是输出到同一个output/路径结果某次模式写错了把前一天好不容易调出来的结果直接覆盖了再想复现当时的场景就得重新调半天从那以后不管多小的任务都带日期目录。结语的一点亲身经验做完这次第一次作业我最大的收获不是会写了几个算子而是终于把 Spark 从“抽象的大数据框架”变成了“熟悉的工具”。环境搭好以后它其实就是一个能自动分布式处理的表格计算器和 Pandas 的差距只在于它是运行时懒加载、按分区组织内存、支持更大规模的横向扩展。后面再学习分区器、Shuffle 优化、检查点这些东西我都觉得顺理成章因为前面的基础链路在脑海里已经有了物理实感。如果你正卡在第一次作业的某个环节记住几乎所有问题都是环境问题或路径问题代码本身反而很少出大错。把环境跑通一次把流程走完一遍你就已经比一大半还在观望的人领先了。