Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
via `spec.webservers.roleConfig.trustedProxies` ([#835]).
- Added airflow `3.3.1` ([#865]).
- Add `/ready` endpoint to the operator Deployment, which reports the CRD installation status ([#868]).
- Webservers now have a default affinity to the OPA Pods when OPA authorization is configured ([#872]).

### Changed

Expand Down Expand Up @@ -81,6 +82,7 @@
[#865]: https://github.com/stackabletech/airflow-operator/pull/865
[#868]: https://github.com/stackabletech/airflow-operator/pull/868
[#870]: https://github.com/stackabletech/airflow-operator/pull/870
[#872]: https://github.com/stackabletech/airflow-operator/pull/872

## [26.7.0] - 2026-07-21

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@ The default affinities created by the operator are:

1. Co-locate all the Airflow Pods (weight 20)
2. Distribute all Pods within the same role (worker, webserver, scheduler) (weight 70)
3. If OPA authorization is configured: co-locate the webservers with the OPA Pods (weight 50).
Only the webservers are configured with the OPA auth manager and thus have the affinity.

== Kubernetes executors

Expand Down
12 changes: 8 additions & 4 deletions rust/operator-binary/src/controller/validate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,8 @@ pub fn validate_cluster(
},
);

let default_config = AirflowConfig::default_config(&airflow.name_any(), &role);
let default_config =
AirflowConfig::default_config(&airflow.name_any(), &role, airflow.get_opa_config());

let mut group_configs = BTreeMap::new();
for (rolegroup_name, rolegroup) in &resolved_role.role_groups {
Expand Down Expand Up @@ -431,7 +432,8 @@ mod tests {
let role = cluster
.get_role(&AirflowRole::Webserver)
.expect("webserver role");
let default_config = AirflowConfig::default_config("airflow", &AirflowRole::Webserver);
let default_config =
AirflowConfig::default_config("airflow", &AirflowRole::Webserver, None);
let rolegroup = role.role_groups.get("default").expect("default role group");

let validated = validate_role_group(
Expand Down Expand Up @@ -530,7 +532,8 @@ mod tests {
let role = cluster
.get_role(&AirflowRole::Scheduler)
.expect("scheduler role");
let default_config = AirflowConfig::default_config("airflow", &AirflowRole::Scheduler);
let default_config =
AirflowConfig::default_config("airflow", &AirflowRole::Scheduler, None);
let rolegroup = role.role_groups.get("default").expect("default role group");

let validated = validate_role_group(
Expand Down Expand Up @@ -607,7 +610,8 @@ mod tests {
let role = cluster
.get_role(&AirflowRole::Webserver)
.expect("webserver role");
let default_config = AirflowConfig::default_config("airflow", &AirflowRole::Webserver);
let default_config =
AirflowConfig::default_config("airflow", &AirflowRole::Webserver, None);
let rolegroup = role.role_groups.get("default").expect("default role group");

let validated = validate_role_group(
Expand Down
117 changes: 88 additions & 29 deletions rust/operator-binary/src/crd/affinity.rs
Original file line number Diff line number Diff line change
@@ -1,31 +1,61 @@
use stackable_operator::{
commons::affinity::{
StackableAffinityFragment, affinity_between_cluster_pods, affinity_between_role_pods,
commons::{
affinity::{
StackableAffinityFragment, affinity_between_cluster_pods, affinity_between_role_pods,
},
opa::OpaConfig,
},
k8s_openapi::api::core::v1::{PodAffinity, PodAntiAffinity},
};

use crate::crd::{APP_NAME, AirflowRole};

/// Used for all [`AirflowRole`]s besides executors.
pub fn get_affinity(cluster_name: &str, role: &AirflowRole) -> StackableAffinityFragment {
get_affinity_for_role(cluster_name, &role.to_string())
pub fn get_affinity(
cluster_name: &str,
role: &AirflowRole,
opa_config: Option<&OpaConfig>,
) -> StackableAffinityFragment {
let opa_config = match role {
// Only the webserver is configured with the OPA auth manager, so only the webserver is
// co-located with the OPA Pods.
AirflowRole::Webserver => opa_config,
AirflowRole::Scheduler
| AirflowRole::Worker
| AirflowRole::DagProcessor
| AirflowRole::Triggerer => None,
};
get_affinity_for_role(cluster_name, &role.to_string(), opa_config)
}

/// There is no [`AirflowRole`] for executors (only for workers), so let's have a special case here.
pub fn get_executor_affinity(cluster_name: &str) -> StackableAffinityFragment {
get_affinity_for_role(cluster_name, "executor")
get_affinity_for_role(cluster_name, "executor", None)
}

fn get_affinity_for_role(cluster_name: &str, role: &str) -> StackableAffinityFragment {
fn get_affinity_for_role(
cluster_name: &str,
role: &str,
opa_config: Option<&OpaConfig>,
) -> StackableAffinityFragment {
// Built before the `let`s below, which shadow the helper functions with their results.
let affinity_to_opa_pods = opa_config.map(|opa_config| {
affinity_between_role_pods(
"opa",
&opa_config.config_map_name, // The discovery cm has the same name as the OpaCluster itself
"server",
50,
)
});
let affinity_between_cluster_pods = affinity_between_cluster_pods(APP_NAME, cluster_name, 20);
let affinity_between_role_pods = affinity_between_role_pods(APP_NAME, cluster_name, role, 70);

let mut pod_affinities = vec![affinity_between_cluster_pods];
pod_affinities.extend(affinity_to_opa_pods);

StackableAffinityFragment {
pod_affinity: Some(PodAffinity {
preferred_during_scheduling_ignored_during_execution: Some(vec![
affinity_between_cluster_pods,
]),
preferred_during_scheduling_ignored_during_execution: Some(pod_affinities),
required_during_scheduling_ignored_during_execution: None,
}),
pod_anti_affinity: Some(PodAntiAffinity {
Expand Down Expand Up @@ -90,6 +120,10 @@ mod tests {
redis:
host: airflow-redis-master
credentialsSecretName: airflow-redis-credentials
authorization:
opa:
configMapName: simple-opa
package: airflow
webservers:
roleGroups:
default:
Expand All @@ -111,36 +145,61 @@ mod tests {
let resolved_role = airflow
.get_role(&role)
.expect("the role is defined in the test cluster");
let default_config = AirflowConfig::default_config(&airflow.name_any(), &role);
let default_config =
AirflowConfig::default_config(&airflow.name_any(), &role, airflow.get_opa_config());
let rolegroup = resolved_role
.role_groups
.get("default")
.expect("the 'default' role group is defined in the test cluster");

let mut expected_pod_affinities = vec![WeightedPodAffinityTerm {
pod_affinity_term: PodAffinityTerm {
label_selector: Some(LabelSelector {
match_expressions: None,
match_labels: Some(BTreeMap::from([
("app.kubernetes.io/name".to_string(), "airflow".to_string()),
(
"app.kubernetes.io/instance".to_string(),
"airflow".to_string(),
),
])),
}),
topology_key: "kubernetes.io/hostname".to_string(),
..PodAffinityTerm::default()
},
weight: 20,
}];
// Only the webserver is configured with the OPA auth manager.
if role == AirflowRole::Webserver {
expected_pod_affinities.push(WeightedPodAffinityTerm {
pod_affinity_term: PodAffinityTerm {
label_selector: Some(LabelSelector {
match_expressions: None,
match_labels: Some(BTreeMap::from([
("app.kubernetes.io/name".to_string(), "opa".to_string()),
(
"app.kubernetes.io/instance".to_string(),
"simple-opa".to_string(),
),
(
"app.kubernetes.io/component".to_string(),
"server".to_string(),
),
])),
}),
topology_key: "kubernetes.io/hostname".to_string(),
..PodAffinityTerm::default()
},
weight: 50,
});
}

let expected: StackableAffinity = StackableAffinity {
node_affinity: None,
node_selector: None,
pod_affinity: Some(PodAffinity {
required_during_scheduling_ignored_during_execution: None,
preferred_during_scheduling_ignored_during_execution: Some(vec![
WeightedPodAffinityTerm {
pod_affinity_term: PodAffinityTerm {
label_selector: Some(LabelSelector {
match_expressions: None,
match_labels: Some(BTreeMap::from([
("app.kubernetes.io/name".to_string(), "airflow".to_string()),
(
"app.kubernetes.io/instance".to_string(),
"airflow".to_string(),
),
])),
}),
topology_key: "kubernetes.io/hostname".to_string(),
..PodAffinityTerm::default()
},
weight: 20,
},
]),
preferred_during_scheduling_ignored_during_execution: Some(expected_pod_affinities),
}),
pod_anti_affinity: Some(PodAntiAffinity {
required_during_scheduling_ignored_during_execution: None,
Expand Down
18 changes: 16 additions & 2 deletions rust/operator-binary/src/crd/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -415,6 +415,16 @@ impl HasStatusCondition for v1alpha2::AirflowCluster {
}

impl v1alpha2::AirflowCluster {
/// The OPA config, if OPA authorization is configured.
pub fn get_opa_config(&self) -> Option<&OpaConfig> {
self.spec
.cluster_config
.authorization
.as_ref()
.and_then(|authorization| authorization.opa.as_ref())
.map(|opa| &opa.opa)
}

/// The name of the group-listener provided for a specific role.
/// Webservers will use this group listener so that only one load balancer
/// is needed for that role.
Expand Down Expand Up @@ -967,11 +977,15 @@ pub struct AirflowConfig {
}

impl AirflowConfig {
pub(crate) fn default_config(cluster_name: &str, role: &AirflowRole) -> AirflowConfigFragment {
pub(crate) fn default_config(
cluster_name: &str,
role: &AirflowRole,
opa_config: Option<&OpaConfig>,
) -> AirflowConfigFragment {
AirflowConfigFragment {
resources: default_resources(role),
logging: product_logging::spec::default_logging(),
affinity: get_affinity(cluster_name, role),
affinity: get_affinity(cluster_name, role, opa_config),
graceful_shutdown_timeout: Some(match role {
AirflowRole::Webserver
| AirflowRole::Scheduler
Expand Down
Loading