From af4bcb3055b38cc738ea4c39398096e868e890c8 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Wed, 3 Jun 2026 12:44:55 +0100 Subject: [PATCH 1/4] add initial config job --- dags/configuration.py | 87 +++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 87 insertions(+) create mode 100644 dags/configuration.py diff --git a/dags/configuration.py b/dags/configuration.py new file mode 100644 index 0000000..a8987c8 --- /dev/null +++ b/dags/configuration.py @@ -0,0 +1,87 @@ +""" +Module containing DAG for configuration data + +Outputs: + - configuration data in parquet_datasets bucket + - configuration data in digital land db +""" + +import boto3 +from datetime import timedelta +from airflow import DAG +from airflow.operators.python import PythonOperator, Param +from airflow.providers.amazon.aws.operators.emr import EmrServerlessStartJobOperator +from emr_dags_utils import get_secrets +from utils import get_config, dag_default_args + +config = get_config() + +with DAG( + dag_id="configuration", + default_args=dag_default_args, + description="Run digital-land-builder and upload files to S3", + schedule=None, + catchup=False, + params={ + "cpu": Param(default=8192, type="integer"), + "memory": Param(default=32768, type="integer"), + }, + render_template_as_native_obj=True, + is_paused_upon_creation=False, +) as dag: + + ENV = config["env"] + EXECUTION_ROLE_ARN = get_secrets("emr_execution_role", ENV) + S3_BUCKET = f"{ENV}-pd-batch-jobs-codepackage-bucket" + S3_LOG_BUCKET = f"{ENV}-pd-batch-jobs-logs-bucket" + S3_TASKS_ENTRY_POINT = f"s3://{S3_BUCKET}/pkg/entry_script/run_config.py" + S3_WHEEL_FILE = f"s3://{S3_BUCKET}/pkg/whl_pkg/pyspark_jobs-0.1.0-py3-none-any.whl" + S3_LOG_URI = f"s3://{S3_LOG_BUCKET}/" + S3_DATA_PATH = f"s3://{ENV}-collection-data/" + + def get_tasks_emr_application_id(**context): + app_name = f"{ENV}-pd-batch-emrsl-application" + client = boto3.client("emr-serverless", region_name="eu-west-2") + response = client.list_applications(maxResults=50) + for app in response.get("applications", []): + if app["name"] == app_name: + app_id = app["id"] + context["ti"].xcom_push(key="application_id", value=app_id) + return app_id + raise ValueError(f"EMR application '{app_name}' not found") + + get_tasks_app_id = PythonOperator( + task_id="get-tasks-emr-app-id", + python_callable=get_tasks_emr_application_id, + execution_timeout=timedelta(minutes=2), + dag=dag, + ) + + assemble_tasks_emr_task = EmrServerlessStartJobOperator( + task_id="assemble-tasks", + application_id='{{ task_instance.xcom_pull(task_ids="get-tasks-emr-app-id", key="application_id") }}', + execution_role_arn=EXECUTION_ROLE_ARN, + job_driver={ + "sparkSubmit": { + "entryPoint": S3_TASKS_ENTRY_POINT, + "entryPointArguments": [ + "--env", + ENV, + "--collection-data-path", + S3_DATA_PATH, + "--parquet-datasets-path", + f"s3://{ENV}-parquet-datasets", + ], + "sparkSubmitParameters": f"--py-files {S3_WHEEL_FILE} " "--conf spark.serializer=org.apache.spark.serializer.KryoSerializer", + } + }, + configuration_overrides={"monitoringConfiguration": {"s3MonitoringConfiguration": {"logUri": S3_LOG_URI}}}, + name="assemble-tasks-job", + wait_for_completion=True, + aws_conn_id="aws_default", + waiter_max_attempts=180, + waiter_delay=60, + execution_timeout=timedelta(hours=3), + ) + + get_tasks_app_id >> assemble_tasks_emr_task From 7edb924081038db8bccc36ee95bee278e0b62307 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Wed, 3 Jun 2026 15:23:16 +0100 Subject: [PATCH 2/4] format --- dags/configuration.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/dags/configuration.py b/dags/configuration.py index a8987c8..6a9215d 100644 --- a/dags/configuration.py +++ b/dags/configuration.py @@ -6,13 +6,14 @@ - configuration data in digital land db """ -import boto3 from datetime import timedelta + +import boto3 from airflow import DAG -from airflow.operators.python import PythonOperator, Param +from airflow.operators.python import Param, PythonOperator from airflow.providers.amazon.aws.operators.emr import EmrServerlessStartJobOperator from emr_dags_utils import get_secrets -from utils import get_config, dag_default_args +from utils import dag_default_args, get_config config = get_config() From 2d175ff08a6aafd4c5b1bc4c7dc4d65dacbc8737 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Wed, 3 Jun 2026 16:48:30 +0100 Subject: [PATCH 3/4] remove Param as its not being used --- dags/configuration.py | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/dags/configuration.py b/dags/configuration.py index 6a9215d..7a84cf4 100644 --- a/dags/configuration.py +++ b/dags/configuration.py @@ -10,7 +10,7 @@ import boto3 from airflow import DAG -from airflow.operators.python import Param, PythonOperator +from airflow.operators.python import PythonOperator from airflow.providers.amazon.aws.operators.emr import EmrServerlessStartJobOperator from emr_dags_utils import get_secrets from utils import dag_default_args, get_config @@ -20,13 +20,9 @@ with DAG( dag_id="configuration", default_args=dag_default_args, - description="Run digital-land-builder and upload files to S3", + description="run processes related to our configuration files", schedule=None, catchup=False, - params={ - "cpu": Param(default=8192, type="integer"), - "memory": Param(default=32768, type="integer"), - }, render_template_as_native_obj=True, is_paused_upon_creation=False, ) as dag: From 3c741b2ee403aa5f574cff5eac187b4c4c2d67ba Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Thu, 4 Jun 2026 20:09:13 +0100 Subject: [PATCH 4/4] add debug option to script --- dags/configuration.py | 17 +++++++++-------- 1 file changed, 9 insertions(+), 8 deletions(-) diff --git a/dags/configuration.py b/dags/configuration.py index 7a84cf4..1f093f0 100644 --- a/dags/configuration.py +++ b/dags/configuration.py @@ -10,6 +10,7 @@ import boto3 from airflow import DAG +from airflow.models.param import Param from airflow.operators.python import PythonOperator from airflow.providers.amazon.aws.operators.emr import EmrServerlessStartJobOperator from emr_dags_utils import get_secrets @@ -23,6 +24,9 @@ description="run processes related to our configuration files", schedule=None, catchup=False, + params={ + "debug": Param(default=False, type="boolean", description="Enable debug logging for the Spark job"), + }, render_template_as_native_obj=True, is_paused_upon_creation=False, ) as dag: @@ -61,14 +65,11 @@ def get_tasks_emr_application_id(**context): job_driver={ "sparkSubmit": { "entryPoint": S3_TASKS_ENTRY_POINT, - "entryPointArguments": [ - "--env", - ENV, - "--collection-data-path", - S3_DATA_PATH, - "--parquet-datasets-path", - f"s3://{ENV}-parquet-datasets", - ], + "entryPointArguments": ( + f"{{{{ ['--env', '{ENV}', '--collection-data-path', '{S3_DATA_PATH}', " + f"'--parquet-datasets-path', 's3://{ENV}-parquet-datasets'] " + f"+ (['--debug'] if params.debug else []) }}}}" + ), "sparkSubmitParameters": f"--py-files {S3_WHEEL_FILE} " "--conf spark.serializer=org.apache.spark.serializer.KryoSerializer", } },