随着云原生技术的全面普及,越来越多的企业开始将 Spark 作业从传统的 YARN 集群迁移到 Kubernetes 上。Spark 从 2.3 版本开始官方支持 Kubernetes 作为集群管理器(即 Spark on Kubernetes,简称 Spark-on-K8s),到 3.x 版本已经趋于成熟稳定。相比 YARN,K8s 提供了更细粒度的资源隔离、更灵活的弹性伸缩能力,以及与云原生生态系统的无缝集成。本文将系统性地讲解 Spark on Kubernetes 的架构原理、部署方式、生产级配置调优、监控告警体系以及常见问题的排查方案,帮助你在生产环境中稳定运行大规模 Spark 作业。
一、Spark on Kubernetes 架构原理
Spark on Kubernetes 的运行模式与 YARN / Standalone 有本质区别。在 K8s 环境下,Spark 不再依赖一个长期运行的 Resource Manager,而是直接与 Kubernetes API Server 交互,通过创建 Pod 来运行 Executor。理解这套架构是后续所有调优的基础。
1.1 核心组件交互流程
当用户提交一个 Spark 作业到 K8s 集群时,整个执行流程如下:
- Spark Driver:以 Pod 形式运行在 K8s 集群内(cluster mode)或运行在提交端(client mode),负责解析作业、生成 DAG、调度 Task、并直接调用 K8s API 创建 Executor Pod。
- Executor Pod:由 Driver 动态创建,每个 Pod 对应一个或多个 Task 执行器,作业结束后被自动清理。
- Kubernetes API Server:作为 Spark 与 K8s 集群的通信入口,Driver 通过它创建、监听和删除 Executor Pod。
- Kubernetes Scheduler:负责将 Driver 和 Executor Pod 调度到合适的 Node 上运行。
关键区别在于:YARN 模式下 Executor 由 ApplicationMaster 通过 YARN ResourceManager 分配 Container;而 K8s 模式下,Driver 本身就承担了类似 ApplicationMaster 的角色,直接调用 Kubernetes API 创建 Pod,省去了中间的 Resource Manager 层。
1.2 Cluster Mode vs Client Mode
Spark on K8s 支持两种部署模式:
| 模式 | Driver 运行位置 | 适用场景 | 生命周期管理 |
|---|---|---|---|
| Cluster Mode | K8s 集群内 Pod | 生产环境、批处理作业 | 由 K8s 管理,可配置 TTL |
| Client Mode | 提交客户端本地 | 交互式开发、调试 | 客户端断开则 Driver 终止 |
生产环境推荐使用 Cluster Mode,因为 Driver 运行在 K8s 内部,即使提交端断开连接,作业仍可继续执行。同时 Driver Pod 可以被 K8s 重新调度(取决于配置),提高了容错性。
二、生产级部署实战
2.1 环境准备与镜像构建
Spark on K8s 需要一个包含 Spark 运行时的容器镜像。虽然 Spark 官方提供了基础镜像,但生产环境通常需要自定义镜像来包含额外的依赖库。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18 # Dockerfile 示例:构建生产级 Spark 镜像
FROM spark:3.5.1-scala2.12-java17-python3
# 安装系统依赖
USER root
RUN apt-get update && apt-get install -y curl krb5-user && rm -rf /var/lib/apt/lists/*
# 添加企业内部依赖 JAR
COPY libs/*.jar /opt/spark/jars/
# 添加自定义 Python 依赖
RUN pip install --no-cache-dir pyspark==3.5.1 pandas==2.2.0 numpy==1.26.4 pyarrow==15.0.0
# 添加作业代码
COPY jobs/ /opt/spark/work-dir/
USER spark
WORKDIR /opt/spark/work-dir
构建并推送镜像:
1
2 docker build -t registry.internal.com/spark:3.5.1-prod .
docker push registry.internal.com/spark:3.5.1-prod
2.2 RBAC 权限配置
Driver Pod 需要 K8s API 权限来创建和删除 Executor Pod。必须预先配置 ServiceAccount、Role 和 RoleBinding:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39 # spark-rbac.yaml
apiVersion: v1
kind: ServiceAccount
metadata:
name: spark-driver
namespace: spark-jobs
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
name: spark-driver-role
namespace: spark-jobs
rules:
- apiGroups: [""]
resources: ["pods"]
verbs: ["get", "list", "watch", "create", "delete"]
- apiGroups: [""]
resources: ["services"]
verbs: ["get", "create", "delete"]
- apiGroups: [""]
resources: ["configmaps"]
verbs: ["get", "create", "delete", "update"]
- apiGroups: [""]
resources: ["events"]
verbs: ["list", "watch"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
name: spark-driver-binding
namespace: spark-jobs
subjects:
- kind: ServiceAccount
name: spark-driver
namespace: spark-jobs
roleRef:
kind: Role
name: spark-driver-role
apiGroup: rbac.authorization.k8s.io
2.3 提交作业的完整命令
以下是一个生产级 Spark 作业提交命令的完整示例,包含了资源配置、动态分配、存储后端等关键参数:
1 spark-submit --master k8s://https://k8s-api.internal.com:6443 --deploy-mode cluster --name etl-daily-pipeline --class com.company.ETLJob --conf spark.kubernetes.namespace=spark-jobs --conf spark.kubernetes.container.image=registry.internal.com/spark:3.5.1-prod --conf spark.kubernetes.authenticate.driver.serviceAccountName=spark-driver --conf spark.kubernetes.authenticate.submission.caCertFile=/etc/k8s/ca.crt --conf spark.kubernetes.authenticate.submission.clientCertFile=/etc/k8s/client.crt --conf spark.kubernetes.authenticate.submission.clientKeyFile=/etc/k8s/client.key --conf spark.kubernetes.driver.request.cores=4 --conf spark.kubernetes.driver.request.memory=8g --conf spark.kubernetes.driver.limit.memory=12g --conf spark.kubernetes.executor.request.cores=4 --conf spark.kubernetes.executor.request.memory=16g --conf spark.kubernetes.executor.limit.memory=20g --conf spark.executor.instances=20 --conf spark.dynamicAllocation.enabled=true --conf spark.dynamicAllocation.minExecutors=5 --conf spark.dynamicAllocation.maxExecutors=50 --conf spark.dynamicAllocation.shuffleTracking.enabled=true --conf spark.kubernetes.executor.podNamePrefix=etl-exec --conf spark.kubernetes.allocation.batch.size=5 --conf spark.kubernetes.allocation.batch.delay=1s --conf spark.kubernetes.node.selector.nodegroup=spark --conf spark.kubernetes.driver.node.selector.nodegroup=spark-driver --conf spark.kubernetes.driver.tolerations[0].key=spark-driver --conf spark.kubernetes.driver.tolerations[0].operator=Exists --conf spark.kubernetes.executor.tolerations[0].key=spark-executor --conf spark.kubernetes.executor.tolerations[0].operator=Exists --conf spark.local.dirs=/tmp/spark-local --conf spark.sql.adaptive.enabled=true --conf spark.sql.adaptive.coalescePartitions.enabled=true --conf spark.serializer=org.apache.spark.serializer.KryoSerializer --conf spark.sql.shuffle.partitions=200 local:///opt/spark/work-dir/etl-job.jar --date 2026-09-07
这个命令涵盖了生产环境的几乎所有关键配置维度:命名空间隔离、镜像指定、RBAC 认证、Driver/Executor 资源配额、动态分配、污点容忍、节点选择器、自适应查询。每一个参数都有其对应的业务含义,下文将逐一解析。
三、关键调优参数深度解析
3.1 内存配置:Request vs Limit
K8s 的内存模型与 YARN 有本质区别。在 YARN 中只有一个内存参数,而在 K8s 中有
1 | request |
和
1 | limit |
两个值:
- request.memory:调度依据,K8s 调度器据此判断节点是否有足够资源放置 Pod。
- limit.memory:Pod 可使用的内存上限,超过会被 OOMKilled。
- Spark 内存的分配:Spark 实际使用的内存是
1request.memory - memoryOverhead
,其中 overhead 默认为 request 的 10%。
关键公式:
1 | spark.executor.memory = request.memory - spark.executor.memoryOverhead |
。如果 limit > request,多出的空间实际不会被 Spark 利用(Spark 按 executor.memory 申请),但可以缓冲防止 OOM。生产建议设置 limit 比 request 高 20%-30%:
1
2
3
4
5 # 推荐配置策略
# Executor request.memory = 16g, limit.memory = 20g
# memoryOverhead = 16g * 0.15 = 2.4g(调高 overhead 用于堆外内存)
# Spark 堆内存 = 16g - 2.4g = 13.6g
--conf spark.kubernetes.executor.request.memory=16g --conf spark.kubernetes.executor.limit.memory=20g --conf spark.executor.memory=13g --conf spark.executor.memoryOverhead=2g
3.2 动态分配与 Shuffle Tracking
在 K8s 上使用动态分配时,最大的挑战是 Shuffle 数据的存储。YARN 模式下有 External Shuffle Service,而在原生 K8s 模式下,Shuffle 数据存储在 Executor Pod 的本地磁盘上,Executor 被回收后 Shuffle 数据会丢失。
从 Spark 3.0 开始,引入了 Shuffle Tracking 机制来解决这个问题。它通过跟踪每个 Executor 的 Shuffle 输出状态,只有在确认该 Executor 的所有 Shuffle 数据都已被消费完毕后,才会回收该 Executor:
1
2
3
4
5
6 --conf spark.dynamicAllocation.enabled=true
--conf spark.dynamicAllocation.minExecutors=5
--conf spark.dynamicAllocation.maxExecutors=50
--conf spark.dynamicAllocation.shuffleTracking.enabled=true
--conf spark.dynamicAllocation.shuffleTracking.timeout=3600s
--conf spark.dynamicAllocation.executorAllocationRatio=0.8
1 | shuffleTracking.timeout |
设置了一个保守的超时时间,即使无法确认 Shuffle 数据是否消费完毕,也会在超时后回收 Executor。生产环境建议设置为作业最长运行时间的 1.5 倍。另外也可以使用 External Shuffle Service on K8s(通过 DaemonSet 部署),将 Shuffle 数据持久化到共享存储,实现真正的 Executor 无状态回收。
3.3 Pod 调度优化:批量分配与亲和性
当 Executor 数量较大时(如 100+),K8s API 调用会成为瓶颈。Spark 默认一次创建一个 Executor Pod,对于大规模作业来说初始化非常缓慢。通过批量分配配置可以显著改善:
1
2
3 # 批量分配:一次创建5个Pod,间隔1秒
--conf spark.kubernetes.allocation.batch.size=5
--conf spark.kubernetes.allocation.batch.delay=1s
对于 IO 密集型作业,建议配置 Pod 反亲和性,让 Executor 分散到不同的 Node 上,避免单个节点的网络和磁盘 IO 成为瓶颈:
1 --conf spark.kubernetes.executor.podTemplateFile=/path/to/executor-template.yaml
executor-template.yaml 中可以定义反亲和规则:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27 apiVersion: v1
kind: Pod
metadata:
labels:
spark-role: executor
spec:
affinity:
podAntiAffinity:
preferredDuringSchedulingIgnoredDuringExecution:
- weight: 100
podAffinityTerm:
labelSelector:
matchExpressions:
- key: spark-role
operator: In
values:
- executor
topologyKey: kubernetes.io/hostname
containers:
- name: spark-kubernetes-executor
volumeMounts:
- name: spark-local-dir
mountPath: /data
volumes:
- name: spark-local-dir
emptyDir:
medium: "" # 不使用 tmpfs,避免占用内存
四、监控与可观测性体系
4.1 Spark UI 的访问方式
在 Cluster Mode 下,Driver 运行在 K8s 内部,Spark UI 默认无法直接从外部访问。有以下三种方案:
方案一:Service + Port-Forward(开发调试)
1
2
3
4 export DRIVER_POD=$(kubectl get pods -n spark-jobs -l spark-app-name=etl-daily-pipeline -l spark-role=driver -o jsonpath='{.items[0].metadata.name}')
kubectl port-forward -n spark-jobs $DRIVER_POD 4040:4040
# 然后浏览器访问 http://localhost:4040
方案二:Ingress + Service(生产推荐)
Spark 3.3+ 原生支持创建 Driver UI Service 和 Ingress,只需配置几个参数:
1
2
3
4
5
6
7 --conf spark.kubernetes.driver.service.enabled=true
--conf spark.kubernetes.ui.service.enabled=true
--conf spark.kubernetes.ui.service.type=ClusterIP
--conf spark.kubernetes.ui.ingress.enabled=true
--conf spark.kubernetes.ui.ingress.url-format="spark-ui.internal.com/{driver-pod-name}"
--conf spark.kubernetes.ui.ingress.className=nginx
--code>
4.2 Prometheus 指标采集
Spark 内置了 Prometheus 指标暴露端点。配置
1 | metrics.properties |
后,Spark UI 会同时暴露
1 | /metrics/prometheus |
端点:
1
2
3
4
5
6 # metrics.properties
*.sink.prometheus.class=org.apache.spark.metrics.sink.PrometheusSink
*.sink.prometheus.protocol=http
driver.sink.prometheus.endpoint=/metrics/prometheus
executor.sink.prometheus.endpoint=/metrics/prometheus
*.sink.prometheus.period=10
配合 Prometheus 的 ServiceMonitor 或 PodMonitor 自动发现 Driver 和 Executor Pod 的指标端点:
1
2
3
4
5
6
7
8
9
10
11
12
13 apiVersion: monitoring.coreos.com/v1
kind: PodMonitor
metadata:
name: spark-driver-metrics
namespace: spark-jobs
spec:
selector:
matchLabels:
spark-role: driver
podMetricsEndpoints:
- port: driver-ui # Spark UI 端口映射
path: /metrics/prometheus
interval: 15s
4.3 关键监控指标与告警规则
生产环境需要重点监控的 Spark 指标:
| 指标名称 | 含义 | 告警阈值建议 |
|---|---|---|
| spark_executor_alive_workers | 存活的 Executor 数量 | 低于 minExecutors 的 50% |
| spark_executor_removed_total | 被移除的 Executor 总数 | 5 分钟内超过 maxExecutors 的 20% |
| spark_driver_failed_tasks | Driver 失败 Task 数 | 持续增长超过 100 |
| spark_executor_memory_used_bytes | Executor 内存使用量 | 超过 limit 的 85% |
| spark_executor_disk_used_bytes | Shuffle 磁盘使用量 | 超过节点磁盘的 80% |
建议的 Prometheus 告警规则示例:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21 groups:
- name: spark-alerts
rules:
- alert: SparkExecutorOOMRisk
expr: |
spark_executor_memory_used_bytes /
on(pod) kube_pod_container_resource_limits{resource="memory"}
> 0.85
for: 5m
labels:
severity: warning
annotations:
summary: "Spark Executor 内存使用率超过 85%"
description: "Pod {{ $labels.pod }} 内存使用率: {{ $value | humanizePercentage }}"
- alert: SparkExecutorHighRemovalRate
expr: rate(spark_executor_removed_total[5m]) > 0.5
for: 2m
labels:
severity: critical
annotations:
summary: "Spark Executor 被频繁移除"
五、生产环境常见问题与解决方案
5.1 Executor Pod 创建失败(ImagePullBackOff)
这是最常见的初次部署问题,原因通常是镜像仓库认证未配置。在 K8s 集群中创建 Secret 并通过 spark-submit 传入:
1
2
3
4
5 # 创建镜像拉取 Secret
kubectl create secret docker-registry regcred --docker-server=registry.internal.com --docker-username=pull-user --docker-password=pull-password --namespace=spark-jobs
# spark-submit 中指定
--conf spark.kubernetes.container.imagePullSecrets=regcred
5.2 Driver Pod 频繁 OOMKilled
Driver OOM 通常是因为
1 | spark.driver.memory |
设置不足,或者收集了大量 Task 结果到 Driver。排查思路:
- 检查
1limit.memory
与
1spark.driver.memory的差距,确保 overhead 足够
- 如果是 collect() 导致的,改用
1toLocalIterator()
分批拉取
- 开启
1spark.driver.maxResultSize
限制,默认 1g,可按需调整
- 检查是否有大量 Broadcast Variable 堆积在 Driver 端
5.3 Shuffle 数据丢失导致作业失败
当 Executor 被回收或 Pod 被驱逐(如节点资源压力),Shuffle 数据丢失会导致下游 Task 失败。解决方案:
1
2
3
4
5
6
7
8
9
10
11
12 # 方案一:增大 shuffleTracking 超时
--conf spark.dynamicAllocation.shuffleTracking.timeout=7200s
# 方案二:开启 External Shuffle Service(需部署 DaemonSet)
--conf spark.shuffle.service.enabled=true
--conf spark.shuffle.service.port=7337
# 方案三:使用远程 Shuffle Service(推荐生产环境)
# 部署 Apache Celeborn 或 AWS EMR Shuffle Service
--conf spark.shuffle.manager=org.apache.spark.sql.execution.celeborn.CelebornShuffleManager
--conf spark.celeborn.master.endpoints=celeborn-master:9097
--conf spark.celeborn.client.spark.shuffle.writer=HASH
Apache Celeborn 是一个专门为 Spark/Flink 设计的远程 Shuffle Service,将 Shuffle 数据写入独立的服务集群,彻底解耦了 Executor 生命周期与 Shuffle 数据。对于大规模生产环境(100+ Executor、TB 级 Shuffle 数据),强烈推荐采用此方案。
5.4 K8s API 限流导致作业卡死
大规模作业(创建 200+ Pod)时,Driver 会高频调用 K8s API,可能触发 API Server 的限流(HTTP 429)。建议:
- 增大批量分配参数
1spark.kubernetes.allocation.batch.size
(上限 30)
- 增大
1spark.kubernetes.allocation.batch.delay
(如 2s)
- 配置 API Server 客户端的 QPS 和 Burst:
1--conf spark.kubernetes.driver.apiQPS=50 --conf spark.kubernetes.driver.apiBurst=100
- 如有多个作业并发,考虑引入 Volcano 或 Yunikorn 进行 K8s 层面的作业级调度
六、成本优化与最佳实践
6.1 Spot/Preemptible 实例利用
在云环境(AWS/GCP/Azure)中,Executor 可以运行在 Spot 实例上,成本可降低 60%-70%。但需要处理 Pod 被抢占的情况:
1
2
3
4
5
6
7
8
9 # nodeSelector 选择 Spot 节点
--conf spark.kubernetes.executor.node.selector.instance-type=spot
# 开启 Task 级别重试
--conf spark.task.maxFailures=4
--conf spark.stage.maxConsecutiveAttempts=4
# 配置 Pod Disruption Budget 保护关键 Executor
# 通过 Pod Template 注入
6.2 资源配额管理
对于多租户的 K8s 集群,建议为 Spark 作业创建独立的 Namespace 并设置 ResourceQuota:
1
2
3
4
5
6
7
8
9
10
11 apiVersion: v1
kind: ResourceQuota
metadata:
name: spark-quota
namespace: spark-jobs
spec:
hard:
requests.cpu: "500"
requests.memory: 2000Gi
pods: "200"
services: "100"
同时建议配置 PriorityClass,让高优先级作业能抢占低优先级作业的资源:
1
2
3
4
5
6
7 apiVersion: scheduling.k8s.io/v1
kind: PriorityClass
metadata:
name: spark-high-priority
value: 1000
globalDefault: false
description: "高优先级 Spark 作业"
6.3 作业提交最佳实践总结
综合以上内容,生产环境 Spark on K8s 的最佳实践清单:
- 统一镜像管理:构建包含业务依赖的基础镜像,禁止在运行时动态下载依赖
- 资源 request/limit 合理设置:limit 比 request 高 20%-30%,memoryOverhead 调高至 15%
- 动态分配 + Shuffle Tracking:开启动态分配,设置合理的 min/max Executors
- 批量 Pod 分配:批量 size 设为 5-10,避免 API 限流
- 节点亲和与污点:使用专用 NodeGroup,通过 nodeSelector 和 toleration 隔离
- 远程 Shuffle Service:大规模作业使用 Celeborn 等远程 Shuffle
- 完整监控体系:Prometheus + Grafana + 告警,覆盖 Executor 存活、内存、磁盘
- Pod Template 灵活配置:通过 Pod Template 文件管理复杂配置,避免在 spark-submit 中堆叠参数
结语
Spark on Kubernetes 已成为云原生大数据架构的标准方案之一。相比传统 YARN,它在资源利用率、弹性伸缩、多租户隔离方面具有显著优势,但也引入了 Shuffle 数据管理、API 调用限流、Pod 调度等新的挑战。本文从架构原理到部署实战,从调优参数到监控告警,系统性地覆盖了生产环境的核心实践。关键在于:理解 K8s 调度模型与 Spark 资源模型的差异,合理配置 request/limit,引入远程 Shuffle Service 解决数据可靠性问题,并建立完善的监控告警体系。随着 Spark 3.5+ 对 K8s 支持的持续增强以及 Volcano、Celeborn 等生态项目的成熟,Spark on Kubernetes 已经完全具备承载 PB 级数据处理生产作业的能力。
汤不热吧