惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Last Week in AI
Last Week in AI
D
DataBreaches.Net
腾讯CDC
Recent Announcements
Recent Announcements
有赞技术团队
有赞技术团队
A
About on SuperTechFans
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
Google DeepMind News
Google DeepMind News
Microsoft Security Blog
Microsoft Security Blog
云风的 BLOG
云风的 BLOG
罗磊的独立博客
月光博客
月光博客
MyScale Blog
MyScale Blog
U
Unit 42
Martin Fowler
Martin Fowler
Stack Overflow Blog
Stack Overflow Blog
T
Tailwind CSS Blog
Engineering at Meta
Engineering at Meta
N
Netflix TechBlog - Medium
G
Google Developers Blog
博客园 - 【当耐特】
D
Docker
I
InfoQ
雷峰网
雷峰网

祈雨的笔记

安全多方计算MPC spark原理解析 kueue执行源码分析 spark on k8s执行源码分析 系统压测遇到的缓存击穿问题 我的世界PC与安卓联机 蚂蚁金服流量投放平台的AIG改造 G1大对象致Old区占用率高 日志打印导致接口响应率下跌分析 Groovy加载类导致OOM分析 ERROR日志打印导致CPU满载 记OceanBase死锁超时 应用发版期间服务响应超时 Ark Serverless初探 系统优化复盘一二三 The user specified as a definer does not exist Kong网关初探 API网关选型调研 CPU火焰图常用工具 配置中心选型调研 root操作Nginx导致用户组错误 基于Proxifier使用代理 FastJSON字段智能匹配踩坑 Nacos初探 记一次Nginx服务器CPU满荷载故障 基于券系统分库分表的思考 limit不参与SQL成本计算致索引失效 Linux常用性能监控命令 golang低版本http2偶现400 hostname in certificate didn't match
spark-operator源码解析
祈雨的笔记 · 2024-07-28 · via 祈雨的笔记

Apache Spark的Kubernetes Operator旨在使指定和运行Spark应用程序变得像在Kubernetes上运行其他工作负载一样简单且符合惯例。它使用 Kubernetes 自定义资源来指定、运行和显示Spark应用程序的状态。

具体来说,用户使用sparkctl(或kubectl)创建一个SparkApplication对象。SparkApplication控制器通过 API 服务器的观察器接收对象,创建携带spark-submit参数的提交,并将提交发送给提交运行器。提交运行器提交应用程序运行并创建应用程序的驱动程序 pod。驱动程序 pod 启动后会创建执行器 pod。在应用程序运行时,Spark pod 监视器会监视应用程序的 pod,并将 pod 的状态更新发送回控制器,然后控制器会相应地更新应用程序的状态。

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
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127

