欢迎光临
我们一直在努力

Kubernetes大数据处理架构与Spark On K8s实践

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为大数据处理提供了弹性、可扩展的运行环境。通过合理配置资源、存储和调度策略,可以高效处理大规模数据。结合监控和容错机制,可以确保作业的可靠性和可观测性。

赞(0)
未经允许不得转载:171主机测评 » Kubernetes大数据处理架构与Spark On K8s实践
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址