From f6ce4ebc1fcebe90a4f5b43743adefc3a887d06b Mon Sep 17 00:00:00 2001 From: kena vyas Date: Wed, 14 May 2025 11:54:53 +0100 Subject: [PATCH 1/4] create performance DAG --- dags/performance_dag.py | 56 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 56 insertions(+) create mode 100644 dags/performance_dag.py diff --git a/dags/performance_dag.py b/dags/performance_dag.py new file mode 100644 index 00000000..98f9d5dd --- /dev/null +++ b/dags/performance_dag.py @@ -0,0 +1,56 @@ +from datetime import timedelta +from airflow import DAG +from airflow.operators.python import PythonOperator +from utils import dag_default_args, get_config, setup_configure_dag_callable +from airflow.providers.amazon.aws.operators.ecs import EcsRunTaskOperator + +config = get_config() +ecs_cluster = f"{config['env']}-cluster" +collection_task_name = f"{config['env']}-mwaa-collection-task" + +with DAG( + dag_id="build-performance-dataset", + default_args=dag_default_args, + description="Generate provision-quality parquet and upload to S3", + schedule=None, + catchup=False, +) as dag: + + configure_dag_task = PythonOperator( + task_id="configure-dag", + python_callable=setup_configure_dag_callable(config, collection_task_name), + dag=dag, + ) + + build_performance_task = EcsRunTaskOperator( + task_id="build-performance-task", + dag=dag, + execution_timeout=timedelta(minutes=60), + cluster=ecs_cluster, + task_definition=collection_task_name, + launch_type="FARGATE", + overrides={ + "containerOverrides": [ + { + "name": collection_task_name, + "command": ["./build-performance.sh"], + "environment": [ + {"name": "ENVIRONMENT", "value": config["env"]}, + { + "name": "COLLECTION_DATASET_BUCKET_NAME", + "value": "'{{ task_instance.xcom_pull(task_ids=\"configure-dag\", key=\"collection-dataset-bucket-name\") | string }}'" + }, + ], + } + ] + }, + network_configuration={ + "awsvpcConfiguration": '{{ task_instance.xcom_pull(task_ids="configure-dag", key="aws_vpc_config") }}' + }, + awslogs_group='{{ task_instance.xcom_pull(task_ids="configure-dag", key="collection-task-log-group") }}', + awslogs_region='{{ task_instance.xcom_pull(task_ids="configure-dag", key="collection-task-log-region") }}', + awslogs_stream_prefix='performance-build', + awslogs_fetch_interval=timedelta(seconds=1) + ) + + configure_dag_task >> build_performance_task From d33fb8408423ffac250550821353f0898820fb38 Mon Sep 17 00:00:00 2001 From: kena vyas Date: Wed, 14 May 2025 12:14:22 +0100 Subject: [PATCH 2/4] add parameters --- dags/performance_dag.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/dags/performance_dag.py b/dags/performance_dag.py index 98f9d5dd..a8d03d51 100644 --- a/dags/performance_dag.py +++ b/dags/performance_dag.py @@ -1,6 +1,7 @@ from datetime import timedelta from airflow import DAG from airflow.operators.python import PythonOperator +from airflow.models.param import Param from utils import dag_default_args, get_config, setup_configure_dag_callable from airflow.providers.amazon.aws.operators.ecs import EcsRunTaskOperator @@ -14,6 +15,11 @@ description="Generate provision-quality parquet and upload to S3", schedule=None, catchup=False, + params={ + "cpu": Param(default=8192, type="integer"), + "memory": Param(default=32768, type="integer"), + }, + is_paused_upon_creation=False, ) as dag: configure_dag_task = PythonOperator( From 2a26055904aadef01592bfdb18f837629a7f36fe Mon Sep 17 00:00:00 2001 From: kena vyas Date: Wed, 14 May 2025 12:26:52 +0100 Subject: [PATCH 3/4] add render_template_as_native_obj --- dags/performance_dag.py | 1 + 1 file changed, 1 insertion(+) diff --git a/dags/performance_dag.py b/dags/performance_dag.py index a8d03d51..7072da3c 100644 --- a/dags/performance_dag.py +++ b/dags/performance_dag.py @@ -19,6 +19,7 @@ "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 9151c9963296c735b076e7bf7893bf91a182e505 Mon Sep 17 00:00:00 2001 From: kena vyas Date: Wed, 14 May 2025 12:50:31 +0100 Subject: [PATCH 4/4] update awslogs_stream_prefix --- dags/performance_dag.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/dags/performance_dag.py b/dags/performance_dag.py index 7072da3c..1f471d40 100644 --- a/dags/performance_dag.py +++ b/dags/performance_dag.py @@ -56,7 +56,7 @@ }, awslogs_group='{{ task_instance.xcom_pull(task_ids="configure-dag", key="collection-task-log-group") }}', awslogs_region='{{ task_instance.xcom_pull(task_ids="configure-dag", key="collection-task-log-region") }}', - awslogs_stream_prefix='performance-build', + awslogs_stream_prefix='{{ task_instance.xcom_pull(task_ids="configure-dag", key="collection-task-log-stream-prefix") }}', awslogs_fetch_interval=timedelta(seconds=1) )