
Spark Standalone 模式Apache Spark 内置集群管理器的部署、配置与高可用实战指南【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/sparkSpark Standalone 是 Apache Spark 自带的轻量级集群管理器无需依赖 YARN 或 Kubernetes 即可让 Spark 应用跑在自建的 Master/Worker 集群上尤其适合中小规模集群、开发测试环境以及希望快速上手 Spark 集群的用户。本文以仓库 docs/spark-standalone.md 为骨架结合sbin/启动脚本与 Master 源码 的实现细节完整讲解集群的安装、手动与脚本化启动、全部 Master/Worker 配置项、应用提交spark-submit 与 REST API、资源调度、监控日志以及 ZooKeeper/本地文件系统两种高可用方案。架构总览Standalone 在 Spark 部署体系中的位置Spark 支持多种集群管理器。除了 YARN 集群管理器 和 Kubernetes 之外Spark 还提供了一个极简的 standalone 部署模式由一个 Master 进程负责调度决策多个 Worker 进程负责提供计算资源应用通过spark://HOST:PORT这样的 URL 接入集群。上图展示了 Spark 应用运行时的核心组件关系Driver Program 通过 SparkContext 向 Cluster Manager 申请资源Cluster Manager 将 Worker Node 上的资源以 Executor 的形式分配出来在 Standalone 模式下这个 Cluster Manager 角色就由 Master 进程扮演。你可以手动逐个启动 Master 和 Worker也可以使用仓库sbin/目录下提供的启动脚本一键拉起整个集群这些守护进程也可以全部跑在同一台机器上用于测试。从源码结构看Standalone 的 Master 实现在 core/src/main/scala/org/apache/spark/deploy/master/Master.scalaWorker 实现在同目录的 worker 包中二者通过 Spark 的 RPC 框架通信。Master 内部维护了workers、apps、waitingApps、waitingDrivers等集合分别追踪已注册的工作节点、已注册的应用、等待调度的应用和驱动。安全须知Standalone 模式默认不启用身份认证等安全特性。当集群暴露在公网或不可信网络中时必须采取措施保护集群访问防止未经授权的应用在集群上运行。部署前请务必阅读 Spark Security 以及本文后续的「配置网络安全端口」小节。安装 Spark Standalone 到集群安装 Standalone 模式非常简单只需在每个节点上放置一份编译好的 Spark 发行包即可。你可以直接使用每个版本随附的预编译包也可以参考 Building Spark 自行构建。手动启动集群启动 Master在集群中的一台机器上执行./sbin/start-master.sh启动成功后Master 会打印出自己的spark://HOST:PORTURL这个 URL 有两个用途一是让 Worker 连接它二是作为SparkContext的master参数传给应用。你也可以在 Master 的 Web UI 上找到这个 URL默认地址是http://localhost:8080。从 sbin/start-master.sh 的源码可以看到脚本在未显式设置环境变量时会使用这些默认值端口默认7077SPARK_MASTER_PORT为空时Web UI 端口默认8080SPARK_MASTER_WEBUI_PORT为空时绑定主机名默认取hostname -f的结果Solaris 下用check-hostname实际启动的进程类为org.apache.spark.deploy.master.Master脚本注释特别说明这个类名会被下游SparkSubmit匹配不能随意改动。启动 Worker 并接入 Master在一个或多个节点上执行./sbin/start-worker.sh master-spark-URL启动 Worker 后回到 Master 的 Web UI默认http://localhost:8080应当能看到新节点出现在列表中并显示它的 CPU 核数与内存内存默认扣掉 1 GiB 留给操作系统。对应的启动脚本是 sbin/start-worker.sh其中进程类为org.apache.spark.deploy.worker.Worker并支持通过SPARK_WORKER_INSTANCES环境变量在一台机器上启动多个 Worker 实例该变量自 Spark 3.0 起被标记为已废弃。Master/Worker 通用命令行参数以下命令行参数可同时传给 Master 和 Worker带仅 Worker标注的除外参数含义-h HOST,--host HOST监听的主机名-p PORT,--port PORT服务监听端口默认Master 为 7077Worker 为随机端口--webui-port PORTWeb UI 端口默认Master 为 8080Worker 为 8081-c CORES,--cores CORES允许 Spark 应用使用的机器总 CPU 核数默认全部可用核仅 Worker-m MEM,--memory MEM允许 Spark 应用使用的机器总内存格式如1000M或2G默认机器总内存减 1 GiB仅 Worker-d DIR,--work-dir DIR用于临时空间与作业输出日志的目录默认SPARK_HOME/work仅 Worker--properties-file FILE自定义 Spark 属性文件的路径默认conf/spark-defaults.conf集群启动脚本conf/workers 文件要使用脚本一键拉起 Standalone 集群需要在 Spark 目录下创建conf/workers文件逐行列出所有需要启动 Worker 的机器主机名可参考 conf/workers.template模板默认内容为localhost。如果conf/workers不存在脚本默认只在单机localhost上启动 Worker便于测试。需要注意Master 机器需要通过 ssh 访问每台 Worker 机器。默认情况下 ssh 是并行执行的因此要求配置免密登录基于私钥。如果没有免密配置可以设置环境变量SPARK_SSH_FOREGROUND改为串行方式逐个输入密码。sbin 下的全部脚本这些脚本基于 Hadoop 的部署脚本编写位于SPARK_HOME/sbin完整的启停能力如下脚本作用sbin/start-master.sh在脚本所在机器上启动一个 Master 实例sbin/start-workers.sh在conf/workers中列出的每台机器上启动 Workersbin/start-worker.sh在脚本所在机器上启动一个 Worker 实例sbin/start-connect-server.sh启动 Spark Connect 服务器sbin/start-history-server.sh启动 History Server用于查看已完成应用的日志需要开启事件日志见 Monitoring and Instrumentationsbin/start-all.sh同时启动一个 Master 和一批 Worker行为同上sbin/stop-master.sh停止由start-master.sh启动的 Mastersbin/stop-worker.sh停止脚本所在机器上的所有 Worker 实例sbin/stop-workers.sh停止conf/workers中列出的所有机器上的 Workersbin/stop-connect-server.sh停止所有 Spark Connect 服务器实例sbin/stop-history-server.sh停止 History Serversbin/stop-all.sh停止 Master 和所有 Workersbin/decommission-worker.sh优雅下线一个 Worker让正在运行的任务完成、shuffle 数据迁移完毕后再退出注意这些脚本必须在你想运行 Spark Master 的那台机器上执行而不是在本地机器上。从 sbin/start-workers.sh 可以看到它实际上是通过workers.sh远程执行start-worker.sh来逐个拉起 WorkerMaster 地址取自SPARK_MASTER_HOST:SPARK_MASTER_PORT默认spark://hostname -f:7077。Windows 说明启动脚本目前不支持 Windows。要在 Windows 上运行 Spark 集群请手动逐个启动 Master 和 Worker。conf/spark-env.sh 环境变量可以通过在conf/spark-env.sh中设置环境变量进一步配置集群。先复制conf/spark-env.sh.template创建该文件模板中已经逐项注释了各变量的含义见 conf/spark-env.sh.template然后复制到所有 Worker 机器上才能使设置生效。可用设置如下环境变量含义SPARK_MASTER_HOST将 Master 绑定到特定的主机名或 IP 地址例如公网地址SPARK_MASTER_PORT让 Master 使用不同端口默认7077SPARK_MASTER_WEBUI_PORTMaster Web UI 端口默认8080SPARK_MASTER_OPTS仅作用于 Master 的配置属性形如-Dxy默认无。可选属性见下节SPARK_LOCAL_DIRSSpark 的临时空间目录包括 map 输出文件和落盘 RDD。应放在系统的快速本地磁盘上也可以是跨多块磁盘的逗号分隔目录列表SPARK_LOG_DIR日志文件存放位置默认SPARK_HOME/logsSPARK_LOG_MAX_FILES最大日志文件数量默认5SPARK_PID_DIRpid 文件存放位置默认/tmpSPARK_WORKER_CORES允许 Spark 应用使用的机器总核数默认全部可用核SPARK_WORKER_MEMORY允许 Spark 应用使用的机器总内存如1000m、2g默认总内存减 1 GiB。注意每个应用个体的内存由它自己的spark.executor.memory属性配置SPARK_WORKER_PORTWorker 监听端口默认随机SPARK_WORKER_WEBUI_PORTWorker Web UI 端口默认8081SPARK_WORKER_DIR运行应用的目录包含日志与临时空间默认SPARK_HOME/workSPARK_WORKER_OPTS仅作用于 Worker 的配置属性形如-Dxy默认无。可选属性见下文SPARK_DAEMON_MEMORYMaster 和 Worker 守护进程自身分配的内存默认1gSPARK_DAEMON_JAVA_OPTSMaster 和 Worker 守护进程自身的 JVM 选项形如-Dxy默认无SPARK_DAEMON_CLASSPATHMaster 和 Worker 守护进程自身的 classpath默认无SPARK_PUBLIC_DNSSpark Master 和 Worker 的公网 DNS 名称默认无SPARK_MASTER_OPTS 支持的系统属性属性名默认值含义引入版本spark.master.ui.port8080Master Web UI 端点的端口号1.1.0spark.master.ui.title无Master UI 页面标题未设置时默认使用Spark Master at master url4.0.0spark.master.ui.decommission.allow.modeLOCALMaster Web UI 的/workers/kill端点行为LOCAL表示仅允许 Master 所在机器的本机 IP 调用DENY表示完全禁用该端点ALLOW表示允许任意 IP 调用3.1.0spark.master.ui.historyServerUrl无Spark History Server 的运行 URL注意它假定所有 Spark 作业共享 History Server 访问的同一事件日志位置4.0.0spark.master.rest.enabledtrue是否启用 Master REST API 端点1.3.0spark.master.rest.host无Master REST API 端点的主机4.0.0spark.master.rest.port6066Master REST API 端点的端口号1.3.0spark.master.rest.filters无应用于 Master REST API 的过滤器类名列表逗号分隔4.0.0spark.master.useAppNameAsAppId.enabledfalse实验性为 true 时Spark Master 使用用户提供的 appName 作为 appId4.0.0spark.deploy.retainedApplications200UI 中最多展示的已完成应用数更早的应用会被丢弃以维持该上限0.8.0spark.deploy.retainedDrivers200UI 中最多展示的已完成驱动数更早的驱动会被丢弃1.1.0spark.deploy.spreadOutDriverstrueStandalone 集群管理器是否将驱动分散到各节点还是尽量集中到尽可能少的节点上。分散通常更利于 HDFS 数据本地性集中对计算密集型负载更高效4.0.0spark.deploy.spreadOutAppstrue同上针对应用Executor 分配的分散/集中策略0.6.1spark.deploy.defaultCoresInt.MaxValue应用未设置spark.cores.max时默认获得的核数。若不设置应用默认获得全部可用核。在共享集群上建议调低防止用户默认抢占整个集群0.9.0spark.deploy.maxExecutorRetries10连续 Executor 失败次数上限超过后 Standalone 集群管理器移除故障应用。只要还有运行中的 Executor应用就不会被移除若应用连续失败超过该次数、期间没有 Executor 成功启动且当前无运行中 Executor则会被标记为失败并移除。设为-1可禁用自动移除1.6.3spark.deploy.maxDriversInt.MaxValue最多可运行的驱动数量4.0.0spark.deploy.appNumberModulo无应用编号取模。默认情况下app-yyyyMMddHHmmss-9999的下一个编号是app-yyyyMMddHHmmss-10000若取模为 10000则会变为app-yyyyMMddHHmmss-0000。大多数情况下前缀app-yyyyMMddHHmmss在创建 10000 个应用的过程中已经变化4.0.0spark.deploy.driverIdPatterndriver-%s-%04d基于 JavaString.format的驱动 ID 生成模式如driver-20231031224459-0019。请谨慎确保生成的 ID 唯一4.0.0spark.deploy.appIdPatternapp-%s-%04d基于 JavaString.format的应用 ID 生成模式如app-20231031224509-0008。请谨慎确保生成的 ID 唯一4.0.0spark.worker.timeout60若 Master 在多少秒内未收到 Worker 心跳即认为该 Worker 丢失0.6.2spark.dead.worker.persistence15UI 中保留已死亡 Worker 信息的迭代次数。默认情况下死亡 Worker 从最后一次心跳起可见(15 1) * spark.worker.timeout秒0.8.0spark.worker.resource.{name}.amount无Worker 上某类资源的使用量3.0.0spark.worker.resource.{name}.discoveryScript无资源发现脚本路径Worker 启动时用它查找特定资源脚本输出格式应与ResourceInformation类一致3.0.0spark.worker.resourcesFile无资源文件路径Worker 启动时用它查找各类资源。文件内容格式如[{id:{componentName: spark.worker, resourceName:gpu}, addresses:[0,1,2]}]。若某资源未在资源文件中找到则回退到发现脚本若发现脚本也找不到Worker 将启动失败3.0.0以上配置在 Master 源码中有直接对应。例如 Master.scala 中Master 启动时即读取了driverIdPattern、appIdPattern、workerTimeoutMs、retainedApplications、retainedDrivers、maxDrivers、recoveryMode、maxExecutorRetries等配置而restServerEnabled与restServerBoundPort则对应 REST API 的启用与端口绑定逻辑见 Master.scala。SPARK_WORKER_OPTS 支持的系统属性属性名默认值含义引入版本spark.worker.initialRegistrationRetries6以短间隔5 到 15 秒重连注册的重试次数4.0.0spark.worker.maxRegistrationRetries16最大重连次数。超过initialRegistrationRetries后重连间隔变为 30 到 90 秒4.0.0spark.worker.cleanup.enabledtrue周期性清理 Worker/应用目录。仅影响 Standalone 模式YARN 机制不同且只清理已停止应用的目录。若spark.shuffle.service.db.enabled为 true建议启用此配置1.0.0spark.worker.cleanup.interval180030 分钟Worker 清理本地旧应用工作目录的时间间隔秒1.0.0spark.worker.cleanup.appDataTtl6048007 天每个 Worker 上保留应用工作目录的秒数TTL。应用日志和 jar 都会下载到应用工作目录作业频繁时磁盘空间会迅速被占满TTL 应根据可用磁盘空间设置1.0.0spark.shuffle.service.db.enabledtrue将外部 Shuffle 服务状态存到本地磁盘重启外部 shuffle 服务时可自动重载当前 Executor 信息。仅影响 Standalone 模式YARN 始终开启此行为。建议同时启用spark.worker.cleanup.enabled确保状态最终被清理。此配置未来可能移除3.0.0spark.shuffle.service.db.backendROCKSDB当spark.shuffle.service.db.enabled为 true 时指定 shuffle 服务状态存储的磁盘存储类型支持ROCKSDB和已废弃的LEVELDB。原有 RocksDB/LevelDB 数据不会自动转换存储类型3.4.0spark.storage.cleanupFilesAfterExecutorExittrueExecutor 退出后清理 Worker 目录中的非 shuffle 文件如临时 shuffle 块、缓存的 RDD/broadcast 块、溢出文件等。与spark.worker.cleanup.enabled不重叠前者清理死亡 Executor 本地目录中的非 shuffle 文件后者清理已停止且超时应用的全部文件/子目录。仅影响 Standalone 模式2.4.0spark.worker.ui.compressedLogFileLengthCacheSize100压缩日志文件无法直接得知解压后大小Spark 会缓存压缩日志的解压后大小此属性控制缓存大小2.0.2spark.worker.idPatternworker-%s-%s-%d基于 JavaString.format的 Worker ID 生成模式如worker-20231109183042-[fe80::1%lo0]-39729。请谨慎确保生成的 ID 唯一4.0.0资源分配与配置概览建议先阅读 configuration 页面 中的「Custom Resource Scheduling and Configuration Overview」一节。本文只讨论 Spark Standalone 特有的资源调度部分它分为两块一是配置 Worker 的资源二是配置特定应用的资源分配。Worker 必须配置一组可用资源才能分配给 Executor。使用spark.worker.resource.{resourceName}.amount控制每类资源在 Worker 上的总量同时必须通过spark.worker.resourcesFile或spark.worker.resource.{resourceName}.discoveryScript指定 Worker 如何发现所分配的资源两者的区别与格式见上表。应用侧唯一的特殊情况是client 模式下的 Driver可通过spark.driver.resourcesFile或spark.driver.resource.{resourceName}.discoveryScript指定 Driver 使用的资源。如果 Driver 与其他 Driver 运行在同一台主机上请确保资源文件或发现脚本只返回不与同节点其他 Driver 冲突的资源。注意提交应用时用户无需指定发现脚本因为 Worker 启动每个 Executor 时会把分配给它的资源一并传下去。连接应用与集群要在一个 Spark 集群上运行应用只需把 Master 的spark://IP:PORTURL 传给 SparkContext 构造函数val sc new SparkContext(new SparkConf().setMaster(spark://IP:PORT).setAppName(MyApp))若想在集群上运行交互式 Spark Shell./bin/spark-shell --master spark://IP:PORT还可以追加选项--total-executor-cores numCores控制 spark-shell 在集群上使用的核数。客户端属性Standalone 模式特有的客户端配置属性如下属性名默认值含义引入版本spark.standalone.submit.waitAppCompletionfalse在 Standalone 集群模式下控制客户端是否等待应用完成后再退出。为true时客户端进程保持存活并轮询 Driver 状态否则提交完成后客户端立即退出3.1.0启动 Spark 应用Spark 提交协议spark-submitspark-submit脚本 是把编译好的 Spark 应用提交到集群最直接的方式。对 Standalone 集群Spark 支持两种部署模式client 模式Driver 在与提交客户端相同的进程中启动cluster 模式Driver 在集群内某个 Worker 进程中启动客户端进程在完成提交职责后即退出不等待应用结束。如果应用通过 spark-submit 启动应用 jar 会自动分发到所有 Worker 节点。额外的依赖 jar 需要通过--jars参数用逗号分隔指定如--jars jar1,jar2。应用配置或执行环境相关设置参见 Spark Configuration。另外Standalone 的 cluster 模式支持应用以非零退出码结束时自动重启。使用该特性只需在提交时传入--supervise标志。若要杀掉反复失败的应用可执行./bin/spark-class org.apache.spark.deploy.Client kill master url driver IDDriver ID 可以通过 Standalone Master Web UI 在http://master url:8080找到。REST API当spark.master.rest.enabled启用时默认 trueSpark Master 额外提供 REST API地址格式为http://[host:port]/[version]/submissions/[action]其中host是 Master 主机port由spark.master.rest.port指定默认 6066version是协议版本目前为v1action为下列动作之一命令HTTP 方法描述引入版本createPOST通过cluster模式创建 Spark Driver。自 4.0.0 起Spark Master 支持对 Spark 属性和环境变量的值做服务端变量替换1.3.0killPOST杀死单个 Spark Driver1.3.0killallPOST杀死所有运行中的 Spark Driver4.0.0statusGET查询 Spark 作业状态1.3.0clearPOST清除已完成的 Driver 和应用4.0.0以下是用curl通过 REST API 提交pi.py的完整示例$ curl -XPOST http://IP:PORT/v1/submissions/create \ --header Content-Type:application/json;charsetUTF-8 \ --data { appResource: , sparkProperties: { spark.master: spark://master:7077, spark.app.name: Spark Pi, spark.driver.memory: 1g, spark.driver.cores: 1, spark.jars: }, clientSparkVersion: , mainClass: org.apache.spark.deploy.SparkSubmit, environmentVariables: { }, action: CreateSubmissionRequest, appArgs: [ /opt/spark/examples/src/main/python/pi.py, 10 ] }上面create请求的响应示例{ action : CreateSubmissionResponse, message : Driver successfully submitted as driver-20231124153531-0000, serverSparkVersion : 4.0.0, submissionId : driver-20231124153531-0000, success : true }当 Master 通过spark.master.rest.filtersorg.apache.spark.ui.JWSFilter和spark.org.apache.spark.ui.JWSFilter.param.secretKeyBASE64URL-ENCODED-KEY要求 HTTPAuthorization头时curl需要携带相应请求头$ curl -XPOST http://IP:PORT/v1/submissions/create \ --header Authorization: Bearer USER-PROVIDED-WEB-TOKEN-SIGNED-BY-THE-SAME-SHARED-KEY ...对于sparkProperties和environmentVariables可以使用服务端环境变量的占位符大括号包裹的变量名会在 Master 端被替换例如... sparkProperties: { spark.hadoop.fs.s3a.endpoint: {{AWS_ENDPOINT_URL}}, spark.hadoop.fs.s3a.endpoint.region: {{AWS_REGION}} }, environmentVariables: { AWS_CA_BUNDLE: {{AWS_CA_BUNDLE}} }, ...资源调度FIFO 与核数上限Standalone 集群模式目前只支持跨应用简单的 FIFO 调度器。为了允许多个并发用户可以控制每个应用使用的最大资源数。默认情况下应用会获取集群的全部核这只有在同时只跑一个应用时才合理。可以通过在 SparkConf 中设置spark.cores.max来限制核数例如val conf new SparkConf() .setMaster(...) .setAppName(...) .set(spark.cores.max, 10) val sc new SparkContext(conf)此外可以在 Master 进程上配置spark.deploy.defaultCores为未设置spark.cores.max的应用修改默认核数从无限改为有限值。在conf/spark-env.sh中添加export SPARK_MASTER_OPTS-Dspark.deploy.defaultCoresvalue这对用户未单独配置核数上限的共享集群非常有用。从源码看该属性对应 Master.scala 中的defaultCores它是spark.deploy.defaultCoresDEFAULT_CORES的取值而spreadOutAppsMaster.scala则控制调度时是把应用尽量分散到各节点还是集中到少数节点这正是原文档中spark.deploy.spreadOutApps与spark.deploy.spreadOutDrivers两项配置在调度算法中的落点。Executor 调度每个 Executor 分配的核数可配置。当显式设置spark.executor.cores时如果 Worker 有足够核数和内存同一应用的多个 Executor 可以在同一个 Worker 上启动否则默认情况下每个 Executor 会占用 Worker 上的全部可用核此时同一轮调度中每个 Worker 上每个应用只能启动一个 Executor。阶段级调度Stage Level SchedulingStandalone 支持阶段级调度关闭动态分配时用户可以在阶段级别指定不同的任务资源需求并使用启动时申请到的同一批 Executor开启动态分配时Master 为一个应用分配 Executor 时会按照多个 ResourceProfile 的 id 顺序调度id 较小的 ResourceProfile 先被调度。通常这无关紧要Spark 总是先完成一个阶段再开始下一个但在 job server 之类的场景中可能会产生影响需要留意。调度时只会从内置 Executor 资源中取 executor memory 和 executor cores其余自定义资源取自 ResourceProfile其他内置 Executor 资源如offHeap和memoryOverhead不生效。提交应用时基础默认 profile 会基于 Spark 配置创建基础默认 profile 的 executor memory 和 executor cores 可以传播到自定义 ResourceProfile但其他自定义资源不能传播。注意事项如 Dynamic Resource Allocation 所述启用动态分配且未显式指定每个 Executor 的核数时Spark 可能申请远超预期的 Executor 数量。因此使用阶段级调度时强烈建议为每个资源 profile 显式设置 executor cores。监控与日志Standalone 模式提供基于 Web 的用户界面来监控集群。Master 和每个 Worker 都有自己的 Web UI展示集群与作业统计信息。默认通过 8080 端口访问 Master 的 Web UI端口可通过配置文件或命令行选项修改。每个作业的详细日志输出还会写入每个 Worker 节点的工作目录默认SPARK_HOME/work。每个作业会看到两个文件stdout和stderr包含它写入控制台的全部输出。要跨已完成应用追踪和回看日志请启用事件日志并启动 History Server。可挂起应用Held Applications可挂起的应用会向 Master 上报自己当前是否处于挂起状态Master Web UI 会相应地标注应用状态例如RUNNING (held, draining 2 executors)。尚未退出的 Executor 仍在完成其正在运行的任务当不再有 Executor 时挂起即完成。Master 的/json/端点在每个应用的holdsupported、held和draining字段中上报同样的信息。只有 Driver 上报为可挂起的应用才会被标注前提条件见 Web UI。这类应用在Kill按钮旁边还会出现Hold按钮挂起期间显示Resume按钮。挂起会停止为应用申请新 Executor并优雅下线正在运行的 Executor应用保留 Driver 和已写入的 shuffle 输出但缓存的块会丢失恢复后需要重新计算。挂起和恢复需要与 kill 相同的修改权限。按钮由spark.ui.holdEnabled双重控制在 Master 上设为 false 会隐藏所有应用的按钮在每个应用上应用自身的设置随注册信息一起上报——禁用该功能的应用不显示按钮、挂起请求会被拒绝但挂起状态仍然可见。相同的控制在 Driver Web UI 上同样可用。与 Hadoop 并行运行可以简单地让 Spark 作为独立服务运行在与现有 Hadoop 集群相同的机器上。要从 Spark 访问 Hadoop 数据只需使用hdfs://URL通常是hdfs://namenode:9000/path具体 URL 可在 Hadoop Namenode 的 Web UI 上找到。或者也可以为 Spark 单独搭建集群通过网络访问 HDFS——这比本地磁盘访问慢但在同一局域网内例如在 Hadoop 的每个机架放置几台 Spark 机器通常不是问题。配置网络安全端口一般来说Spark 集群及其服务不会部署在公网上它们属于私有服务只应在其所属组织内部网络中可访问。访问 Spark 服务所用主机和端口的权限应限制在需要访问的源主机上。这一点对使用 Standalone 资源管理器的集群尤其重要因为它们不像其他资源管理器那样支持细粒度的访问控制。完整端口清单参见 security 页面。高可用High Availability默认情况下Standalone 调度集群对 Worker 故障是有弹性的Spark 本身可以把失败的工作转移到其他 Worker。但调度决策由 Master 做出这默认构成了单点故障如果 Master 崩溃就无法再创建新应用。为规避这一点Spark 提供了两种高可用方案。基于 ZooKeeper 的备用 Master概览利用 ZooKeeper 提供领导者选举和部分状态存储可以在集群中启动多个连接同一 ZooKeeper 实例的 Master。其中一个会被选举为 leader其余保持 standby 状态。当前 leader 宕机后另一个 Master 会被选举出来恢复旧 Master 的状态然后继续调度。整个恢复过程从第一个 leader 宕机算起大约需要 1 到 2 分钟。注意这个延迟只影响调度新应用——Master 故障转移时已在运行的应用不受影响。配置要启用该恢复模式可在 spark-env 中通过SPARK_DAEMON_JAVA_OPTS配置spark.deploy.recoveryMode及相关的spark.deploy.zookeeper.*配置。一个常见的坑如果集群中有多个 Master但没有正确配置它们使用 ZooKeeperMaster 们将无法发现彼此并认为各自都是 leader这会导致集群状态不健康所有 Master 独立调度。细节搭建好 ZooKeeper 集群后启用高可用很直接只需在不同节点上以相同的 ZooKeeper 配置ZooKeeper URL 和目录启动多个 Master 进程即可Master 可以随时增删。要调度新应用或向集群添加 Worker它们需要知道当前 leader 的 IP 地址。最简单的方式是在原来只传单个 Master 的地方传入 Master 列表。例如将 SparkContext 指向spark://host1:port1,host2:port2SparkContext 会尝试向两个 Master 注册——如果host1宕机这个配置依然有效因为会找到新 leaderhost2。需要区分向 Master 注册与正常运行启动时应用或 Worker 必须找到并注册到当前 leader Master。但成功注册后它就在系统中即存储在 ZooKeeper 中。如果发生故障转移新 leader 会联系所有已注册的应用和 Worker通知它们领导权的变化因此它们甚至不需要在启动时知道新 Master 的存在。基于这一特性新 Master 可以随时创建只需确保新应用和 Worker 在它成为 leader 时能找到并注册即可。从源码看ZOOKEEPER恢复模式由 Master.scala 中onStart()的recoveryMode分支创建ZooKeeperRecoveryModeFactory实现对应核心类包括 ZooKeeperPersistenceEngine.scala 与 ZooKeeperLeaderElectionAgent.scala。Master 被选举为 leader 后会从持久化引擎读取已存储的应用、驱动与 Worker 数据进入RECOVERING状态并完成恢复见 Master.scala。基于本地文件系统的单节点恢复概览ZooKeeper 是生产级高可用的最佳选择但如果只是想能在 Master 宕机后重启它FILESYSTEM模式就够了。应用和 Worker 注册时足够的状态会被写入指定目录这样 Master 进程重启后可以恢复这些状态。配置在 spark-env 中通过SPARK_DAEMON_JAVA_OPTS配置如下系统属性系统属性默认值含义引入版本spark.deploy.recoveryModeNONE恢复模式设为FILESYSTEM启用基于文件系统的单节点恢复ROCKSDB启用基于 RocksDB 的单节点恢复ZOOKEEPER使用基于 ZooKeeper 的恢复CUSTOM通过额外的spark.deploy.recoveryMode.factory配置提供自定义提供者类。NONE是默认值禁用此恢复模式0.8.1spark.deploy.recoveryDirectorySpark 存储恢复状态的目录从 Master 视角可访问。注意若更改spark.deploy.recoveryMode或spark.deploy.recoveryCompressionCodec该目录应手动清空0.8.1spark.deploy.recoveryCompressionCodec无持久化引擎的压缩编解码器none默认、lz4、lzf、snappy、zstd。目前仅FILESYSTEM模式支持此配置4.0.0spark.deploy.recoveryTimeout无恢复过程超时。默认值与spark.worker.timeout相同4.0.0spark.deploy.recoveryMode.factory实现StandaloneRecoveryModeFactory接口的类1.2.0spark.deploy.zookeeper.url无当spark.deploy.recoveryMode为ZOOKEEPER时设置要连接的 ZooKeeper URL0.8.1spark.deploy.zookeeper.dir无当spark.deploy.recoveryMode为ZOOKEEPER时设置存储恢复状态的 ZooKeeper 目录0.8.1对应地Master.scala 中FILESYSTEM、ROCKSDB、CUSTOM分支分别创建对应的恢复工厂同目录下的 FileSystemPersistenceEngine.scala、RocksDBPersistenceEngine.scala 与 RecoveryModeFactory.scala 构成了这些恢复模式的实现基础。细节该方案可以与进程监视/管理器如 monit配合使用也可以只用于通过重启进行手动恢复。尽管文件系统恢复看起来总比完全不恢复要好但该模式在某些开发或实验场景下可能不够理想。特别是用stop-master.sh杀掉 Master 并不会清理其恢复状态因此每次启动新 Master 都会进入恢复模式。如果它需要等待所有之前注册的 Worker/客户端超时启动时间可能增加最多 1 分钟。虽然未正式支持但可以把 NFS 目录挂载为恢复目录如果原 Master 节点完全宕机可以在另一节点启动 Master正确恢复所有已注册的 Worker/应用效果等同于 ZooKeeper 恢复。不过未来的应用仍需能发现新 Master 才能注册。小结Spark Standalone 模式用一套 Master/Worker 守护进程提供了开箱即用的资源管理能力通过 sbin/start-master.sh 与 sbin/start-worker.sh 可手动搭建通过 sbin/start-all.sh 与conf/workers文件可脚本化批量拉起Master 与 Worker 各自的SPARK_MASTER_OPTS/SPARK_WORKER_OPTS系统属性覆盖了从 Web UI 端口、应用保留数量、Executor 重试上限到自定义资源发现、工作目录清理的全部细节应用接入支持spark-submit协议与 REST API 两条路径高可用则提供 ZooKeeper 多 Master 选举与 FILESYSTEM/ROCKSDB 单节点恢复两种选择。部署时请务必结合 Spark Security 与本文的端口配置建议在网络边界上做好访问控制。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考