尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Windows下从零搭建Spark本地集群:环境配置与实战避坑指南

Windows下从零搭建Spark本地集群:环境配置与实战避坑指南 想在Windows上跑Spark很多人第一反应是这玩意儿不是跑在Linux服务器上的吗。我一开始也这么想直到有次接了个数据分析的私活手头只有一台Windows笔记本客户又要求用Spark做用户复购率分析硬着头皮在Windows上折腾了整整两天踩了一堆坑之后才算跑通。后来复盘发现其实Windows下搭建Spark本地集群并没有想象中那么难难的是网上大部分教程要么默认你有Linux基础要么跳过了Windows特有的那些坑。这篇就把我从零搭建的完整过程、每一步背后的逻辑、以及那些教程里不会写的坑全部摊开讲清楚。Spark本地集群在Windows上的搭建核心要解决三件事Java环境的版本匹配、Hadoop相关依赖的Windows适配、以及Spark自身的配置调优。这三件事任何一件出问题你看到的报错都会让你怀疑人生。下面按照实际操作的顺序把每个环节拆开讲。1. 先把Java环境这件事说透1.1 为什么JDK版本选错后面全是坑Spark是用Scala写的跑在JVM上所以Java环境是地基。但这里有个很多人不知道的细节Spark不同版本对JDK的要求是不一样的。Spark 3.0到3.2官方推荐JDK 8Spark 3.3开始正式支持JDK 11和17但如果你用的是Spark 2.x系列那基本只能锁死在JDK 8。我一开始图省事装了JDK 17结果跑Spark 3.2的时候各种IllegalAccessError查了半天才发现是模块化系统导致的反射访问限制。后来换成JDK 8世界瞬间安静了。具体怎么选看这张表Spark版本推荐JDK备注Spark 2.4.xJDK 8只能用8别想别的Spark 3.0-3.2JDK 8官方主推811也能跑但偶有问题Spark 3.3JDK 8/11/1717需要额外加JVM参数对于零基础的朋友我的建议是直接用JDK 8 Spark 3.2.x这个组合。这是经过大量生产验证的稳定搭配网上资料也最多遇到问题好搜。1.2 JDK 8在Windows上的安装与验证去Oracle官网或者用国内镜像下载JDK 8的Windows安装包双击安装。安装路径千万不要带空格和中文比如C:\Program Files\Java\jdk1.8.0_301这种路径里的空格后面配置Spark的时候会让你抓狂。我一般直接装到C:\Java\jdk1.8.0_301。装完之后配置环境变量这一步很多人会漏新建系统变量JAVA_HOME值为C:\Java\jdk1.8.0_301编辑Path变量添加%JAVA_HOME%\bin新建CLASSPATH值为.;%JAVA_HOME%\lib\dt.jar;%JAVA_HOME%\lib\tools.jar配置完打开新的命令行窗口注意一定要新开旧窗口读不到新环境变量执行java -version看到类似java version 1.8.0_301就说明成功了。如果提示找不到命令八成是Path没配对或者你忘了开新窗口。提示如果你电脑上已经装了其他版本的JDK建议先把旧的卸载干净或者至少确保JAVA_HOME指向的是JDK 8。多版本共存虽然可行但对新手来说就是给自己找麻烦。2. Hadoop依赖Windows下最容易被卡住的地方2.1 为什么跑Spark还要装Hadoop这是零基础朋友最困惑的一点我明明只想用Spark为什么还要搞Hadoop原因在于Spark本身不负责文件存储它默认用HDFS作为数据源而Spark的发行包里虽然带了Hadoop的客户端库但在Windows上运行时会调用Hadoop的一些本地库比如winutils.exe和hadoop.dll这些库在Linux上是自带的Windows上需要手动补。如果你不补启动Spark的时候会看到这样的报错java.io.IOException: Could not locate executable null\bin\winutils.exe in the Hadoop binaries.这个报错不会让Spark直接崩溃但会导致一些功能异常比如写文件权限问题。所以这一步必须做。2.2 winutils.exe的获取与放置winutils.exe是Hadoop为Windows编译的工具集官方不直接提供需要从第三方编译版本获取。常见做法是找对应Hadoop版本的winutils包比如Hadoop 3.2.2对应的winutils。下载下来之后目录结构是这样的hadoop-3.2.2/ ├── bin/ │ ├── winutils.exe │ ├── hadoop.dll │ └── ...把这个hadoop-3.2.2文件夹放到一个没有空格和中文的路径下比如C:\hadoop。然后配置环境变量新建HADOOP_HOME值为C:\hadoop编辑Path添加%HADOOP_HOME%\bin还有一个关键操作把hadoop.dll复制到C:\Windows\System32目录下。这一步很多人会漏导致后面报UnsatisfiedLinkError。2.3 验证Hadoop依赖是否生效配置完之后新开命令行执行winutils.exe ls C:\如果能列出C盘的文件列表说明winutils配置成功了。如果提示找不到命令检查Path。这一步做完Windows下跑Spark的最大障碍就扫清了。我当初就是卡在这里整整一个下午因为网上教程都只说下载winutils放到bin目录但没说还要复制dll到System32。3. Spark本体的下载与配置3.1 版本选择与下载去Spark官网的下载页面选择Pre-built for Apache Hadoop 3.2 and later这个版本。为什么选这个而不是Pre-built with user-provided Hadoop因为前者已经打包好了Hadoop客户端省得你自己配。下载下来是个tgz压缩包Windows下用7-Zip或者WinRAR解压。解压路径同样不要有空格和中文我放在C:\spark\spark-3.2.4-bin-hadoop3.2。3.2 环境变量配置新建SPARK_HOME值为C:\spark\spark-3.2.4-bin-hadoop3.2编辑Path添加%SPARK_HOME%\bin3.3 spark-env.sh的Windows适配Spark的配置文件在conf目录下Linux下用的是spark-env.shWindows下需要改成spark-env.cmd。直接复制一份spark-env.sh.template重命名为spark-env.cmd然后编辑内容set JAVA_HOMEC:\Java\jdk1.8.0_301 set HADOOP_HOMEC:\hadoop set SPARK_MASTER_HOSTlocalhost set SPARK_LOCAL_IP127.0.0.1这里SPARK_MASTER_HOST和SPARK_LOCAL_IP是本地集群的关键配置不配的话启动时会尝试用主机名解析Windows下经常解析失败导致启动卡住。3.4 启动验证新开命令行执行spark-shell如果看到Spark的ASCII艺术字和scala提示符恭喜你Spark已经跑起来了。这时候你可以试试val data Seq(1,2,3,4,5) val rdd sc.parallelize(data) rdd.map(x x * 2).collect().foreach(println)能打印出2、4、6、8、10就说明一切正常。注意第一次启动spark-shell会比较慢因为要初始化SparkContext。如果卡在某个地方超过一分钟大概率是网络或者主机名解析问题检查SPARK_LOCAL_IP配置。4. 本地集群模式从单机到伪分布式4.1 本地模式和历史服务器模式的区别前面跑通的spark-shell其实是local模式也就是单JVM运行没有真正的Master和Worker进程。如果你想体验更接近生产环境的集群模式可以启动Spark自带的Standalone集群。Standalone集群需要启动两个进程Master和Worker。在Windows下启动这两个进程需要用到spark-class.cmd脚本。4.2 启动Master和Worker先启动Masterspark-class.cmd org.apache.spark.deploy.master.Master看到Master started之类的日志并且监听在7077端口说明Master起来了。然后新开一个命令行窗口启动Workerspark-class.cmd org.apache.spark.deploy.worker.Worker spark://localhost:7077Worker启动后会向Master注册你可以在Master的日志里看到注册成功的消息。这时候打开浏览器访问http://localhost:8080能看到Spark的Web UI显示一个Worker和它的资源情况。4.3 提交任务到Standalone集群集群起来之后可以用spark-submit提交任务spark-submit.cmd --master spark://localhost:7077 --class org.apache.spark.examples.SparkPi examples/jars/spark-examples_2.12-3.2.4.jar 10这个命令会计算Pi的近似值跑完会输出结果。如果能在Web UI里看到任务执行记录说明集群模式完全跑通了。4.4 内存配置的那些门道Standalone模式下Worker默认使用机器所有可用内存但会预留1GB给系统。你可以通过SPARK_WORKER_MEMORY环境变量来控制set SPARK_WORKER_MEMORY4g这个值设置多少合适我的经验是物理内存的60%到70%。比如你机器有16GB内存设成10g比较稳妥。设太大容易导致系统卡顿设太小任务跑不动。还有一个容易忽略的参数是SPARK_WORKER_CORES控制Worker能用的CPU核数。默认是全部核心但如果你还要用电脑干别的建议留一两个核心set SPARK_WORKER_CORES45. 实战演示用户复购率分析5.1 场景与数据准备光跑通环境没意思得用真实场景验证。这里用一个电商用户复购率分析的例子。复购率的定义是在统计周期内购买次数大于等于2次的用户占总用户的比例。准备一份订单数据CSV字段包括用户ID、订单ID、订单日期、订单金额。数据可以自己造也可以用公开数据集。我一般用Python快速生成一份测试数据import csv import random from datetime import datetime, timedelta users [fU{i:05d} for i in range(1, 1001)] start_date datetime(2024, 1, 1) with open(orders.csv, w, newline, encodingutf-8) as f: writer csv.writer(f) writer.writerow([user_id, order_id, order_date, amount]) order_id 1 for user in users: order_count random.choices([1, 2, 3, 4, 5], weights[50, 25, 15, 7, 3])[0] for _ in range(order_count): date start_date timedelta(daysrandom.randint(0, 180)) amount round(random.uniform(10, 500), 2) writer.writerow([user, fO{order_id:06d}, date.strftime(%Y-%m-%d), amount]) order_id 1这份数据模拟了1000个用户其中大约一半只买了一次另一半有复购行为。5.2 用Spark SQL计算复购率启动spark-shell把数据读进来val df spark.read.option(header, true).option(inferSchema, true).csv(C:\\data\\orders.csv) df.createOrReplaceTempView(orders)然后用SQL计算复购率val repurchaseRate spark.sql( SELECT COUNT(DISTINCT user_id) AS total_users, COUNT(DISTINCT CASE WHEN order_count 2 THEN user_id END) AS repurchase_users, ROUND(COUNT(DISTINCT CASE WHEN order_count 2 THEN user_id END) * 100.0 / COUNT(DISTINCT user_id), 2) AS repurchase_rate FROM ( SELECT user_id, COUNT(order_id) AS order_count FROM orders GROUP BY user_id ) ) repurchaseRate.show()跑出来的结果应该显示总用户数1000复购用户数大约500复购率50%左右。5.3 按月份分析复购趋势更进一步看看每个月的复购率变化val monthlyTrend spark.sql( SELECT DATE_FORMAT(order_date, yyyy-MM) AS month, COUNT(DISTINCT user_id) AS monthly_users, COUNT(DISTINCT CASE WHEN order_count 2 THEN user_id END) AS repurchase_users, ROUND(COUNT(DISTINCT CASE WHEN order_count 2 THEN user_id END) * 100.0 / COUNT(DISTINCT user_id), 2) AS repurchase_rate FROM ( SELECT user_id, order_date, COUNT(order_id) OVER (PARTITION BY user_id) AS order_count FROM orders ) GROUP BY DATE_FORMAT(order_date, yyyy-MM) ORDER BY month ) monthlyTrend.show()这里用到了窗口函数COUNT OVER这是Spark SQL里非常实用的功能可以在不聚合的情况下给每行附加统计信息。5.4 把结果写成Parquet格式分析完的结果通常要落地存储。Parquet是Spark生态里最推荐的列式存储格式压缩比高、查询快repurchaseRate.write.mode(overwrite).parquet(C:\\data\\output\\repurchase_rate) monthlyTrend.write.mode(overwrite).parquet(C:\\data\\output\\monthly_trend)写完之后可以读回来验证val check spark.read.parquet(C:\\data\\output\\repurchase_rate) check.show()6. 那些教程不会告诉你的坑6.1 路径问题空格和中文是万恶之源Windows用户习惯把东西放在桌面或者我的文档里这些路径都带中文。Spark和Hadoop对中文路径的支持很差经常出现文件找不到或者乱码。所有涉及Spark的路径一律用纯英文、无空格的路径比如C:\spark、C:\data。6.2 端口冲突8080被占用的处理Spark Master的Web UI默认用8080端口但这个端口太常用了Tomcat、Jenkins、各种开发工具都爱用。如果启动Master时报端口被占用改spark-env.cmdset SPARK_MASTER_WEBUI_PORT8081Worker的Web UI端口是8081如果也冲突同样可以改。6.3 防火墙拦截Worker注册不上Windows防火墙有时候会拦截Master和Worker之间的通信导致Worker启动后一直注册不上。如果遇到这种情况临时关闭防火墙测试一下确认是防火墙问题后给Java进程添加例外规则。6.4 内存溢出OOM的排查思路跑大一点的数据集时可能会遇到OutOfMemoryError。排查思路是先看是Driver OOM还是Executor OOM日志里会写Driver OOM的话调大spark.driver.memoryExecutor OOM的话调大spark.executor.memory或者增加分区数让每个分区数据量小一点在spark-shell里可以这样设置spark.conf.set(spark.sql.shuffle.partitions, 200)默认的200个分区对于小数据量来说太多了可以调小到10或者20减少任务调度开销。6.5 日志太吵调整日志级别Spark默认会打印大量INFO日志刷屏严重。在conf目录下创建log4j2.properties文件内容rootLogger.level WARN rootLogger.appenderRef.stdout.ref console这样只打印WARN及以上级别的日志清爽很多。7. 从本地集群到生产环境的距离本地集群跑通之后你可能会想这套东西能直接上生产吗答案是能但需要改不少东西。本地Standalone集群最大的问题是单点故障——Master挂了整个集群就废了。生产环境通常用YARN或者Kubernetes做资源调度这两个在Windows上跑起来都比较费劲所以真实生产环境基本还是Linux。但Windows本地集群的价值在于开发调试。你可以在Windows上写好Spark任务本地测试通过后再提交到Linux集群跑。这样开发效率比直接连远程集群高得多也不用担心把生产环境搞崩。如果要在Windows上做更接近生产的开发可以考虑用Docker跑一个Linux容器里面装完整的Spark集群。这样既保留了Windows的开发便利又有了Linux的运行环境。不过这是另一个话题了先把本地集群玩明白再说。8. 几个提升开发效率的小技巧用spark-shell交互式开发虽然方便但代码写多了不好管理。我一般用IDEA或者VS Code写Scala脚本然后用spark-submit提交。这样代码可以版本管理也方便复用。另外spark-shell支持--jars参数加载外部依赖比如MySQL驱动spark-shell.cmd --jars C:\drivers\mysql-connector-java-8.0.28.jar加载后就可以直接读写MySQLval jdbcDF spark.read.format(jdbc) .option(url, jdbc:mysql://localhost:3306/test) .option(dbtable, orders) .option(user, root) .option(password, 123456) .load()还有一个实用技巧把常用的配置写进spark-defaults.conf这样每次启动都自动生效不用重复设置。比如spark.sql.shuffle.partitions 20 spark.driver.memory 2g spark.sql.adaptive.enabled truespark.sql.adaptive.enabled是Spark 3.0引入的自适应查询执行能根据运行时统计信息动态调整执行计划对性能提升很明显建议默认开启。我在实际使用中最大的体会是Windows下玩Spark环境配置占了80%的时间真正写代码反而很快。所以第一次搭建的时候耐心把环境搞扎实后面就一劳永逸了。另外遇到报错不要慌先看日志里的Caused by那才是根因表面的报错信息往往只是连锁反应。
返回列表