Skip to content

Commit fe99463

Browse files
committed
Adding support for multiple git-sync repos
1 parent 49848cf commit fe99463

3 files changed

Lines changed: 55 additions & 28 deletions

File tree

crate-hashes.json

Lines changed: 12 additions & 7 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

rust/operator-binary/src/crd/mod.rs

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,7 @@ pub const STACKABLE_LOG_DIR: &str = "/stackable/log";
7070
pub const LOG_CONFIG_DIR: &str = "/stackable/app/log_config";
7171
pub const AIRFLOW_HOME: &str = "/stackable/airflow";
7272
pub const AIRFLOW_CONFIG_FILENAME: &str = "webserver_config.py";
73+
pub const AIRFLOW_DAGS_FOLDER: &str = "/stackable/app/allDAGs";
7374

7475
pub const TEMPLATE_VOLUME_NAME: &str = "airflow-executor-pod-template";
7576
pub const TEMPLATE_LOCATION: &str = "/templates";
@@ -594,11 +595,19 @@ impl AirflowRole {
594595
format!(
595596
"cp -RL {CONFIG_PATH}/{AIRFLOW_CONFIG_FILENAME} {AIRFLOW_HOME}/{AIRFLOW_CONFIG_FILENAME}"
596597
),
598+
format!("mkdir {AIRFLOW_DAGS_FOLDER}"),
597599
// graceful shutdown part
598600
COMMON_BASH_TRAP_FUNCTIONS.to_string(),
599601
remove_vector_shutdown_file_command(STACKABLE_LOG_DIR),
600602
];
601603

604+
for (i, _) in airflow.spec.cluster_config.dags_git_sync.iter().enumerate() {
605+
command.push(
606+
format!("ln -s /stackable/app/git-{i}/current /stackable/app/allDAGs/current-{i}")
607+
.to_string(),
608+
)
609+
}
610+
602611
if resolved_product_image.product_version.starts_with("3.") {
603612
// Start-up commands have changed in 3.x.
604613
// See https://airflow.apache.org/docs/apache-airflow/3.0.1/installation/upgrading_to_airflow3.html#step-6-changes-to-your-startup-scripts and

rust/operator-binary/src/env_vars.rs

Lines changed: 34 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -15,8 +15,8 @@ use stackable_operator::{
1515

1616
use crate::{
1717
crd::{
18-
AirflowExecutor, AirflowRole, ExecutorConfig, LOG_CONFIG_DIR, STACKABLE_LOG_DIR,
19-
TEMPLATE_LOCATION, TEMPLATE_NAME,
18+
AIRFLOW_DAGS_FOLDER, AirflowExecutor, AirflowRole, ExecutorConfig, LOG_CONFIG_DIR,
19+
STACKABLE_LOG_DIR, TEMPLATE_LOCATION, TEMPLATE_NAME,
2020
authentication::{
2121
AirflowAuthenticationClassResolved, AirflowClientAuthenticationDetailsResolved,
2222
},
@@ -154,12 +154,11 @@ pub fn build_airflow_statefulset_envs(
154154
);
155155
}
156156

157-
let dags_folder = get_dags_folder(git_sync_resources);
158157
env.insert(
159158
AIRFLOW_CORE_DAGS_FOLDER.into(),
160159
EnvVar {
161160
name: AIRFLOW_CORE_DAGS_FOLDER.into(),
162-
value: Some(dags_folder),
161+
value: Some(AIRFLOW_DAGS_FOLDER.to_owned()),
163162
..Default::default()
164163
},
165164
);
@@ -288,23 +287,28 @@ pub fn build_airflow_statefulset_envs(
288287
Ok(transform_map_to_vec(env))
289288
}
290289

291-
pub fn get_dags_folder(git_sync_resources: &git_sync::v1alpha1::GitSyncResources) -> String {
292-
let git_sync_count = git_sync_resources.git_content_folders.len();
293-
if git_sync_count > 1 {
294-
tracing::warn!(
295-
"There are {git_sync_count} git-sync entries: Only the first one will be considered.",
296-
);
297-
}
298-
290+
pub fn get_dags_folder(git_sync_resources: &git_sync::v1alpha1::GitSyncResources) -> Vec<String> {
291+
// let git_sync_count = git_sync_resources.git_content_folders.len();
292+
// if git_sync_count > 1 {
293+
// tracing::warn!(
294+
// "There are {git_sync_count} git-sync entries: Only the first one will be considered.",
295+
// );
296+
// }
297+
let mut git_folders = Vec::<String>::new();
299298
// If DAG provisioning via git-sync is not configured, set a default value
300299
// so that PYTHONPATH can refer to it. N.B. nested variables need to be
301300
// resolved, so that /stackable/airflow is used instead of $AIRFLOW_HOME.
302301
// see https://airflow.apache.org/docs/apache-airflow/stable/configurations-ref.html#dags-folder
303-
git_sync_resources
304-
.git_content_folders_as_string()
305-
.first()
306-
.cloned()
307-
.unwrap_or("/stackable/airflow/dags".to_string())
302+
// git_sync_resources
303+
// .git_content_folders_as_string()
304+
// .first()
305+
// .cloned()
306+
// .unwrap_or("/stackable/airflow/dags".to_string())
307+
// TODO: Might need check weather path is correct.
308+
for folder in git_sync_resources.git_content_folders_as_string() {
309+
git_folders.push(folder)
310+
}
311+
git_folders
308312
}
309313

310314
// This set of environment variables is a standard set that is not dependent on any
@@ -314,7 +318,17 @@ fn static_envs(
314318
) -> BTreeMap<String, EnvVar> {
315319
let mut env: BTreeMap<String, EnvVar> = BTreeMap::new();
316320

317-
let dags_folder = get_dags_folder(git_sync_resources);
321+
let dags_folders = get_dags_folder(git_sync_resources);
322+
let mut dag_python_path = String::new();
323+
324+
// TODO: Might be there is a better solution to this
325+
for (i, dags_folder) in dags_folders.iter().enumerate() {
326+
dag_python_path.push_str(dags_folder);
327+
// Can't append ":" if it's last entry
328+
if i != (dags_folders.len() - 1) {
329+
dag_python_path.push_str(":");
330+
}
331+
}
318332

319333
env.insert(
320334
PYTHONPATH.into(),
@@ -323,7 +337,7 @@ fn static_envs(
323337
// dependencies can be found: this must be the actual path and not a variable.
324338
// Also include the airflow site-packages by default (for airflow and kubernetes classes etc.)
325339
name: PYTHONPATH.into(),
326-
value: Some(format!("{LOG_CONFIG_DIR}:{dags_folder}")),
340+
value: Some(format!("{LOG_CONFIG_DIR}:{dag_python_path}")),
327341
..Default::default()
328342
},
329343
);
@@ -407,12 +421,11 @@ pub fn build_airflow_template_envs(
407421

408422
// the config map also requires the dag-folder location as this will be passed on
409423
// to the pods started by airflow.
410-
let dags_folder = get_dags_folder(git_sync_resources);
411424
env.insert(
412425
AIRFLOW_CORE_DAGS_FOLDER.into(),
413426
EnvVar {
414427
name: AIRFLOW_CORE_DAGS_FOLDER.into(),
415-
value: Some(dags_folder),
428+
value: Some(AIRFLOW_DAGS_FOLDER.to_owned()),
416429
..Default::default()
417430
},
418431
);

0 commit comments

Comments
 (0)