Commit 6b65d38c authored by BO ZHANG's avatar BO ZHANG 🏀
Browse files

add real run_l1_pipeline

parent adcb6e0e
Loading
Loading
Loading
Loading
+27 −9
Original line number Diff line number Diff line
@@ -3,22 +3,41 @@ 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
from csst_dag.cli._run import run_l1_pipeline


# 添加默认参数
default_args = dict(
    params=dict(
        # data parameters
        dataset="csst-msc-c9-25sqdeg-v3",
        instrument="MSC",
        obs_type="WIDE",
        obs_group="W2",
        obs_id="10100232366",
        detector="09",
        pmapname="csst_000094.pmap",
        ref_cat="trilegal_093",
        prc_status=None,
        qc_status=None,
        # task parameters
        batch_id="airflow-batch",
        priority="1",
        # DAG parameters
        pmapname="csst_000094.pmap",
        ref_cat="trilegal_093",
        # submit
        verbose=True,
        submit=False,
        final_prc_status=-2,
        force=False,
        top_n=-1,
        # select DAGs
        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,
)
@@ -41,6 +60,9 @@ def dispatch_tasks(**context):
    import csst_dag

    print(os.environ)

    run_l1_pipeline(**params)

    return "great"


@@ -54,7 +76,7 @@ with DAG(
    tags={"MSC", "L1", "conditional trigger"},
) as dag:

    process_task = PythonOperator(
    opt_dispatch_tasks = PythonOperator(
        task_id="dispatch_tasks",
        python_callable=dispatch_tasks,
    )
@@ -63,9 +85,5 @@ with DAG(
        task_id="print_date",
        bash_command="date",
    )
    t2 = BashOperator(
        task_id="sleep",
        bash_command="sleep 5",
        retries=3,
    )
    t1 >> t2 >> process_task

    t1 >> opt_dispatch_tasks