func (c *Controller) syncSparkApplication(key string) error {
namespace, name, err := cache.SplitMetaNamespaceKey(key)
if err != nil {
return fmt.Errorf("failed to get the namespace and name from key %s: %v", key, err)
}
app, err := c.getSparkApplication(namespace, name)
if err != nil {
return err
}
if app == nil {

return nil
}
if !app.DeletionTimestamp.IsZero() {
c.handleSparkApplicationDeletion(app)
return nil
}

appCopy := app.DeepCopy()


v1beta2.SetSparkApplicationDefaults(appCopy)


switch appCopy.Status.AppState.State {
case v1beta2.NewState:
c.recordSparkApplicationEvent(appCopy)
if err := c.validateSparkApplication(appCopy); err != nil {
appCopy.Status.AppState.State = v1beta2.FailedState
appCopy.Status.AppState.ErrorMessage = err.Error()
} else {
appCopy = c.submitSparkApplication(appCopy)
}
case v1beta2.SucceedingState:
if !shouldRetry(appCopy) {
appCopy.Status.AppState.State = v1beta2.CompletedState
c.recordSparkApplicationEvent(appCopy)
} else {
if err := c.deleteSparkResources(appCopy); err != nil {
glog.Errorf("failed to delete resources associated with SparkApplication %s/%s: %v",
appCopy.Namespace, appCopy.Name, err)
return err
}
appCopy.Status.AppState.State = v1beta2.PendingRerunState
}
case v1beta2.FailingState:
if !shouldRetry(appCopy) {
appCopy.Status.AppState.State = v1beta2.FailedState
c.recordSparkApplicationEvent(appCopy)
} else if isNextRetryDue(appCopy.Spec.RestartPolicy.OnFailureRetryInterval, appCopy.Status.ExecutionAttempts, appCopy.Status.TerminationTime) {
if err := c.deleteSparkResources(appCopy); err != nil {
glog.Errorf("failed to delete resources associated with SparkApplication %s/%s: %v",
appCopy.Namespace, appCopy.Name, err)
return err
}
appCopy.Status.AppState.State = v1beta2.PendingRerunState
}
case v1beta2.FailedSubmissionState:
if !shouldRetry(appCopy) {

appCopy.Status.AppState.State = v1beta2.FailedState
c.recordSparkApplicationEvent(appCopy)
} else if isNextRetryDue(appCopy.Spec.RestartPolicy.OnSubmissionFailureRetryInterval, appCopy.Status.SubmissionAttempts, appCopy.Status.LastSubmissionAttemptTime) {
if c.validateSparkResourceDeletion(appCopy) {
c.submitSparkApplication(appCopy)
} else {
if err := c.deleteSparkResources(appCopy); err != nil {
glog.Errorf("failed to delete resources associated with SparkApplication %s/%s: %v",
appCopy.Namespace, appCopy.Name, err)
return err
}
}
}
case v1beta2.InvalidatingState:

if err := c.deleteSparkResources(appCopy); err != nil {
glog.Errorf("failed to delete resources associated with SparkApplication %s/%s: %v",
appCopy.Namespace, appCopy.Name, err)
return err
}
c.clearStatus(&appCopy.Status)
appCopy.Status.AppState.State = v1beta2.PendingRerunState
case v1beta2.PendingRerunState:
glog.V(2).Infof("SparkApplication %s/%s is pending rerun", appCopy.Namespace, appCopy.Name)
if c.validateSparkResourceDeletion(appCopy) {
glog.V(2).Infof("Resources for SparkApplication %s/%s successfully deleted", appCopy.Namespace, appCopy.Name)
c.recordSparkApplicationEvent(appCopy)
c.clearStatus(&appCopy.Status)
appCopy = c.submitSparkApplication(appCopy)
}
case v1beta2.SubmittedState, v1beta2.RunningState, v1beta2.UnknownState:
if err := c.getAndUpdateAppState(appCopy); err != nil {
return err
}
case v1beta2.CompletedState, v1beta2.FailedState:
if c.hasApplicationExpired(app) {
glog.Infof("Garbage collecting expired SparkApplication %s/%s", app.Namespace, app.Name)
err := c.crdClient.SparkoperatorV1beta2().SparkApplications(app.Namespace).Delete(context.TODO(), app.Name, metav1.DeleteOptions{GracePeriodSeconds: int64ptr(0)})
if err != nil && !errors.IsNotFound(err) {
return err
}
return nil
}
if err := c.getAndUpdateExecutorState(appCopy); err != nil {
return err
}
}

if appCopy != nil {
err = c.updateStatusAndExportMetrics(app, appCopy)
if err != nil {
glog.Errorf("failed to update SparkApplication %s/%s: %v", app.Namespace, app.Name, err)
return err
}

if state := appCopy.Status.AppState.State; state == v1beta2.CompletedState ||
state == v1beta2.FailedState {
if err := c.cleanUpOnTermination(app, appCopy); err != nil {
glog.Errorf("failed to clean up resources for SparkApplication %s/%s: %v", app.Namespace, app.Name, err)
return err
}
}
}

return nil
}