Loading docker/dags/csst-msc-l1.conditional_trigger.py +26 −16 Original line number Original line Diff line number Diff line Loading @@ -4,11 +4,20 @@ from airflow.providers.standard.operators.python import PythonOperator from datetime import datetime, timedelta from datetime import datetime, timedelta from utils import TRIGGER_DEFAULT_ARGS from utils import TRIGGER_DEFAULT_ARGS from csst_dag.cli._run import run_l1_pipeline from csst_dag.cli._run import run_l1_pipeline import json # 添加默认参数 # 添加默认参数 default_args = dict( default_args = dict( params=dict( params=dict( # select DAGs dag_group="csst-msc-l1", dags=[ "csst-msc-l1-qc0", "csst-msc-l1-mbi", "csst-msc-l1-ast", "csst-msc-l1-sls", ], # data parameters # data parameters dataset="csst-msc-c9-25sqdeg-v3", dataset="csst-msc-c9-25sqdeg-v3", instrument="MSC", instrument="MSC", Loading @@ -16,8 +25,8 @@ default_args = dict( obs_group="W2", obs_group="W2", obs_id="10100232366", obs_id="10100232366", detector="09", detector="09", prc_status=None, prc_status="", qc_status=None, qc_status="", # task parameters # task parameters batch_id="airflow-batch", batch_id="airflow-batch", priority="1", priority="1", Loading @@ -27,28 +36,24 @@ default_args = dict( # submit # submit verbose=True, verbose=True, submit=False, submit=False, final_prc_status=-2, force=False, force=False, top_n=-1, top_n=-1, # select DAGs final_prc_status=-2, dags=[ "csst-msc-l1-qc0", "csst-msc-l1-mbi", "csst-msc-l1-ast", "csst-msc-l1-sls", ], dag_group="csst-l1-pipeline", ), ), **TRIGGER_DEFAULT_ARGS, **TRIGGER_DEFAULT_ARGS, ) ) def dispatch_tasks(**context): def dispatch_tasks(**context): """ # 从上下文获取DAG参数 处理天文观测数据的核心逻辑 """ # 1. 从上下文获取DAG参数 params = context["params"] params = context["params"] params["prc_status"] = ( int(params["prc_status"]) if params["prc_status"] == "" else None ) params["qc_status"] = ( int(params["qc_status"]) if params["qc_status"] == "" else None ) print(f"🚀 获得的上下文为 {context}") print(f"🚀 获得的上下文为 {context}") print(f"🚀 获得的参数为 {params}") print(f"🚀 获得的参数为 {params}") import os import os Loading @@ -63,7 +68,12 @@ def dispatch_tasks(**context): run_l1_pipeline(**params) run_l1_pipeline(**params) return "great" return json.dumps( { "prc_status": params["prc_status"], "qc_status": params["qc_status"], } ) with DAG( with DAG( Loading Loading
docker/dags/csst-msc-l1.conditional_trigger.py +26 −16 Original line number Original line Diff line number Diff line Loading @@ -4,11 +4,20 @@ from airflow.providers.standard.operators.python import PythonOperator from datetime import datetime, timedelta from datetime import datetime, timedelta from utils import TRIGGER_DEFAULT_ARGS from utils import TRIGGER_DEFAULT_ARGS from csst_dag.cli._run import run_l1_pipeline from csst_dag.cli._run import run_l1_pipeline import json # 添加默认参数 # 添加默认参数 default_args = dict( default_args = dict( params=dict( params=dict( # select DAGs dag_group="csst-msc-l1", dags=[ "csst-msc-l1-qc0", "csst-msc-l1-mbi", "csst-msc-l1-ast", "csst-msc-l1-sls", ], # data parameters # data parameters dataset="csst-msc-c9-25sqdeg-v3", dataset="csst-msc-c9-25sqdeg-v3", instrument="MSC", instrument="MSC", Loading @@ -16,8 +25,8 @@ default_args = dict( obs_group="W2", obs_group="W2", obs_id="10100232366", obs_id="10100232366", detector="09", detector="09", prc_status=None, prc_status="", qc_status=None, qc_status="", # task parameters # task parameters batch_id="airflow-batch", batch_id="airflow-batch", priority="1", priority="1", Loading @@ -27,28 +36,24 @@ default_args = dict( # submit # submit verbose=True, verbose=True, submit=False, submit=False, final_prc_status=-2, force=False, force=False, top_n=-1, top_n=-1, # select DAGs final_prc_status=-2, dags=[ "csst-msc-l1-qc0", "csst-msc-l1-mbi", "csst-msc-l1-ast", "csst-msc-l1-sls", ], dag_group="csst-l1-pipeline", ), ), **TRIGGER_DEFAULT_ARGS, **TRIGGER_DEFAULT_ARGS, ) ) def dispatch_tasks(**context): def dispatch_tasks(**context): """ # 从上下文获取DAG参数 处理天文观测数据的核心逻辑 """ # 1. 从上下文获取DAG参数 params = context["params"] params = context["params"] params["prc_status"] = ( int(params["prc_status"]) if params["prc_status"] == "" else None ) params["qc_status"] = ( int(params["qc_status"]) if params["qc_status"] == "" else None ) print(f"🚀 获得的上下文为 {context}") print(f"🚀 获得的上下文为 {context}") print(f"🚀 获得的参数为 {params}") print(f"🚀 获得的参数为 {params}") import os import os Loading @@ -63,7 +68,12 @@ def dispatch_tasks(**context): run_l1_pipeline(**params) run_l1_pipeline(**params) return "great" return json.dumps( { "prc_status": params["prc_status"], "qc_status": params["qc_status"], } ) with DAG( with DAG( Loading