ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Spark单机版环境搭建与本地模式实战:从零跑通核心API

Spark单机版环境搭建与本地模式实战:从零跑通核心API 很多初学者第一次接触Spark看到“集群”、“分布式”、“Driver”、“Executor”这些词就直接被劝退了。实际上大多数人在真正工作需要搭建集群之前需要的只是一个能在自己电脑上把RDD、DataFrame、Spark SQL、运行机制都跑通的环境。这篇内容对应的2.2.2.2小节核心就是教你用最低成本把Spark单机版环境搭起来然后在本地模式Local模式下验证Spark的核心能力。这个环境适合谁两类人最需要一是正在学Spark课程、需要完成作业和实验的学生二是准备面试、想在本地快速验证某些API行为的开发者。它不需要多高的硬件配置不依赖公司服务器也不用配Hadoop集群装好就能跑。我见过太多人一上来就折腾三台虚拟机搭完全分布式最后卡在SSH免密、NTP同步、端口不通这些问题上几天下来Spark的API一行都没写过。我的建议很直接先搭单机版跑通核心案例搞清楚数据是怎么被分片、怎么被并行处理的再去碰集群。这个顺序能帮你省下大量时间也符合从易到难的学习曲线。1. 单机版Spark环境到底在解决什么问题先别急着看安装命令想明白这个环境的价值后面遇到报错才不会慌。1.1 学习Spark为什么建议先跑通本地版Spark实际生产中几乎都是集群部署但在学习阶段单机版环境有一个无法替代的优势你启动的是完整的Spark运行时包括Driver、Executor、DAGScheduler、TaskScheduler这些核心组件Spark内部照样按分布式计算的流程去切分任务、调度执行只不过所有进程都跑在同一台机器上。这意味着什么你写一段rdd.map(...).filter(...).reduceByKey(...)的代码Spark会真实地经历“生成RDD血缘…划分Stage…生成Task…调度到Executor执行”的全过程。虽然数据量不大但执行机制和集群模式下是完全一样的。我经常跟学员打比方单机版Spark就像一台飞行模拟器虽然没真正上天但仪表盘、操作逻辑、故障响应都和真机一致。另外本地模式对数据源的支持也足够丰富。你可以直接读写本地文件系统的CSV、JSON、Parquet可以连MySQL、PostgreSQL甚至可以通过JDBC连到达梦这类国产数据库。很多人在开发阶段先用本地模式写业务逻辑测通了再丢到集群上跑全量数据这个工作流在业界非常常见。1.2 本地模式Local与集群模式的本质区别本地模式在Spark里叫Local模式核心特征是Driver和Executor运行在同一个JVM进程中。你用spark-shell启动时看到SparkContext web UI available at http://192.168.x.x:4040那个Web UI就是本地SparkContext提供的。从执行引擎的角度看关键区别在于资源调度方式。集群模式下Driver向YARN或Standalone的Master申请资源然后由Master在集群的多个节点上启动Executor容器Local模式则完全跳过外部调度器由SparkContext在本地直接创建Executor线程。在Local模式下spark.executor.instances这个参数是无效的实际并发度由local[N]或local[*]中的N决定。举个简单的例子你在本地模式用spark.sparkContext.defaultParallelism查看默认并行度local[*]通常会返回机器CPU核数local[2]就返回2。在集群模式下这个值取决于总核心数除以每个Executor的核心数再乘Executor个数。理解这个差异你就明白为什么本地代码测试结果和集群跑出来的Spark UI日志看起来不太一样。2. Spark单机版环境的选型与搭建这一节直接给到我验证过多次的安装组合和步骤。下面所有方案都在主流的CentOS 7.9和Ubuntu 20.04上实测过Windows 10/11也可以用同样思路只是环境变量路径写法略作调整。2.1 版本选择JDK/Scala/Spark的兼容性本地版最忌讳的就是版本乱配。很多报错的根本原因不是代码问题而是Spark和JDK、Scala版本不匹配。我建议默认采用这个组合JDK1.8对应Java 8或JDK 11。Spark 3.2以上官方声明支持JDK 8/11/17但国内绝大多数教材和生产环境仍以JDK 8为主如果你不是有特殊需求选JDK 8最稳。Spark3.3.x或3.4.x。这两个版本资料多、坑少遇到问题网上几乎都能搜到答案。Scala如果你不需要自己编译源码可以不单独安装Scala因为Spark发行包已经内置了编译好的Scala类库。只有当你想在本地直接用Scala写Spark代码时才需要额外装Scala。Hadoop单机版跑纯Spark任务也不需要装Hadoop除非你想练习Spark读写HDFS。很多教程让你先装Hadoop再装Spark是为了后续扩展到集群模式做准备如果只搭本地环境可以跳过省掉大量配置工作。版本对应关系可以用下面这张表来辅助判断组件推荐版本说明JDK1.88u202以上Spark 3.x兼容性最好Spark3.3.2或3.4.1预编译版选pre-built for Apache Hadoop 3.3Hadoop不需要可选需要读写HDFS时才装Scala不装也可以Spark内置akka、scala-library等依赖下载Spark时选择名称类似spark-3.3.2-bin-hadoop3.tgz的包这里的hadoop3只是标记编译时依赖的Hadoop版本并不表示运行时要额外装Hadoop。2.2 Linux环境准备与角色规划第一次接触Spark的同学我建议不要在Windows上直接跑虽然Windows上也能跑但会有文件权限、winutils.exe、本地网络模拟之类的问题排查起来非常烦躁。准备一个干净的Linux虚拟机或者直接用云服务器整个过程会顺畅很多。安装基础依赖并创建专用用户可选但推荐避免直接用root导致后续权限问题)# Ubuntu/Debian sudo apt-get update sudo apt-get install -y wget tar # CentOS/RHEL sudo yum install -y wget tar然后下载并解压Sparkcd /opt sudo wget https://archive.apache.org/dist/spark/spark-3.3.2/spark-3.3.2-bin-hadoop3.2.tgz sudo tar -zxvf spark-3.3.2-bin-hadoop3.2.tgz sudo mv spark-3.3.2-bin-hadoop3.2 spark配置环境变量编辑~/.bashrc或/etc/profileexport JAVA_HOME/usr/lib/jvm/java-1.8.0-openjdk export SPARK_HOME/opt/spark export PATH$PATH:$JAVA_HOME/bin:$SPARK_HOME/bin:$SPARK_HOME/sbin让环境变量生效source ~/.bashrc然后验证java -version spark-shell --version如果spark-shell --version能正常打出Spark版本信息说明基础环境已经通了。这里有一个细节$SPARK_HOME/conf/spark-env.sh在默认情况下不存在只有spark-env.sh.template模板。本地模式可以不配这个文件但建议显式指定JAVA_HOME因为某些系统上Spark脚本自动找JDK时会找错版本。创建配置文件cp $SPARK_HOME/conf/spark-env.sh.template $SPARK_HOME/conf/spark-env.sh然后在文件里加上export JAVA_HOME/usr/lib/jvm/java-1.8.0-openjdk export SPARK_LOCAL_HOSTNAME127.0.0.1SPARK_LOCAL_HOSTNAME这个参数经常被忽略如果本机hostname映射有问题Spark无法回连本地节点就会出现明明启动了却卡在Initializing SparkContext或Connecting to ...的问题。提前设成127.0.0.1能绕开这个坑。2.3 快速起一个Local模式的Spark Shell环境配好后先用交互式Shell验证是否真的能跑cd /opt/spark ./bin/spark-shell --master local[2]启动日志里如果出现SparkContext: Running Spark version 3.3.2、Using Sparks default log4j profile并且最后进入scala提示符说明环境OK。这时随手跑一个测试val rdd sc.parallelize(Seq(1,2,3,4,5)) rdd.map(_ * 2).collect()预期输出是Array(2, 4, 6, 8, 10)。如果你能在本地看到这个结果Spark单机版环境的核心链路已经走通接下来要考虑的就是怎么用它做正经事。3. 用本地模式跑通三个高频案例从API练习到数据分析工具装好的下一步是找场景练手。这一节我挑三个一定会用到的实操场景分别覆盖RDD API、提交作业和Spark SQL读外部数据源。3.1 场景一用spark-shell验证RDD与DataFrame API对初学者来说spark-shell是最合适的数据结构练习场。RDD虽然在实际业务中已经逐渐被DataFrame取代但很多算子思想仍然是Spark的底层逻辑尤其是reduceByKey、broadcast这些概念面试必问。我在练习时建议从“单词计数”这个经典案例开始。虽然烂大街但它完整覆盖了从加载数据到聚合输出的所有核心步骤val lines sc.textFile(file:///home/user/input/words.txt) val words lines.flatMap(_.split( )) val pairs words.map(word (word, 1)) val counts pairs.reduceByKey(_ _) counts.collect().foreach(println)注意这里textFile要写file://前缀如果不写Spark默认去HDFS上找文件而你本地又没配Hadoop就会报FileNotFoundException。这是初学者最容易踩的坑之一。等你把RDD的算子流程理清后再切到Spark的核心抽象DataFrameval df spark.read.option(header, true).csv(file:///home/user/input/sales.csv) df.printSchema() df.groupBy(category).sum(amount).show()DataFrame比RDD更直观因为它带着Schema信息Spark SQL优化器Catalyst可以直接优化执行计划你的代码跑起来也更高效。3.2 场景二用spark-submit提交一个完整的数据分析任务实际开发中代码不会写在spark-shell里而是打包后通过spark-submit提交。写一个简单的Python脚本或Scala程序模拟统计订单收入。以Python为例文件名为order_stats.pyfrom pyspark.sql import SparkSession spark SparkSession.builder \ .appName(OrderStats) \ .master(local[2]) \ .getOrCreate() df spark.read.option(header, True) \ .csv(file:///home/user/input/orders.csv) df.createOrReplaceTempView(orders) result spark.sql( SELECT region, SUM(amount) AS total_amount, COUNT(*) AS order_cnt FROM orders GROUP BY region ORDER BY total_amount DESC ) result.show()然后使用spark-submit提交注意要和python环境匹配如果是pyspark的话需要先安装pyspark库cd /opt/spark ./bin/spark-submit \ --master local[2] \ --name OrderStats \ /home/user/py_scripts/order_stats.py你会看到Spark在日志里显示作业分成了若干Job和Stage最终在控制台打印出每个地区的总金额。这个过程和你将来在集群上提交作业的流程完全一致唯一差别只是--master从local[2]换成了yarn。3.3 场景三Spark SQL连接外部数据源与JDBC适配本地模式跑通了纯文件读写的例子之后试一下连外部数据源。这是很多开发者在实际工作中真正要面对的业务数据存在数据库里Spark需要把它拉出来做分析。我用MySQL做示例。假设MySQL里有张订单表你需要先下载MySQL的JDBC驱动然后把驱动放到Spark的jars目录或者通过--packages参数动态引入./bin/spark-shell --master local[2] --packages mysql:mysql-connector-java:8.0.28然后可以执行val jdbcDF spark.read .format(jdbc) .option(url, jdbc:mysql://localhost:3306/testdb) .option(dbtable, orders) .option(user, root) .option(password, yourpassword) .load() jdbcDF.groupBy(user_id).count().show()同样的逻辑换成达梦数据库DM时只需要把JDBC URL和驱动类改成dm.jdbc.driver.DmDriver使用jdbc:dm://localhost:5236格式的URL。这几年国产化替换推进很多单位要求从Oracle/MySQL切换到达梦Spark本身不支持数据源直连达梦因为它不是内置的DataSource实现但JDBC是通用标准通过驱动就能适配。本地模式你先验证好驱动能不能加载、数据能不能读写上生产时才不至于手忙脚乱。Redis也是热词里频繁出现的数据源。Spark没有内置Redis连接器常见做法是用spark-redis这个开源库或者直接在mapPartitions里创建Redis连接来做读写。本地模式下建议先在spark-shell里用小批量数据验证序列化逻辑和连接池配置确认没问题再把吞吐量提上去。4. 本地版Spark的执行机制与故障现场排查环境跑起来只是第一步我更希望你能在这套环境里真正理解Spark运行时。下面聊聊本地模式里最好观察、也最容易出问题的几个点。4.1 Driver、Executor和并行度的分配逻辑Spark应用程序跑起来后有个进程叫Driver它负责解析代码、生成执行计划、把任务分发到Executor上去。本地模式的Driver和Executor都在同一个JVM进程里所以你在jps命令里只会看到一个SparkSubmit进程。这个模型有个非常直观的好处你可以在Spark的Web UI默认4040端口里看到一份“五脏俱全”的执行视图。当你在本地跑了一个groupBy操作打开http://localhost:4040/jobs/能看到Job、Stage、Task的层级关系甚至能看到每个Task读取了多少数据、耗时多少毫秒。我很推荐一个练习方法在spark-shell里跑一段代码后马上打开Web UI看它的执行计划。比如rdd.map(...).filter(...).collect()会触发一个Job里面包含至少两个Stage你会清楚地看到shuffle发生在哪个环节。这种直观体验比死记硬背“宽依赖和窄依赖的区别”强一百倍。并行度设置方面本地模式用local[2]就意味着Executor线程池里有2个线程同一时刻最多有2个Task并行执行。如果你有大文件要处理建议把spark.default.parallelism和文件分片数对齐。我通常会在测试脚本里加上spark-submit --master local[4] --conf spark.default.parallelism4 ...这样执行时的分区数预期比较明确不会出现一个几KB的小文件也被拆成几十个Task的荒唐情况。4.2 本地模式常见的OOM和性能陷阱单机版环境虽然简单但“能用”和“用得舒服”之间有一段路要走。最常见的坑就是OOM堆内存溢出。比如你在本地启动spark-shell默认分配1G内存然后一股脑加载了大文件再反复做shuffle很快会看到java.lang.OutOfMemoryError: Java heap space。原因很简单Driver和Executor共享一份堆内存数据量一旦超了阈值就会被打爆。解决思路有几种调大Spark进程内存。在spark-env.sh里设SPARK_DRIVER_MEMORY4g或者提交时指定--driver-memory 4g。控制分区数避免每个分区内数据过于庞大。比如repartition(100)把数据打散到更多分区里。减少collect回Drvier的数据量。很多人写代码习惯最后collect()把所有结果拉到本地打印数据量一大就OOM。正确做法是先filter或limit缩小结果集或者只take(10)看看前几条。另外一个本地模式特有的性能陷阱是文件读取路径。如果读取的是NFS或虚拟机共享目录里的文件I/O等待会被拉满一个简单的count可能跑几十秒。这时候先把文件用head或cp复制到本地磁盘性能往往能提升一个数量级。5. 单机调试中典型报错的排查与一手心得本地环境虽然简化了分布式链路但该出的问题一点不少。我把自己踩过的、以及在指导别人过程中高频遇到的报错整理在这里方便你按图索骥。5.1 高频Error速查报错信息原因解决办法Exception in thread main java.lang.NoClassDefFoundError依赖缺失比如MySQL驱动没放进来检查jars目录或使用--packages引入SparkException: A master URL must be set in your configuration没指定master加--master local[2]或代码里加.master(local[2])FileNotFoundException: file:///...路径不对或者文件放在HDFS的错觉确认文件路径前缀和现实目录位置OutOfMemoryError: Java heap space堆内存不足调大driver-memory并检查是否有过多collect()Cannot run program python3运行pyspark时操作系统没装Python安装python3并配置PYSPARK_PYTHON环境变量Port 4040 in use本地4040端口被占用之前Spark进程没退出换端口或用jps看进程并killjava.net.BindException: Cannot assign requested addressSPARK_LOCAL_HOSTNAME和实际地址不匹配在spark-env.sh里设SPARK_LOCAL_HOSTNAME127.0.0.1这里有一个大家最容易忽略的问题跑完一批作业后不要直接关终端先看一眼4040页面和Driver日志。因为很多报错只会留在日志文件里控制台只Print堆栈的关键几行。完整的日志路径一般在$SPARK_HOME/logs/下文件名类似spark-user-org.apache.spark.deploy.SparkSubmit-1-localhost.localdomain.out。5.2 排查思路与进阶调试技巧排查Spark问题我的基本思路可以归纳为三步看Web UI、看日志、看数据血缘。先说Web UI。不管是本地模式还是集群模式4040页面的Event Timeline和Stage详情页能直接告诉你每个Stage的输入记录数、shuffle读写量、GC时间。比如某个Stage的输入是10万条Shuffle写却达到1GB这说明发生了严重的数据膨胀通常是由crossJoin或笛卡尔积引起的。再说日志。前面提到Spark的日志分两部分框架日志和应用日志。框架日志默认打印在控制台应用日志往往靠log4j里面配置的Appender输出。初次排查时可以在代码里加点临时打印用println输出每个阶段的行数。我自己调试时常用print(after filter count:, df.filter(...).count())这种“土办法”在本地阶段比任何分布式Debug工具都好使因为数据量小跑一轮也就几秒。最后说数据血缘。遇到结果不对的情况不要在结果层挠头往前一步一步看每条转换是不是都符合预期。用RDD时可以通过rdd.toDebugString看血缘用DataFrame时用df.explain(true)看完整执行计划。这个习惯最好在单机版就养成等以后上了集群数据量大到你没法“打印看中间结果”时执行计划就是你唯一的线索。6. 3个容易让初学者栽跟头的“非报错”问题排障不只是看报错有些“没报错但结果不对”的情形更隐蔽。这一节单独说说容易被误解的概念。6.1 local[N]、local[*]、local[N, M]到底怎么写很多人在local模式的参数上照搬网上命令出了性能问题却不知道原因。local[2]代表用2个线程跑Executorslocal[*]代表用本机所有CPU逻辑核心数local[2, 4]的第二个参数是表示要重试的最大失败次数默认是1。在笔者的建议里日常练习就直接写local[*]让Spark自动感知CPU核心性能最优还省心。6.2 本地跑通了集群上为什么还是挂这是个很常见的灵魂拷问。根本原因是单机版掩盖了网络和并发问题。本地模式下所有数据都走本地磁盘和内存不会遇到跨节点传输的网络延迟Executor生命周期随Application结束而结束不存在长期驻留的进程管理问题Driver所在机器就是运行任务的主机不用考虑跨节点回连。所以本地测试通过不等于集群万无一失。常见的坑包括代码里硬编码了本机文件路径如file:///home/user/data.csv上了集群却读不到使用了第三方库打包时没有带--jars没有给Executor设置足够的内存导致集群上数据量大后OOM。这些问题我在带人做项目时几乎每次都会见到。本地版是用来跑通逻辑的不是用来模拟集群运维的心里要有这根弦。6.3 本地环境如何以“半集群”方式过渡如果想让本地环境稍微靠近生产一个性价比很高的折中方案是在单机上用Standalone模式启动Spark集群。在这个模式下Master和Worker进程虽然都在同一台机器但已经有了独立的资源调度过程。你可以用start-master.sh启动Master用start-worker.sh启动Worker再通过spark-submit --master spark://127.0.0.1:7077提交作业。这样做的好处是你能练习集群模式下Driver和Executor的分离、观察Worker上Executor的启动过程、体验Web UI里Master和Worker的状态展示。坏处是它比纯Local模式会多一些网络栈开销和进程管理复杂度。我的建议是初学阶段先用Local模式把API和业务逻辑跑熟当你想研究资源调度、Executor启动参数时再启动Standalone练手。这种过渡方式还有一个容易被忽视的好处如果将来你要部署到YARN或Kubernetes上很多概念是相通的比如Executor的资源规格、Driver与Executor的健康检查、任务的日志聚合方式等。你在Standalone模式下积累的经验换到其他资源管理器时并不会白费。7. 把单机版环境用出“集群味”的几个习惯写到这里想分享几个我每次搭Spark环境都会做的事。这些事不影响跑通教程但能避免后续反复踩坑。第一把常用启动参数写进脚本而不是手动敲。我在/opt/spark下建了一个run_spark_shell.shexport SPARK_HOME/opt/spark export SPARK_LOCAL_IP127.0.0.1 $SPARK_HOME/bin/spark-shell \ --master local[*] \ --driver-memory 4G \ --conf spark.sql.shuffle.partitions4 \ --conf spark.ui.port4040每次进入Spark Shell时执行这个脚本保证Driver内存和shuffle分区数一致避免不同的调试会话之间因为参数不一致得出互相矛盾的结果。第二本地数据样本要能覆盖边界情况。我见过有人用一百行测试数据把分析逻辑跑通了结果一换真实数据就出错原因是代码里没处理NULL值或字段里的空格。建议你在本地准备两套数据一套是教材里的玩具数据用于快速验证逻辑另一套是“脏数据样本”刻意包含NULL、空字符串、超长字段和重复主键。这样跑下来的逻辑才更经得起推敲。第三经常查看Spark的默认配置。执行下面这条命令Spark会把所有配置项的默认值按前缀分组列出来./bin/spark-shell --conf spark.sql.shuffle.partitions4 -e spark.conf.getAll.sorted.foreach(println)每天花几分钟翻阅这些配置你会慢慢明白spark.sql.shuffle.partitions为什么默认是200spark.executor.memory在集群模式下怎么被YARN回收spark.serializer从Java序列化换成Kryo后能带来多少性能收益。这些认知积累起来比单纯背面试题管用得多。8. 从本地向真实业务落地的扩展建议环境能跑、代码能执行只完成了第一阶段。分享几个更进阶的方向按“需要时再深入”的顺序排列。先提一个搜索热词里的“spark数据分析案例”。如果你只想更熟练地使用Spark建议自行找一份多表数据关联的练习比如“订单明细商品维度用户维度”的三表关联再统计不同维度的GMV、复购率。不要满足于一两个groupBy在线业务的分析场景往往有多个聚合条件和时间窗口把这些逻辑在本地用Spark SQL实现一遍你自然能体会到DataFrame API和SQL各自的优劣。接下来是“spark OOM”。本地模式内存受限恰好是学习OOM处理的好场景。刻意调低--driver-memory到512M这种极限值然后跑一个groupByKey操作观察Spark什么时候抛出Job aborted due to stage failure再尝试用repartition、reduceByKey替代groupByKey等方案解决。这个过程会让你对shuffle的代价、内存堆叠的机制有更深的理解比看一百篇OOM文章都有用。最后是“spark集群搭建”。如果本地练习已经不能满足需求下一步可以规划在3台虚拟机上搭Standalone或YARN集群。那时你会真正面临配置分发、日志收集、资源队列等运维问题。回头再来看这篇单机版环境的内容你会发现很多概念像是地基没有这层地基直接盖集群的楼很容易塌。我个人在实际操作中的体会是单机版环境用得好不好不在于你抄了多少命令而在于你是否愿意盯着Web UI去观察一次Job从提交到结束的完整履历以及是否舍得花时间把一次OOM排查到底。Spark的学习曲线确实陡但它给了你一个足够宽容的起点。把这套本地环境玩明白后面的路会顺很多。
返回列表