Skip to content

Commit 05a3440

Browse files
committed
refactor: Introduce build aggregator for SparkApplication
1 parent 5b31a9e commit 05a3440

2 files changed

Lines changed: 133 additions & 95 deletions

File tree

rust/operator-binary/src/spark_k8s_controller.rs

Lines changed: 33 additions & 95 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,11 @@ use std::sync::Arc;
33
use snafu::{ResultExt, Snafu};
44
use stackable_operator::{
55
builder::{self},
6+
k8s_openapi::api::{
7+
batch::v1::Job,
8+
core::v1::{ConfigMap, ServiceAccount},
9+
rbac::v1::RoleBinding,
10+
},
611
kube::{
712
ResourceExt,
813
core::{DeserializeGuard, error_boundary},
@@ -115,6 +120,18 @@ impl ReconcilerError for Error {
115120
}
116121
}
117122

123+
/// Every Kubernetes resource produced by the build step for a SparkApplication.
124+
///
125+
/// Built without a Kubernetes client: all references are already dereferenced and validated by
126+
/// this point, so the only errors possible during assembly are resource-construction failures.
127+
pub struct SparkResources {
128+
pub service_account: ServiceAccount,
129+
pub role_binding: RoleBinding,
130+
/// Driver pod-template, executor pod-template, and submit-job ConfigMaps (in that order).
131+
pub config_maps: Vec<ConfigMap>,
132+
pub job: Job,
133+
}
134+
118135
pub async fn reconcile(
119136
spark_application: Arc<DeserializeGuard<v1alpha1::SparkApplication>>,
120137
ctx: Arc<Ctx>,
@@ -147,117 +164,38 @@ pub async fn reconcile(
147164
.context(ValidateSparkApplicationSnafu)?;
148165

149166
let spark_application = &validated.spark_application;
150-
let opt_s3conn = &validated.cluster_config.s3_connection;
151-
let logdir = &validated.cluster_config.log_dir;
152-
let resolved_product_image = &validated.resolved_product_image;
153167
// This is the final version of the spark app to reconcile.
154168
// No more mutating operations after this point (except for status).
155169
tracing::debug!("reconciling spark application [{spark_application:?}]");
156170

157-
let (serviceaccount, rolebinding) =
158-
build::resource::serviceaccount::build_spark_role_serviceaccount(&validated)?;
159-
client
160-
.apply_patch(SPARK_CONTROLLER_NAME, &serviceaccount, &serviceaccount)
161-
.await
162-
.context(ApplyServiceAccountSnafu)?;
163-
client
164-
.apply_patch(SPARK_CONTROLLER_NAME, &rolebinding, &rolebinding)
165-
.await
166-
.context(ApplyRoleBindingSnafu)?;
167-
168-
let env_vars = spark_application.env(opt_s3conn, logdir);
169-
170-
let driver_config = spark_application
171-
.driver_config()
172-
.context(FailedToResolveConfigSnafu)?;
173-
174-
let driver_config_overrides = spark_application
175-
.spec
176-
.driver
177-
.as_ref()
178-
.map(|driver| driver.config_overrides.clone())
179-
.unwrap_or_default();
180-
181-
let driver_pod_template_config_map = build::resource::config_map::pod_template_config_map(
182-
&validated,
183-
SparkApplicationRole::Driver,
184-
&driver_config,
185-
&driver_config_overrides,
186-
&env_vars,
187-
&serviceaccount,
188-
)?;
189-
client
190-
.apply_patch(
191-
SPARK_CONTROLLER_NAME,
192-
&driver_pod_template_config_map,
193-
&driver_pod_template_config_map,
194-
)
195-
.await
196-
.context(ApplyApplicationSnafu)?;
197-
198-
let executor_config = spark_application
199-
.executor_config()
200-
.context(FailedToResolveConfigSnafu)?;
171+
let resources = build::build(&validated)?;
201172

202-
let executor_config_overrides = spark_application
203-
.spec
204-
.executor
205-
.as_ref()
206-
.map(|executor| executor.config.config_overrides.clone())
207-
.unwrap_or_default();
208-
209-
let executor_pod_template_config_map = build::resource::config_map::pod_template_config_map(
210-
&validated,
211-
SparkApplicationRole::Executor,
212-
&executor_config,
213-
&executor_config_overrides,
214-
&env_vars,
215-
&serviceaccount,
216-
)?;
173+
// Apply the ServiceAccount and RoleBinding first, then the ConfigMaps, and finally the Job:
174+
// the Job runs under the ServiceAccount and mounts the ConfigMaps, so they must exist first.
217175
client
218176
.apply_patch(
219177
SPARK_CONTROLLER_NAME,
220-
&executor_pod_template_config_map,
221-
&executor_pod_template_config_map,
178+
&resources.service_account,
179+
&resources.service_account,
222180
)
223181
.await
224-
.context(ApplyApplicationSnafu)?;
225-
226-
let job_commands = spark_application
227-
.build_command(opt_s3conn, logdir, &resolved_product_image.image)
228-
.context(BuildCommandSnafu)?;
229-
230-
let submit_config = spark_application
231-
.submit_config()
232-
.context(SubmitConfigSnafu)?;
233-
234-
let submit_config_overrides = spark_application
235-
.spec
236-
.job
237-
.as_ref()
238-
.map(|job| job.config_overrides.clone())
239-
.unwrap_or_default();
240-
241-
let submit_job_config_map =
242-
build::resource::config_map::submit_job_config_map(&validated, &submit_config_overrides)?;
182+
.context(ApplyServiceAccountSnafu)?;
243183
client
244184
.apply_patch(
245185
SPARK_CONTROLLER_NAME,
246-
&submit_job_config_map,
247-
&submit_job_config_map,
186+
&resources.role_binding,
187+
&resources.role_binding,
248188
)
249189
.await
250-
.context(ApplyApplicationSnafu)?;
251-
252-
let job = build::resource::job::spark_job(
253-
&validated,
254-
&serviceaccount,
255-
&env_vars,
256-
&job_commands,
257-
&submit_config,
258-
)?;
190+
.context(ApplyRoleBindingSnafu)?;
191+
for config_map in &resources.config_maps {
192+
client
193+
.apply_patch(SPARK_CONTROLLER_NAME, config_map, config_map)
194+
.await
195+
.context(ApplyApplicationSnafu)?;
196+
}
259197
client
260-
.apply_patch(SPARK_CONTROLLER_NAME, &job, &job)
198+
.apply_patch(SPARK_CONTROLLER_NAME, &resources.job, &resources.job)
261199
.await
262200
.context(ApplyApplicationSnafu)?;
263201

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,2 +1,102 @@
11
pub mod pod;
22
pub mod resource;
3+
4+
use snafu::ResultExt;
5+
6+
use crate::{
7+
crd::roles::SparkApplicationRole,
8+
spark_k8s_controller::{
9+
BuildCommandSnafu, FailedToResolveConfigSnafu, Result, SparkResources, SubmitConfigSnafu,
10+
validate::ValidatedSparkApplication,
11+
},
12+
};
13+
14+
/// Builds every Kubernetes resource for the given validated SparkApplication.
15+
pub fn build(validated: &ValidatedSparkApplication) -> Result<SparkResources> {
16+
let spark_application = &validated.spark_application;
17+
let opt_s3conn = &validated.cluster_config.s3_connection;
18+
let logdir = &validated.cluster_config.log_dir;
19+
let resolved_product_image = &validated.resolved_product_image;
20+
21+
let (service_account, role_binding) =
22+
resource::serviceaccount::build_spark_role_serviceaccount(validated)?;
23+
24+
let env_vars = spark_application.env(opt_s3conn, logdir);
25+
26+
let driver_config = spark_application
27+
.driver_config()
28+
.context(FailedToResolveConfigSnafu)?;
29+
30+
let driver_config_overrides = spark_application
31+
.spec
32+
.driver
33+
.as_ref()
34+
.map(|driver| driver.config_overrides.clone())
35+
.unwrap_or_default();
36+
37+
let driver_pod_template_config_map = resource::config_map::pod_template_config_map(
38+
validated,
39+
SparkApplicationRole::Driver,
40+
&driver_config,
41+
&driver_config_overrides,
42+
&env_vars,
43+
&service_account,
44+
)?;
45+
46+
let executor_config = spark_application
47+
.executor_config()
48+
.context(FailedToResolveConfigSnafu)?;
49+
50+
let executor_config_overrides = spark_application
51+
.spec
52+
.executor
53+
.as_ref()
54+
.map(|executor| executor.config.config_overrides.clone())
55+
.unwrap_or_default();
56+
57+
let executor_pod_template_config_map = resource::config_map::pod_template_config_map(
58+
validated,
59+
SparkApplicationRole::Executor,
60+
&executor_config,
61+
&executor_config_overrides,
62+
&env_vars,
63+
&service_account,
64+
)?;
65+
66+
let job_commands = spark_application
67+
.build_command(opt_s3conn, logdir, &resolved_product_image.image)
68+
.context(BuildCommandSnafu)?;
69+
70+
let submit_config = spark_application
71+
.submit_config()
72+
.context(SubmitConfigSnafu)?;
73+
74+
let submit_config_overrides = spark_application
75+
.spec
76+
.job
77+
.as_ref()
78+
.map(|job| job.config_overrides.clone())
79+
.unwrap_or_default();
80+
81+
let submit_job_config_map =
82+
resource::config_map::submit_job_config_map(validated, &submit_config_overrides)?;
83+
84+
let job = resource::job::spark_job(
85+
validated,
86+
&service_account,
87+
&env_vars,
88+
&job_commands,
89+
&submit_config,
90+
)?;
91+
92+
Ok(SparkResources {
93+
service_account,
94+
role_binding,
95+
config_maps: vec![
96+
driver_pod_template_config_map,
97+
executor_pod_template_config_map,
98+
submit_job_config_map,
99+
],
100+
job,
101+
})
102+
}

0 commit comments

Comments
 (0)