Skip to content

Commit 71d04c5

Browse files
committed
refactor: consolidate logging for connect
1 parent 53829d8 commit 71d04c5

2 files changed

Lines changed: 46 additions & 4 deletions

File tree

rust/operator-binary/src/connect/controller/build/executor.rs

Lines changed: 42 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,7 @@
1-
use std::collections::{BTreeMap, HashMap};
1+
use std::{
2+
collections::{BTreeMap, HashMap},
3+
str::FromStr,
4+
};
25

36
use snafu::{ResultExt, Snafu};
47
use stackable_operator::{
@@ -15,14 +18,23 @@ use stackable_operator::{
1518
},
1619
kube::ResourceExt,
1720
product_logging::framework::{VECTOR_CONFIG_FILE, calculate_log_volume_size_limit},
18-
v2::{builder::pod::container::new_container_builder, role_utils::JavaCommonConfig},
21+
v2::{
22+
builder::pod::container::{EnvVarSet, new_container_builder},
23+
product_logging::framework::vector_container,
24+
role_group_utils::ResourceNames,
25+
role_utils::JavaCommonConfig,
26+
types::operator::{RoleGroupName, RoleName},
27+
},
1928
};
2029

2130
use crate::{
2231
connect::{
2332
common::{self, SparkConnectRole, object_name},
2433
controller::validate::ValidatedSparkConnectServer,
25-
crd::{SparkConnectContainer, v1alpha1},
34+
crd::{
35+
CONNECT_EXECUTOR_ROLE_NAME, DEFAULT_SPARK_CONNECT_GROUP_NAME, SparkConnectContainer,
36+
v1alpha1,
37+
},
2638
s3,
2739
},
2840
crd::constants::{
@@ -155,7 +167,33 @@ pub fn executor_pod_template(
155167
.context(AddVolumeSnafu)?;
156168
}
157169

158-
let mut result = template.add_container(container.build()).build_template();
170+
template.add_container(container.build());
171+
172+
// Vector log-aggregation sidecar (symmetric with the server), added when the executor enables
173+
// the Vector agent.
174+
if let Some(vector_log_config) = &validated.executor_logging.vector_container {
175+
// The Vector sidecar's `CLUSTER_NAME`/`ROLE_NAME`/`ROLE_GROUP_NAME` log-metadata env vars.
176+
// These do NOT affect resource naming: Spark Connect keeps its `{cluster}-{role}` names.
177+
let vector_resource_names = ResourceNames {
178+
cluster_name: validated.name.clone(),
179+
role_name: RoleName::from_str(CONNECT_EXECUTOR_ROLE_NAME)
180+
.expect("CONNECT_EXECUTOR_ROLE_NAME is a valid role name"),
181+
role_group_name: RoleGroupName::from_str(DEFAULT_SPARK_CONNECT_GROUP_NAME)
182+
.expect("DEFAULT_SPARK_CONNECT_GROUP_NAME is a valid role group name"),
183+
};
184+
185+
template.add_container(vector_container(
186+
&SparkConnectContainer::Vector.to_container_name(),
187+
resolved_product_image,
188+
vector_log_config,
189+
&vector_resource_names,
190+
&VOLUME_MOUNT_NAME_CONFIG,
191+
&VOLUME_MOUNT_NAME_LOG,
192+
EnvVarSet::new(),
193+
));
194+
}
195+
196+
let mut result = template.build_template();
159197

160198
// Merge user provided pod spec if any
161199
result.merge_from(validated.executor_overrides.pod_overrides.clone());

rust/operator-binary/src/connect/controller/validate.rs

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -142,6 +142,7 @@ pub struct ValidatedSparkConnectServer {
142142
pub server_logging: ValidatedLogging,
143143
pub executor_config: v1alpha1::ExecutorConfig,
144144
pub executor_overrides: ValidatedOverrides,
145+
pub executor_logging: ValidatedLogging,
145146
}
146147

147148
/// User-provided overrides for a role, captured during validation so the resource builders never
@@ -341,6 +342,8 @@ pub fn validate(
341342

342343
let server_logging =
343344
validate_logging(&server_config.logging, &vector_aggregator_config_map_name)?;
345+
let executor_logging =
346+
validate_logging(&executor_config.logging, &vector_aggregator_config_map_name)?;
344347

345348
Ok(ValidatedSparkConnectServer {
346349
metadata: scs.meta().clone(),
@@ -359,5 +362,6 @@ pub fn validate(
359362
server_logging,
360363
executor_config,
361364
executor_overrides,
365+
executor_logging,
362366
})
363367
}

0 commit comments

Comments
 (0)