ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Spark集群部署与实战:环境准备、ETL与调优避坑指南

Spark集群部署与实战:环境准备、ETL与调优避坑指南 Spark这套东西我在生产环境里从1.x版本一路用到3.x装过的集群少说也有几十套单机伪分布式、Standalone、YARN上跑的都折腾过。最开始的动机特别朴素——用Hadoop MapReduce跑一个简单的日志统计要等好几分钟磁盘反复读写中间还得盯着进度条那体验确实折磨人。后来换到Spark同样的活儿代码量少一半跑起来快一个数量级内存计算这四个字才算是真正理解了。这篇内容适合谁看刚接触大数据、想把Spark跑起来的新手需要给团队搭一套测试环境的开发者还有那些面试前想把集群部署这块补一补的朋友。我会从环境准备讲到集群搭建从spark-shell交互讲到写ETL脚本把版本匹配、内存分配、常见报错这些实战细节都摊开来说。照着我踩过的坑走你大概率能一次跑通不用像我当年那样对着一堆日志抓瞎。1. 先搞明白Spark凭什么值得学别上来就闷头装很多人学Spark的第一步是去官网下载压缩包结果解压完发现跑不起来JDK版本不对、Scala不匹配、Hadoop依赖缺失一路报错。我的建议是先花十分钟把Why想清楚后面装的时候心里有数遇到问题也更容易定位。1.1 Spark到底解决了MapReduce的哪些痛点MapReduce的模型其实很简单Map阶段处理Shuffle阶段落盘Reduce阶段汇总。问题就出在这个落盘上——每一次MapReduce作业的中间结果都要写到HDFS下一次作业再从磁盘读回来。迭代式计算比如机器学习里的梯度下降要跑几十上百轮每轮都落盘读盘I/O开销直接把性能拖垮。Spark的核心改进是把中间结果尽可能留在内存里。它的抽象是RDDResilient Distributed Dataset弹性分布式数据集数据在多个节点上分区存储算子之间形成DAG有向无环图Spark会根据依赖关系把任务切分成Stage能在内存中pipeline执行的就不落盘。这里要说一个很多人误解的点Spark并不是全部在内存。Shuffle操作、缓存溢出、内存不够时的spill都会落盘。它的优势是尽量不落盘而不是绝对不落盘。理解这一点你后面调内存参数的时候才不会想当然。还有一个优点是API友好。同样一个按key聚合的逻辑MapReduce要写Mapper、Reducer、Driver三个类Spark用Scala一行就能表达。开发效率的提升对团队来说是实打实的。1.2 和Hadoop生态的关系到底是替代还是共存经常有人问我学了Spark是不是就不用学Hadoop了这个问题得分开看。Spark本身不负责存储。它需要底层的存储和资源调度。存储这块绝大多数场景用的是HDFS也可以用S3、OSS这类对象存储资源调度这块可以是Spark自带的Standalone模式也可以是YARN或者Kubernetes。所以在典型的Hadoop生态里Spark是作为计算引擎存在的HDFS负责存YARN负责调度Spark负责算三者分工。实际项目里的组合通常是这样的组件角色常见替代HDFS分布式存储S3、OSS、HBaseYARN资源调度Standalone、KubernetesSpark计算引擎MapReduce、FlinkHiveSQL数仓Spark SQLHBase列式存储Cassandra这个概念理清楚你就知道装Spark的时候为什么经常要绑定一个Hadoop版本了。2. 安装前的环境准备与版本选择这一步错后面全白搭我记得第一次装Spark踩的最大的坑就是版本不匹配。当时电脑上装的是JDK 11Spark 2.4死活起不来报了一堆看不懂的反射错误折腾了一下午才发现Spark 2.4对JDK 11支持不完善。所以这一章我会把版本这块说透。2.1 JDK、Scala、Hadoop、Spark四个版本的匹配关系先把结论放前面新手直接照这个表选Spark版本推荐JDK推荐Scala对应Hadoop2.4.xJDK 82.11 / 2.122.6 / 2.73.0 - 3.2JDK 8 / 112.123.2 / 3.33.3 - 3.5JDK 8 / 11 / 172.12 / 2.133.3 / 3.4为什么JDK版本这么关键Spark用到了大量反射和Unsafe APIJDK 9之后模块系统改了很多内部类被封装老版本Spark访问不到就会报IllegalAccessError。这也是我用JDK 8踩坑最多、也最省心的原因——兼容性最好社区资料最多。Scala版本也需要注意。Spark是用Scala写的编译时Scala的主版本必须对上。2.11和2.12编译出的字节码不兼容混用会报NoSuchMethodError。提示如果你只是想在本机学一学又不想折腾版本直接下载Spark官网预编译好的Pre-built for Apache Hadoop 3.3 and later包里面已经把Hadoop客户端依赖打包进去了省心。2.2 Standalone、YARN、Kubernetes三种部署模式怎么选部署模式的选择取决于你的场景Local模式单机跑用来学习和调试。不用起集群进程就在本地速度最快。local[*]表示用所有CPU核心。Standalone模式Spark自带的集群管理器部署最简单适合中小规模集群和测试环境。缺点是资源调度能力比起YARN弱一些。YARN模式跑在Hadoop集群上和Hive、MapReduce共享资源是企业生产环境的主流选择。Kubernetes模式容器化部署适合已有K8s基础设施的团队弹性扩缩容方便。新手我建议从Local开始跑通了再上Standalone最后再考虑YARN。一步到位直接上YARN中间出问题你都不知道是哪一层的问题。2.3 硬件和操作系统的准备清单具体到机器我给个最低配置参考单机学习内存8G起步最好16G。Spark自己就要占几个G浏览器再开一堆标签页很容易把机器拖卡。三节点测试集群每台4核8G起步共12核24G。可以跑中等规模的数据测试。生产环境这个得按数据量和并发任务算后面调优章节我会讲怎么估算。操作系统我推荐Ubuntu或者CentOS这类主流的Linux发行版。Windows上跑Spark能跑但是各种路径分隔符、权限问题会让你怀疑人生强烈建议用Linux虚拟机或者WSL。3. Standalone集群安装实操一步一步来选Standalone模式做演示是因为它最能让你看清Spark集群的组成一个Master节点负责调度多个Worker节点负责执行。理解了这套机制后面理解YARN模式就是换个调度器而已。3.1 三台机器的基础环境配置假设我们有三台机器主机名分别是node1、node2、node3我准备一个表格列一下角色分配主机名角色说明node1Master Worker主节点也兼职算力node2Worker工作节点node3Worker工作节点第一步是配置主机名和hosts映射让三台机器能互相用名字访问# 在每台机器上执行修改主机名 hostnamectl set-hostname node1 # 其余两台改成node2、node3 # 编辑hosts文件三台机器都要加 vim /etc/hosts在hosts文件里追加192.168.1.101 node1 192.168.1.102 node2 192.168.1.103 node3然后是SSH免密登录。Spark的启动脚本要靠SSH去其他节点拉起进程不配免密每台都要输密码脚本直接卡住# 在node1上生成密钥一路回车 ssh-keygen -t rsa # 把公钥拷到三台机器包括自己 ssh-copy-id node1 ssh-copy-id node2 ssh-copy-id node3 # 验证 ssh node2 date能直接返回时间就说明免密通了。再装JDK# 下载的jdk解压到/opt目录 tar -zxvf jdk-8u371-linux-x64.tar.gz -C /opt/配置环境变量编辑/etc/profileexport JAVA_HOME/opt/jdk1.8.0_371 export PATH$JAVA_HOME/bin:$PATH执行source /etc/profile生效然后java -version验证。3.2 Spark解压与核心配置文件详解去Spark官网下载预编译包我这里是spark-3.5.0-bin-hadoop3。解压tar -zxvf spark-3.5.0-bin-hadoop3.tgz -C /opt/ cd /opt/spark-3.5.0-bin-hadoop3/conf复制模板文件cp spark-env.sh.template spark-env.sh cp workers.template workers编辑spark-env.sh这是最关键的文件export JAVA_HOME/opt/jdk1.8.0_371 export SPARK_MASTER_HOSTnode1 export SPARK_MASTER_PORT7077 export SPARK_MASTER_WEBUI_PORT8080 export SPARK_WORKER_CORES2 export SPARK_WORKER_MEMORY4g export SPARK_WORKER_INSTANCES1逐条解释一下为什么这么配。SPARK_MASTER_HOST绑定Master所在机器Worker启动时会按这个地址去注册SPARK_WORKER_CORES是每个Worker能用的CPU核数设成机器实际核心数别设满给系统留点SPARK_WORKER_MEMORY同理4g是保守值机器8G内存的话建议留一半给系统。编辑workers文件写所有Worker节点的主机名node1 node2 node33.3 启动集群与验证在node1上执行启动脚本/opt/spark-3.5.0-bin-hadoop3/sbin/start-all.sh如果脚本报JAVA_HOME is not set说明你SSH到其他节点时环境变量没加载把所有节点的/etc/profile都配好即可。启动完成后用jps看进程# node1上应该看到 Master Worker # node2、node3上应该看到 Worker访问http://node1:8080能看到Spark的Web UI上面列出了所有Worker和它们的状态、核心数、内存。看到三个ALIVE的Worker集群就算搭好了。然后提交一个测试任务验证/opt/spark-3.5.0-bin-hadoop3/bin/spark-submit \ --class org.apache.spark.examples.SparkPi \ --master spark://node1:7077 \ --executor-memory 1g \ --total-executor-cores 2 \ /opt/spark-3.5.0-bin-hadoop3/examples/jars/spark-examples_2.12-3.5.0.jar \ 100任务跑完会打印出Pi的近似值。看到结果集群就真的活了。4. Spark Core与DataFrame上手实战从API到ETL集群搭好只是开始真正干活要写代码。Spark的编程模型这几年变化挺大从最早的RDD到后来的DataFrame、Dataset再到Spark SQL直接写SQL。这一章我带你把几种API都过一遍。4.1 spark-shell交互式环境快速上手spark-shell是学习Spark最好的工具边敲边看结果非常适合调试。/opt/spark-3.5.0-bin-hadoop3/bin/spark-shell --master local[2]启动后有个sc对象就是SparkContextspark对象是SparkSession。先来个WordCount热热身val data List(hello spark, hello world, spark is fast) val rdd sc.parallelize(data) val counts rdd .flatMap(_.split( )) .map(word (word, 1)) .reduceByKey(_ _) counts.collect().foreach(println)这段代码展示了RDD的核心算子flatMap把每行拆成单词map变成键值对reduceByKey按key聚合。collect()是把分布式结果拉回本地注意——数据量大时千万别乱用collect会把Driver内存打爆。4.2 DataFrame与Spark SQL的实操案例RDD的写法偏底层实际项目里现在主流是DataFrame SQL。举个读CSV做统计的例子val df spark.read .option(header, true) .option(inferSchema, true) .csv(/data/orders.csv) df.createOrReplaceTempView(orders) val result spark.sql( SELECT city, COUNT(*) AS order_cnt, SUM(amount) AS total_amount FROM orders WHERE dt 2024-01-01 GROUP BY city ORDER BY total_amount DESC ) result.show()这里有个性能上的点DataFrame走的是Catalyst优化器和Tungsten执行引擎能用上列式存储、谓词下推、代码生成等优化。同样一个聚合DataFrame通常比手写RDD快除非你对RDD的分区策略有非常精细的控制需求。我实测过一个2000万行的订单表做分组聚合DataFrame比同逻辑的RDD版本快了大概30%到40%。所以能用DataFrame就用DataFrameRDD留给需要精细控制的场景。4.3 写一个完整的ETL脚本实战里最常见的任务就是ETL从源读数据清洗转换写到目标。下面这个脚本从Hive表读数据清洗后写回我用注释标了每一步的意图from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, to_date spark SparkSession.builder \ .appName(OrderETL) \ .config(spark.sql.shuffle.partitions, 200) \ .enableHiveSupport() \ .getOrCreate() # 1. 读取源数据 df spark.sql(SELECT * FROM ods.orders WHERE dt 2024-01-01) # 2. 清洗过滤空值处理异常金额 clean df.filter(col(order_id).isNotNull()) \ .withColumn(amount, when(col(amount) 0, 0).otherwise(col(amount))) # 3. 转换日期格式化拆出年月 result clean.withColumn(order_date, to_date(col(create_time))) # 4. 写出按日期分区 result.write.mode(overwrite) \ .partitionBy(dt) \ .saveAsTable(dwd.orders_clean) spark.stop()spark.sql.shuffle.partitions这个参数我再强调一遍默认200小数据量时会生成一堆空任务拖慢速度大数据量时又可能分区太少导致单任务处理过量。经验值是根据数据量算每个分区处理128MB到200MB比较合适比如10G的Shuffle数据200个左右分区是合理的。提交脚本用spark-submitspark-submit \ --master spark://node1:7077 \ --deploy-mode client \ --executor-memory 2g \ --executor-cores 2 \ --num-executors 3 \ order_etl.py5. 常见问题排查与踩坑经验实录装Spark和用Spark的过程中报错是常态。这一章我把这些年攒下来的典型问题和解决方法整理出来。5.1 启动失败、内存溢出、连接超时的排查思路先给一张速查表遇到问题对着找报错关键词常见原因解决方向JAVA_HOME is not set环境变量未加载检查所有节点的/etc/profile并sourceConnection refusedMaster地址或端口错核对spark-env.sh里的MASTER_HOSTOutOfMemoryErrorExecutor内存不足调大executor-memory或减少分区数据量FetchFailedExceptionShuffle读到失败调大spark.reducer.maxSizeInFlight检查网络ClassNotFoundException依赖包没打进去用--jars或--packages指定依赖Container killed by YARN内存超限调大yarn.nodemanager.resource.memory-mb拿内存溢出这个最常见的说。很多人一看OOM就无脑加内存其实得先判断是哪块内存炸了。如果是Driver OOM多半是你把大结果collect回本地了如果是Executor OOM可能是某个分区数据倾斜一个任务处理了远大于平均的数据量。数据倾斜的典型表现是大部分任务秒完就剩一两个卡在那里不动。数据倾斜的处理思路有几种给倾斜的key加随机前缀打散、用map-side join规避Shuffle、开启自适应执行spark.sql.adaptive.enabledtrue让Spark自动处理。我通常先开AQE看能不能自动缓解不行再手动加盐。5.2 内存与并行度调优的参数估算调优这块我总结了一套估算流程。假设集群总共300G内存、120核要跑一个大作业Executor数量先确定每个Executor多少内存。经验值是单个Executor不超过64G超过的话JVM的GC压力会很大。我们取每个Executor 20G那么300 / 20 15个Executor。每个Executor的核心数一般是4到8个取5的话15 × 5 75核够用。并行度计算总核心数75每个核心可以跑2到3个任务那么并行度大概是150到225个。这就是spark.sql.shuffle.partitions该设的值。内存结构也得了解。一个Executor的内存分三块执行内存execution、存储内存storage、用户内存user。执行和存储可以互相借用通过spark.memory.fraction控制比例默认0.6。如果你的作业缓存RDD多可以把存储那块调大一点。提示调参数之前先看Spark UI在Executors页签看GC时间占比。如果GC时间超过任务时间的10%说明内存不够先加内存再调其他参数。5.3 我在生产环境踩过的三个坑第一个坑是序列化。早期我用Java序列化一个对象序列化后特别大Shuffle数据量吓人。后来换成Kryo序列化配置spark.serializerorg.apache.spark.serializer.KryoSerializer并注册自定义类Shuffle数据量直接降了一截作业快了不少。第二个坑是数据本地性。有一次提交任务每个task都在跨节点拉数据慢得离谱。查下来是计算和数据不在同一台机器。Spark的调度有本地性级别PROCESS_LOCAL最快NODE_LOCAL次之RACK_LOCAL再差ANY最差。如果UI上显示大量任务在RACK_LOCAL或ANY级别说明你的数据分布和任务分配不匹配可以考虑调整分区或增加等待时间让Spark调度到更好位置的数据。第三个坑是小文件问题。HDFS上如果有一堆几十KB的小文件每个文件对应一个task光启动任务的调度开销就占大头。解决方法是合并或者用Spark读取后用coalesce/repartition重新分区。我一般读进来先repartition(合适的分区数)再往下走。6. 从学完到用起来我给新手的几条实在建议装完、跑通、能写脚本这只是入门。真正把Spark用起来还需要在真实的业务场景里反复练。我自己带新人的经验是别一上来就啃源码或者啃难懂的内核书。先把官方文档里的Quick Start和Programming Guide过一遍然后把Spark自带的examples跑一遍理解每种算子会产生什么Shuffle。学习路线上我是这么走的先掌握Spark Core的RDD操作理解DAG、Stage、Shuffle这些概念然后转到DataFrame和Spark SQL这是现在的生产力工具最后再看调优和源码。面试里问得最多的是数据倾斜怎么解决、Shuffle原理、内存模型这些都在前面几章覆盖到了。还有一个习惯我觉得特别有用每次调参数之前先看Spark UI把执行计划、Stage耗时、shuffle读写量都摸清楚再动手。凭感觉瞎调参数十次有九次是无效甚至负优化。UI上每一个数字背后都对应着作业真实的行为学会读UI比背参数值重要得多。资源上官方文档永远是最准的。中文博客和教程可以作为理解入门但版本更新后很多细节会变遇到对不上的地方以官方为准。跑数据的路上没有捷径多跑几个真实的作业比看十篇文章都有用。
返回列表