From 6a29e3187e44bffc0d59201ce70e3698a43cdab2 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Tue, 14 Apr 2026 13:12:26 +0100 Subject: [PATCH 1/9] include flood risk zone in the new collections for testing --- dags/new_collection_generator.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index 79cd063..c9ad0c5 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -32,7 +32,7 @@ collections = get_collections_dict(datasets_dict.values()) -filtered_collections = {k: v for k, v in collections.items() if k in ["central-activities-zone", "transport-access-node", "title-boundary"]} +filtered_collections = {k: v for k, v in collections.items() if k in ["central-activities-zone", "transport-access-node", "title-boundary", "flood-risk-zone"]} # read config from file and environment config = get_config() From dbe0a031dd0f62c7daeffd45fa05be9120a76cc2 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Tue, 14 Apr 2026 22:47:06 +0100 Subject: [PATCH 2/9] update missing params that have been added --- dags/new_collection_generator.py | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index c9ad0c5..4b1c360 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -69,6 +69,7 @@ "transform-batch-size": Param(default=200, type="integer"), "incremental-loading-override": Param(default=False, type="boolean"), "regenerate-log-override": Param(default=False, type="boolean"), + "force-reprocessing": Param(default=False, type="boolean"), }, render_template_as_native_obj=True, is_paused_upon_creation=False, @@ -93,6 +94,7 @@ def configure_dag(**kwargs): transform_batch_size = int(kwargs["params"].get("transform-batch-size")) incremental_loading_override = bool(kwargs["params"].get("incremental-loading-override")) regenerate_log_override = bool(kwargs["params"].get("regenerate-log-override")) + force_reprocessing = bool(kwargs["params"].get("force-reprocessing")) # Push values to XCom ti.xcom_push(key="memory", value=memory) @@ -102,6 +104,7 @@ def configure_dag(**kwargs): ti.xcom_push(key="transform-batch-size", value=transform_batch_size) ti.xcom_push(key="incremental-loading-override", value=incremental_loading_override) ti.xcom_push(key="regenerate-log-override", value=regenerate_log_override) + ti.xcom_push(key="force-reprocessing", value="True" if force_reprocessing else "") # add collection_data bucket # add collection bucket name collection_dataset_bucket_name = kwargs["conf"].get(section="custom", key="collection_dataset_bucket_name") @@ -160,6 +163,7 @@ def configure_dag(**kwargs): "value": '\'{{ task_instance.xcom_pull(task_ids="configure-dag", key="incremental-loading-override") | string }}\'', }, {"name": "REGENERATE_LOG_OVERRIDE", "value": '\'{{ task_instance.xcom_pull(task_ids="configure-dag", key="regenerate-log-override") | string }}\''}, + {"name": "REPROCESS", "value": '\'{{ task_instance.xcom_pull(task_ids="configure-dag", key="force-reprocessing") }}\''}, ], }, ] From 2a1e20f82feb3d787d72480e8705f56428e0d8c4 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Sat, 18 Apr 2026 15:01:37 +0100 Subject: [PATCH 3/9] edit EMR serverless job submission to use archives fo python dependencies --- dags/new_collection_generator.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index 4b1c360..4cb089f 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -232,6 +232,7 @@ def get_emr_application_id(**context): S3_LOG_URI = f"s3://{S3_LOG_BUCKET}/" S3_DEPENDENCIES_PATH = f"s3://{S3_BUCKET}/pkg/dependencies/dependencies.zip" S3_DATA_PATH = f"s3://{S3_SOURCE_DATA_PATH}/" + S3_ENVIRONMENT_PATH = f"s3://{S3_BUCKET}/pkg/dependencies/environment.tar.gz" assemble_emr_task = EmrServerlessStartJobOperator( task_id="assemble-emr-job", @@ -252,7 +253,11 @@ def get_emr_application_id(**context): "--parquet-datasets-path", f"s3://{ENV}-parquet-datasets", ], - "sparkSubmitParameters": f"--jars /usr/lib/spark/jars/postgresql-42.7.4.jar --py-files {S3_WHEEL_FILE},{S3_DEPENDENCIES_PATH} " + "sparkSubmitParameters": f"--jars /usr/lib/spark/jars/postgresql-42.7.4.jar " + f"--archives {S3_ENVIRONMENT_PATH}#environment " + f"--py-files {S3_WHEEL_FILE} " + "--conf spark.executorEnv.PYSPARK_PYTHON=./environment/bin/python " + "--conf spark.emr-serverless.driverEnv.PYSPARK_PYTHON=./environment/bin/python " "--conf spark.serializer=org.apache.spark.serializer.KryoSerializer " "--conf spark.kryo.registrator=org.apache.sedona.core.serde.SedonaKryoRegistrator " "--conf spark.sql.extensions=org.apache.sedona.sql.SedonaSqlExtensions", From fd33ab77fc6df16be8371c5d66ed0a8736f4a36d Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Sun, 19 Apr 2026 14:22:49 +0100 Subject: [PATCH 4/9] include python environment in the the driver option as well --- dags/new_collection_generator.py | 1 + 1 file changed, 1 insertion(+) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index 4cb089f..fea5e66 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -258,6 +258,7 @@ def get_emr_application_id(**context): f"--py-files {S3_WHEEL_FILE} " "--conf spark.executorEnv.PYSPARK_PYTHON=./environment/bin/python " "--conf spark.emr-serverless.driverEnv.PYSPARK_PYTHON=./environment/bin/python " + "--conf spark.emr-serverless.driverEnv.PYSPARK_DRIVER_PYTHON=./environment/bin/python " "--conf spark.serializer=org.apache.spark.serializer.KryoSerializer " "--conf spark.kryo.registrator=org.apache.sedona.core.serde.SedonaKryoRegistrator " "--conf spark.sql.extensions=org.apache.sedona.sql.SedonaSqlExtensions", From e172b1867ee68ac21502c13fcabba6b78cfb21d7 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Sun, 19 Apr 2026 21:33:57 +0100 Subject: [PATCH 5/9] remove driver configg for batch jobs --- dags/new_collection_generator.py | 1 - 1 file changed, 1 deletion(-) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index fea5e66..4cb089f 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -258,7 +258,6 @@ def get_emr_application_id(**context): f"--py-files {S3_WHEEL_FILE} " "--conf spark.executorEnv.PYSPARK_PYTHON=./environment/bin/python " "--conf spark.emr-serverless.driverEnv.PYSPARK_PYTHON=./environment/bin/python " - "--conf spark.emr-serverless.driverEnv.PYSPARK_DRIVER_PYTHON=./environment/bin/python " "--conf spark.serializer=org.apache.spark.serializer.KryoSerializer " "--conf spark.kryo.registrator=org.apache.sedona.core.serde.SedonaKryoRegistrator " "--conf spark.sql.extensions=org.apache.sedona.sql.SedonaSqlExtensions", From a62f874f8c06774e569537261365b409033d434f Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Sun, 19 Apr 2026 22:12:02 +0100 Subject: [PATCH 6/9] corect how archives is being passed inn --- dags/new_collection_generator.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index 4cb089f..76482fd 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -254,10 +254,11 @@ def get_emr_application_id(**context): f"s3://{ENV}-parquet-datasets", ], "sparkSubmitParameters": f"--jars /usr/lib/spark/jars/postgresql-42.7.4.jar " - f"--archives {S3_ENVIRONMENT_PATH}#environment " + f"--conf spark.archives={S3_ENVIRONMENT_PATH}#environment " f"--py-files {S3_WHEEL_FILE} " "--conf spark.executorEnv.PYSPARK_PYTHON=./environment/bin/python " "--conf spark.emr-serverless.driverEnv.PYSPARK_PYTHON=./environment/bin/python " + "--conf spark.emr-serverless.driverEnv.PYSPARK_DRIVER_PYTHON=./environment/bin/python " "--conf spark.serializer=org.apache.spark.serializer.KryoSerializer " "--conf spark.kryo.registrator=org.apache.sedona.core.serde.SedonaKryoRegistrator " "--conf spark.sql.extensions=org.apache.sedona.sql.SedonaSqlExtensions", From 0195dd03772229fb063681d0211cd6f71ac41d17 Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Sun, 19 Apr 2026 22:30:31 +0100 Subject: [PATCH 7/9] remove conf and rely on custom image instead --- dags/new_collection_generator.py | 4 ---- 1 file changed, 4 deletions(-) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index 76482fd..fe34389 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -254,11 +254,7 @@ def get_emr_application_id(**context): f"s3://{ENV}-parquet-datasets", ], "sparkSubmitParameters": f"--jars /usr/lib/spark/jars/postgresql-42.7.4.jar " - f"--conf spark.archives={S3_ENVIRONMENT_PATH}#environment " f"--py-files {S3_WHEEL_FILE} " - "--conf spark.executorEnv.PYSPARK_PYTHON=./environment/bin/python " - "--conf spark.emr-serverless.driverEnv.PYSPARK_PYTHON=./environment/bin/python " - "--conf spark.emr-serverless.driverEnv.PYSPARK_DRIVER_PYTHON=./environment/bin/python " "--conf spark.serializer=org.apache.spark.serializer.KryoSerializer " "--conf spark.kryo.registrator=org.apache.sedona.core.serde.SedonaKryoRegistrator " "--conf spark.sql.extensions=org.apache.sedona.sql.SedonaSqlExtensions", From 2ded468a1edcaa58d50b4f405ad2a47648d058da Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Thu, 14 May 2026 09:47:51 +0100 Subject: [PATCH 8/9] correct variabe name for package ecs task --- dags/new_collection_generator.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index fe34389..df064ad 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -293,7 +293,7 @@ def get_emr_application_id(**context): {"name": "COLLECTION_NAME", "value": collection}, {"name": "DATASET_NAME", "value": dataset}, { - "name": "COLLECTION_DATA_BUCKET", + "name": "COLLECTION_DATASET_BUCKET_NAME", "value": '\'{{ task_instance.xcom_pull(task_ids="configure-dag", key="collection-dataset-bucket-name") | string }}\'', }, { From 30c11f17cf7492e08084c4a4eeb7830bead3ca1f Mon Sep 17 00:00:00 2001 From: eveleighoj <35256612+eveleighoj@users.noreply.github.com> Date: Thu, 14 May 2026 12:55:51 +0100 Subject: [PATCH 9/9] revert name change as it was a scrript error --- dags/new_collection_generator.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dags/new_collection_generator.py b/dags/new_collection_generator.py index df064ad..fe34389 100644 --- a/dags/new_collection_generator.py +++ b/dags/new_collection_generator.py @@ -293,7 +293,7 @@ def get_emr_application_id(**context): {"name": "COLLECTION_NAME", "value": collection}, {"name": "DATASET_NAME", "value": dataset}, { - "name": "COLLECTION_DATASET_BUCKET_NAME", + "name": "COLLECTION_DATA_BUCKET", "value": '\'{{ task_instance.xcom_pull(task_ids="configure-dag", key="collection-dataset-bucket-name") | string }}\'', }, {