Kubernetes大数据处理架构与Spark On K8s实践
引言
随着大数据处理需求的增长,将大数据框架部署到Kubernetes已成为趋势。本文将深入探讨Spark on Kubernetes的部署策略和最佳实践。
一、大数据处理架构
1.1 架构设计
┌─────────────────────────────────────────────────────────────────────┐
│ 大数据处理架构 │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ ┌──────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Data │─────▶│ Spark │─────▶│ Storage │ │
│ │ Source │ │ Processing │ │ (HDFS/S3) │ │
│ └──────────┘ └──────────────┘ └──────────────┘ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌──────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Kafka │ │ Kubernetes │ │ Presto │ │
│ │ Flink │ │ Operator │ │ Trino │ │
│ └──────────┘ └──────────────┘ └──────────────┘ │
│ │
└─────────────────────────────────────────────────────────────────────┘
1.2 组件对比
| Spark | 批处理/流处理 | 大规模数据处理 |
| Flink | 实时流处理 | 低延迟实时计算 |
| Presto | SQL查询引擎 | 交互式分析 |
| Trino | 分布式查询 | 多数据源查询 |
二、Spark On Kubernetes部署
2.1 Spark配置
apiVersion: v1
kind: ConfigMap
metadata:
name: spark-config
namespace: spark
data:
spark-defaults.conf: |
spark.master k8s://https://kubernetes.default.svc
spark.kubernetes.namespace spark
spark.kubernetes.container.image spark:3.5.0
spark.kubernetes.authenticate.driver.serviceAccountName spark-driver
spark.kubernetes.authenticate.executor.serviceAccountName spark-executor
spark.executor.instances 3
spark.executor.cores 2
spark.executor.memory 4g
spark.driver.cores 1
spark.driver.memory 2g
2.2 Spark应用提交
spark-submit \\
–master k8s://https://kubernetes.default.svc \\
–deploy-mode cluster \\
–name spark-pi \\
–class org.apache.spark.examples.SparkPi \\
–conf spark.kubernetes.container.image=spark:3.5.0 \\
–conf spark.kubernetes.namespace=spark \\
–conf spark.executor.instances=3 \\
–conf spark.executor.cores=2 \\
–conf spark.executor.memory=4g \\
local:///opt/spark/examples/jars/spark-examples_2.12-3.5.0.jar
三、Spark Operator配置
3.1 安装Spark Operator
kubectl apply -f https://github.com/GoogleCloudPlatform/spark-on-k8s-operator/releases/download/v1.1.1/spark-operator.yaml
kubectl create namespace spark
kubectl apply -f https://raw.githubusercontent.com/GoogleCloudPlatform/spark-on-k8s-operator/master/charts/spark-operator/crds/sparkoperator.k8s.io_sparkapplications.yaml
3.2 SparkApplication配置
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-pi
namespace: spark
spec:
type: Scala
mode: cluster
image: spark:3.5.0
imagePullPolicy: Always
mainClass: org.apache.spark.examples.SparkPi
mainApplicationFile: local:///opt/spark/examples/jars/spark-examples_2.12-3.5.0.jar
sparkVersion: "3.5.0"
restartPolicy:
type: Never
driver:
cores: 1
coreLimit: "1200m"
memory: "512m"
labels:
version: 3.5.0
serviceAccount: spark-driver
executor:
cores: 1
instances: 3
memory: "1g"
labels:
version: 3.5.0
serviceAccount: spark-executor
四、资源管理与调度
4.1 资源配额配置
apiVersion: v1
kind: ResourceQuota
metadata:
name: spark-quota
namespace: spark
spec:
hard:
pods: "100"
requests.cpu: "50"
requests.memory: "100Gi"
limits.cpu: "100"
limits.memory: "200Gi"
4.2 节点选择器
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-data-processing
spec:
executor:
nodeSelector:
nodeType: worker
zone: us-west-2a
tolerations:
– key: "dedicated"
operator: "Equal"
value: "spark"
effect: "NoSchedule"
五、存储配置
5.1 HDFS集成
apiVersion: v1
kind: ConfigMap
metadata:
name: hdfs-config
namespace: spark
data:
hdfs-site.xml: |
<configuration>
<property>
<name>dfs.nameservices</name>
<value>mycluster</value>
</property>
<property>
<name>dfs.ha.namenodes.mycluster</name>
<value>nn1,nn2</value>
</property>
</configuration>
5.2 S3集成
apiVersion: v1
kind: Secret
metadata:
name: spark-s3-secret
namespace: spark
type: Opaque
data:
AWS_ACCESS_KEY_ID: <base64-encoded-key>
AWS_SECRET_ACCESS_KEY: <base64-encoded-secret>
六、监控与日志
6.1 Prometheus监控
apiVersion: monitoring.coreos.com/v1
kind: ServiceMonitor
metadata:
name: spark-monitor
namespace: monitoring
spec:
selector:
matchLabels:
sparkoperator.k8s.io/app-name: spark-pi
endpoints:
– port: http-prometheus
path: /metrics
interval: 15s
6.2 日志采集
apiVersion: v1
kind: ConfigMap
metadata:
name: spark-logging
namespace: spark
data:
log4j.properties: |
log4j.rootCategory=INFO, console
log4j.appender.console=org.apache.log4j.ConsoleAppender
log4j.appender.console.target=System.err
log4j.appender.console.layout=org.apache.log4j.PatternLayout
log4j.appender.console.layout.ConversionPattern=%d{yy/MM/dd HH:mm:ss} %p %c{1}: %m%n
七、作业调度
7.1 CronSparkApplication
apiVersion: sparkoperator.k8s.io/v1beta2
kind: ScheduledSparkApplication
metadata:
name: daily-report
namespace: spark
spec:
schedule: "0 2 * * *"
concurrencyPolicy: Forbid
template:
spec:
type: Python
mode: cluster
image: spark-py:3.5.0
mainApplicationFile: gs://my-bucket/scripts/daily_report.py
sparkVersion: "3.5.0"
executor:
cores: 2
instances: 5
memory: "4g"
7.2 依赖管理
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: spark-job-with-deps
spec:
dependencies:
jars:
– gs://my-bucket/jars/mysql-connector-java.jar
– gs://my-bucket/jars/postgresql.jar
files:
– gs://my-bucket/config/config.properties
pyFiles:
– gs://my-bucket/python/utils.py
八、性能优化
8.1 配置调优
| spark.executor.instances | Executor数量 | 根据数据量调整 |
| spark.executor.cores | 每个Executor核数 | 4-8 |
| spark.executor.memory | Executor内存 | 8-16g |
| spark.sql.shuffle.partitions | Shuffle分区数 | 200-1000 |
| spark.default.parallelism | 默认并行度 | 2-3倍核数 |
8.2 内存管理
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: memory-optimized-job
spec:
driver:
memory: "4g"
javaOptions: "-XX:+UseG1GC -XX:MaxGCPauseMillis=200"
executor:
memory: "16g"
memoryOverhead: "4g"
javaOptions: "-XX:+UseG1GC -XX:MaxGCPauseMillis=200"
九、高可用性配置
9.1 驱动故障恢复
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: ha-spark-job
spec:
restartPolicy:
type: OnFailure
onFailureRetries: 3
onFailureRetryInterval: 10
driver:
affinity:
nodeAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
nodeSelectorTerms:
– matchExpressions:
– key: node-role.kubernetes.io/control-plane
operator: DoesNotExist
9.2 数据容错
apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
name: fault-tolerant-job
spec:
sparkConf:
spark.task.maxFailures: "4"
spark.stage.maxConsecutiveAttempts: "4"
spark.shuffle.service.enabled: "true"
十、常见问题与解决方案
10.1 Executor启动失败
问题分析:
- 资源不足
- 镜像拉取失败
- 权限问题
解决方案:
# 检查Pod状态
kubectl get pods -n spark
# 查看Executor日志
kubectl logs <executor-pod-name> -n spark
# 检查资源配额
kubectl describe resourcequota spark-quota -n spark
10.2 数据倾斜
问题分析:
- 分区不均匀
- Key分布不均
解决方案:
# 增加Shuffle分区
–conf spark.sql.shuffle.partitions=1000
# 使用盐值随机化
df.withColumn("salt", rand() * 10)
10.3 内存溢出
问题分析:
- Executor内存不足
- 数据缓存过多
解决方案:
# 增加Executor内存
–conf spark.executor.memory=16g
# 调整内存比例
–conf spark.memory.fraction=0.8
结论
Spark on Kubernetes为大数据处理提供了弹性、可扩展的运行环境。通过合理配置资源、存储和调度策略,可以高效处理大规模数据。结合监控和容错机制,可以确保作业的可靠性和可观测性。




