运行 Spark Application

本指南演示如何使用 SparkApplication 自定义资源提交 Apache Spark 作业、监控其直至完成,并检查结果。示例使用 Spark 镜像中自带的内置 Spark Pi 示例。

前提条件

  • 已安装 Alauda Build of Spark Operator,且其 CSV 报告为 Succeeded — 参见 安装
  • 已拥有目标集群的 kubectl 访问权限。

1. 创建命名空间和 Spark RBAC

在 cluster 模式下,Spark driver 会创建并管理其 executor pod,因此 driver 的 ServiceAccount 需要具备在其命名空间中管理 pod(以及 Spark 使用的 services / configmaps)的权限。

apiVersion: v1
kind: Namespace
metadata:
  name: spark-demo
---
apiVersion: v1
kind: ServiceAccount
metadata:
  name: spark
  namespace: spark-demo
---
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata:
  name: spark-driver
  namespace: spark-demo
rules:
  - apiGroups: [""]
    resources: ["pods", "services", "configmaps", "persistentvolumeclaims"]
    verbs: ["create", "get", "list", "watch", "delete", "deletecollection", "patch", "update"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: spark-driver
  namespace: spark-demo
roleRef:
  apiGroup: rbac.authorization.k8s.io
  kind: Role
  name: spark-driver
subjects:
  - kind: ServiceAccount
    name: spark
    namespace: spark-demo

保存为 spark-rbac.yaml 并应用:

kubectl apply -f spark-rbac.yaml

2. 提交 SparkApplication

apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
  name: spark-pi
  namespace: spark-demo
spec:
  type: Scala
  mode: cluster
  image: docker.io/library/spark:4.0.1
  imagePullPolicy: IfNotPresent
  mainClass: org.apache.spark.examples.SparkPi
  mainApplicationFile: "local:///opt/spark/examples/jars/spark-examples.jar"
  arguments: ["100"] 
  sparkVersion: "4.0.1"
  restartPolicy:
    type: Never
  driver:
    coreRequest: "500m"
    coreLimit: "1000m"
    memory: 512m
    serviceAccount: spark
    securityContext:
      runAsUser: 185
      runAsGroup: 185
      runAsNonRoot: true
      allowPrivilegeEscalation: false
      capabilities: { drop: ["ALL"] }
      seccompProfile: { type: RuntimeDefault }
  executor:
    instances: 2
    coreRequest: "500m"
    coreLimit: "1000m"
    memory: 512m
    securityContext:
      runAsUser: 185
      runAsGroup: 185
      runAsNonRoot: true
      allowPrivilegeEscalation: false
      capabilities: { drop: ["ALL"] }
      seccompProfile: { type: RuntimeDefault }
  1. spec.image — driver 和 executor 使用的 Apache Spark 运行时镜像。在 air-gapped 集群中,请使用平台 registry 中已重新定位的副本(记录在 CSV 的 relatedImages 中)或你们内部的镜像仓库,例如 build-harbor.alauda.cn/mlops/spark:4.0.1
  2. spec.mainApplicationFile — 应用制品的路径或 URL。local:// 表示镜像内部的路径;Spark 示例 jar 随 Spark 镜像一同提供。
  3. spec.arguments — 传递给应用的参数;对于 Spark Pi,这里是分区数。
  4. spec.driver.serviceAccount — 必须引用第 1 步中创建的 ServiceAccount
  5. spec.driver.securityContext / spec.executor.securityContext — 为满足受限的 Pod Security Admission 策略所必需。Spark 镜像以 UID/GID 185 运行。
  6. spec.executor.instances — 要运行的 executor pod 数量。

保存为 spark-pi.yaml 并应用:

kubectl apply -f spark-pi.yaml

3. 监控应用

operator 会在 .status.applicationState.state 中报告进度,其状态会依次经过 SUBMITTEDRUNNINGCOMPLETED

# high-level status
kubectl get sparkapplication spark-pi -n spark-demo

# just the state field
kubectl get sparkapplication spark-pi -n spark-demo \
  -o jsonpath='{.status.applicationState.state}{"\n"}'

# driver and executor pods (driver pod is named <app-name>-driver)
kubectl get pods -n spark-demo

4. 验证结果

当应用到达 COMPLETED 时,检查 driver 日志中计算出的 Pi 值:

kubectl logs spark-pi-driver -n spark-demo | grep "Pi is roughly"

你应该会看到类似如下的输出:

Pi is roughly 3.1415...

5. 清理

kubectl delete sparkapplication spark-pi -n spark-demo
# optionally remove the demo namespace
kubectl delete namespace spark-demo

调度周期性作业

要按计划运行 Spark 应用,请将相同的应用 spec 包装在 ScheduledSparkApplication 中:

apiVersion: sparkoperator.k8s.io/v1beta2
kind: ScheduledSparkApplication
metadata:
  name: spark-pi-scheduled
  namespace: spark-demo
spec:
  schedule: "@every 5m"
  concurrencyPolicy: Forbid
  successfulRunHistoryLimit: 3
  failedRunHistoryLimit: 1
  template:
    type: Scala
    mode: cluster
    image: docker.io/library/spark:4.0.1
    imagePullPolicy: IfNotPresent
    mainClass: org.apache.spark.examples.SparkPi
    mainApplicationFile: "local:///opt/spark/examples/jars/spark-examples.jar"
    arguments: ["100"]
    sparkVersion: "4.0.1"
    restartPolicy:
      type: Never
    driver:
      coreRequest: "500m"
      coreLimit: "1000m"
      memory: 512m
      serviceAccount: spark
    executor:
      instances: 2
      coreRequest: "500m"
      coreLimit: "1000m"
      memory: 512m
  1. spec.schedule — 用于控制运行何时启动的 cron 表达式(或 @every <duration>)。
  2. spec.concurrencyPolicyAllowForbidReplace;用于控制是否允许重叠运行。
  3. spec.template — 每个调度周期要运行的 SparkApplication spec。

查看该计划及其生成的运行实例:

kubectl get scheduledsparkapplication spark-pi-scheduled -n spark-demo
kubectl get sparkapplication -n spark-demo

常见的 SparkApplication 字段

字段描述
spec.type应用语言:ScalaJavaPythonR
spec.mode提交模式;使用 cluster
spec.imagedriver 和 executor 使用的 Spark 运行时镜像。
spec.mainClassJVM(Scala / Java)应用的主类。
spec.mainApplicationFile应用制品的路径或 URL(local://http(s)://s3a://,等等)。
spec.arguments传递给应用的参数。
spec.sparkVersion镜像中的 Spark 版本(4.0.1)。
spec.sparkConf传递给应用的任意 spark.* 配置。
spec.deps要获取的额外 jarsfilespyFiles 和 Maven 包坐标。
spec.driver / spec.executorpod 设置:cores / coreRequest / coreLimitmemoryserviceAccountinstances(executor)、labels、env、volumes、node selectors、securityContext,等等
spec.restartPolicyNeverOnFailureAlways,以及重试 / 回退设置。

完整 schema 请参见上游 Spark Operator API reference