Loading airflow/docker/dags/csst-msc-l1.trigger.py→airflow/docker/dags/csst-msc-l1.conditional_trigger.py +26 −4 Original line number Diff line number Diff line from airflow import DAG from airflow.providers.standard.operators.bash import BashOperator from airflow.providers.standard.operators.python import PythonOperator from datetime import datetime, timedelta from utils import TRIGGER_DEFAULT_ARGS Loading @@ -20,15 +21,36 @@ default_args = dict( **TRIGGER_DEFAULT_ARGS, ) def process_observation(**context): """ 处理天文观测数据的核心逻辑 """ # 1. 从上下文获取DAG参数 params = context["params"] obs_id = params["obs_id"] detector = params["detector"] print(f"🚀 开始处理观测ID: {obs_id}, 探测器: {detector}") return "great" with DAG( dag_id="csst-msc-l1.trigger", dag_id="csst-msc-l1.conditional-trigger", default_args=default_args, description="CSST MSC L1 pipeline trigger.", description="CSST MSC L1 pipeline conditional trigger.", start_date=datetime(2025, 1, 1), schedule="@daily", catchup=False, tags={"MSC", "L1", "Trigger"}, tags={"MSC", "L1", "conditional trigger"}, ) as dag: process_task = PythonOperator( task_id="process_observation_data", python_callable=process_observation, provide_context=True, # 关键!启用上下文传递 ) t1 = BashOperator( task_id="print_date", bash_command="date", Loading @@ -38,4 +60,4 @@ with DAG( bash_command="sleep 5", retries=3, ) t1 >> t2 t1 >> t2 >> process_observation Loading
airflow/docker/dags/csst-msc-l1.trigger.py→airflow/docker/dags/csst-msc-l1.conditional_trigger.py +26 −4 Original line number Diff line number Diff line from airflow import DAG from airflow.providers.standard.operators.bash import BashOperator from airflow.providers.standard.operators.python import PythonOperator from datetime import datetime, timedelta from utils import TRIGGER_DEFAULT_ARGS Loading @@ -20,15 +21,36 @@ default_args = dict( **TRIGGER_DEFAULT_ARGS, ) def process_observation(**context): """ 处理天文观测数据的核心逻辑 """ # 1. 从上下文获取DAG参数 params = context["params"] obs_id = params["obs_id"] detector = params["detector"] print(f"🚀 开始处理观测ID: {obs_id}, 探测器: {detector}") return "great" with DAG( dag_id="csst-msc-l1.trigger", dag_id="csst-msc-l1.conditional-trigger", default_args=default_args, description="CSST MSC L1 pipeline trigger.", description="CSST MSC L1 pipeline conditional trigger.", start_date=datetime(2025, 1, 1), schedule="@daily", catchup=False, tags={"MSC", "L1", "Trigger"}, tags={"MSC", "L1", "conditional trigger"}, ) as dag: process_task = PythonOperator( task_id="process_observation_data", python_callable=process_observation, provide_context=True, # 关键!启用上下文传递 ) t1 = BashOperator( task_id="print_date", bash_command="date", Loading @@ -38,4 +60,4 @@ with DAG( bash_command="sleep 5", retries=3, ) t1 >> t2 t1 >> t2 >> process_observation