Commit 76ec8e70 authored by BO ZHANG's avatar BO ZHANG 🏀
Browse files

feat(dags): 引入DAG工厂模式重构CSST流水线任务

parent 4062deb0
Loading
Loading
Loading
Loading
+0 −0

Empty file added.

+115 −0
Original line number Diff line number Diff line
from airflow import DAG
from airflow.providers.docker.operators.docker import DockerOperator
from airflow.operators.bash import BashOperator
from docker.types import Mount
from datetime import datetime, timedelta
import os

# 从环境变量获取 Harbor 仓库地址,如果未设置则使用默认值
HARBOR = os.environ.get("HARBOR", "csu-harbor.csst.nao:10443")

# 公共的 DAG 默认参数
COMMON_DEFAULT_ARGS = {
    "owner": "airflow",
    "depends_on_past": False,
    "email": ["bozhang@nao.cas.cn"],
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
}

def create_csst_pipeline_dag(
    dag_id: str,
    description: str,
    tags: list,
    image_name: str,
    default_params_input: str,
    output_subpath: str = "",
    log_tasks: list = None
) -> DAG:
    """
    用于创建 CSST pipeline DAG 的工厂函数,可以有效避免重复代码。
    
    :param dag_id: DAG 的唯一标识
    :param description: DAG 描述
    :param tags: DAG 的标签列表
    :param image_name: Docker 镜像的名称(不带 tag),例如 'csst-msc-l1-qc0'
    :param default_params_input: DAG 参数中 params.input 的默认值(JSON 字符串)
    :param output_subpath: 输出挂载目录的子路径,默认是空字符串。如果需要存入子文件夹如 '/mbi' 则传入 '/mbi'
    :param log_tasks: 运行结束后需要执行的日志处理任务字典列表。
                      格式如:[{"task_id": "...", "bash_command": "..."}]
                      在 bash_command 中可以使用 {BASE_OUTPUT_DIR} 占位符,它会被替换为动态的输出路径。
    """
    args = COMMON_DEFAULT_ARGS.copy()
    
    # 尝试将输入的 JSON 字符串解析为字典,以便作为 params 传递
    import json
    try:
        parsed_input = json.loads(default_params_input)
    except Exception:
        parsed_input = {}
        
    args["params"] = {
        "input": default_params_input,
        **parsed_input
    }

    # Airflow 3 可能会要求 schedule 参数显式传递一个有效的 cron 表达式或者 None 的等价物
    # 我们将其设为 @once 或 None。在 Airflow 3 中通常 None 就足够,但如果仍然失败,可能需要其他配置。
    with DAG(
        dag_id=dag_id,
        default_args=args,
        description=description,
        start_date=datetime(2023, 1, 1),
        schedule="@once",
        catchup=False,
        tags=tags,
    ) as dag:
        
        # 动态的输出目录(基础路径)
        # 使用 params 字典直接获取,而不是通过 fromjson 过滤器
        base_output_dir = f"/tmp/{dag_id}-{{{{ params.get('dag_run_id', run_id) }}}}"
        
        # 实际挂载到容器内的目标宿主机目录
        target_output_dir = f"{base_output_dir}{output_subpath}"
        
        # 1. 创建输出目录
        create_dirs = BashOperator(
            task_id="create_output_dirs",
            bash_command=f"mkdir -p {target_output_dir}"
        )
        
        # 2. 运行 Docker 容器
        run_pipeline = DockerOperator(
            task_id=f"run_{dag_id.replace('-', '_')}",
            image=f"{HARBOR}/csst/{image_name}:{{{{ params.get('docker_images', {{}}).get('{image_name}', 'latest') }}}}",
            command="run {{ params.input }}",
            docker_url="unix://var/run/docker.sock",
            network_mode="bridge",
            auto_remove="force",
            force_pull=True,
            mount_tmp_dir=False,
            mounts=[
                Mount(source=target_output_dir, target="/pipeline/output", type="bind"),
                Mount(source="/home/cham/PycharmProjects/csst-airflow/docker-celery-3.0.5/volumes/env", target="/env", type="bind")
            ],
            env_file="/env/common.env",
        )
        
        # 设置基本依赖
        current_node = run_pipeline
        create_dirs >> current_node
        
        # 3. 追加日志处理等自定义后续任务
        if log_tasks:
            for idx, log_task_def in enumerate(log_tasks):
                bash_cmd = log_task_def["bash_command"].replace("{BASE_OUTPUT_DIR}", base_output_dir)
                log_task = BashOperator(
                    task_id=log_task_def.get("task_id", f"process_log_{idx}"),
                    bash_command=bash_cmd
                )
                current_node >> log_task
                current_node = log_task
                
    return dag
+25 −58
Original line number Diff line number Diff line
from airflow import DAG
from airflow.providers.docker.operators.docker import DockerOperator
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
import sys
import os
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from common.dag_factory import create_csst_pipeline_dag

default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "email": ["bozhang@nao.cas.cn"],
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
    "params": {
        "input": '{"dag_id":"csst-msc-l1-mbi","dag_run_id":"csst-msc-l1-mbi-inttest","dataset":"test-msc-c9-25sqdeg-v3","instrument":"MSC","obs_type":"WIDE","obs_group":"W5","obs_id":"10100131914","detector":"09","batch_id":"inttest", "docker_images":{"csst-msc-l1-mbi":"latest"}}'
    },
}
PARAMS_INPUT = '{"dag_id":"csst-msc-l1-mbi","dag_run_id":"csst-msc-l1-mbi-inttest","dataset":"test-msc-c9-25sqdeg-v3","instrument":"MSC","obs_type":"WIDE","obs_group":"W5","obs_id":"10100131914","detector":"09","batch_id":"inttest", "docker_images":{"csst-msc-l1-mbi":"latest"}}'

with DAG(
# Airflow 需要能够直接在模块的全局命名空间中找到 DAG 对象
# 我们使用 globals() 来确保这一点
dag = create_csst_pipeline_dag(
    dag_id="csst-msc-l1-mbi",
    default_args=default_args,
    description="Run MSC L1 MBI container",
    start_date=datetime(2023, 1, 1),
    schedule=None,
    catchup=False,
    tags={"csst", "msc", "l1", "mbi"},
) as dag:
    # Create main output directory and subdirectory for MBI
    create_dirs = BashOperator(
        task_id="create_output_dirs",
        bash_command="mkdir -p /tmp/csst-msc-l1-mbi-{{ (params.input | fromjson).get('dag_run_id', run_id) }}/mbi"
    )
    
    # Run MBI container
    run_mbi = DockerOperator(
        task_id="run_msc_l1_mbi",
        image="csu-harbor.csst.nao:10443/csst/csst-msc-l1-mbi:{{ (params.input | fromjson).get('docker_images', {}).get('csst-msc-l1-mbi', 'latest') }}",
        command="run {{ params.input }}",
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
        auto_remove=True,
        force_pull=True,
        mount_tmp_dir=False,
        volumes=["/tmp/csst-msc-l1-mbi-{{ (params.input | fromjson).get('dag_run_id', run_id) }}/mbi:/pipeline/output", "/home/cham/PycharmProjects/csst-airflow/docker-celery-3.0.5/volumes/env:/env"],
        env_file="/env/common.env",
        template_fields=["image", "command", "volumes"],
    )
    
    # Concatenate all log files into one
    concatenate_logs = BashOperator(
        task_id="concatenate_logs",
        bash_command="echo '=== Concatenated Pipeline Logs ===' > /tmp/csst-msc-l1-mbi-{{ (params.input | fromjson).get('dag_run_id', run_id) }}/combined.log && find /tmp/csst-msc-l1-mbi-{{ (params.input | fromjson).get('dag_run_id', run_id) }} -name 'pipeline.log' -exec echo '=== Log from {} ===' >> /tmp/csst-msc-l1-mbi-{{ (params.input | fromjson).get('dag_run_id', run_id) }}/combined.log \; -exec cat {} >> /tmp/csst-msc-l1-mbi-{{ (params.input | fromjson).get('dag_run_id', run_id) }}/combined.log \;"
    )
    
    # Display the concatenated log content
    display_combined_logs = BashOperator(
        task_id="display_combined_logs",
        bash_command="cat /tmp/csst-msc-l1-mbi-{{ (params.input | fromjson).get('dag_run_id', run_id) }}/combined.log"
    description="Run CSST MSC L1 MBI container",
    tags=["csst", "msc", "l1", "mbi"],
    image_name="csst-msc-l1-mbi",
    default_params_input=PARAMS_INPUT,
    output_subpath="/mbi",
    log_tasks=[
        {
            "task_id": "concatenate_logs",
            "bash_command": r"echo '=== Concatenated Pipeline Logs ===' > {BASE_OUTPUT_DIR}/combined.log && find {BASE_OUTPUT_DIR} -name 'pipeline.log' -exec echo '=== Log from {} ===' >> {BASE_OUTPUT_DIR}/combined.log \; -exec cat {} >> {BASE_OUTPUT_DIR}/combined.log \;"
        },
        {
            "task_id": "display_combined_logs",
            "bash_command": "cat {BASE_OUTPUT_DIR}/combined.log"
        }
    ]
)
    
    create_dirs >> run_mbi >> concatenate_logs >> display_combined_logs
 No newline at end of file
globals()[dag.dag_id] = dag
+19 −51
Original line number Diff line number Diff line
from airflow import DAG
from airflow.providers.docker.operators.docker import DockerOperator
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
import sys
import os
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from common.dag_factory import create_csst_pipeline_dag

default_args = {
    "owner": "airflow",
    "depends_on_past": False,
    "email": ["bozhang@nao.cas.cn"],
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
    "params": {
        "input": '{"dag_id":"csst-msc-l1-qc0","dag_run_id":"csst-msc-l1-qc0-inttest","dataset":"test-msc-c9-25sqdeg-v3","instrument":"MSC","obs_type":"WIDE","obs_group":"W5","obs_id":"10100131914","detector":"09","batch_id":"inttest", "docker_images":{"csst-msc-l1-qc0":"latest"}}'
    },
}
PARAMS_INPUT = '{"dag_id":"csst-msc-l1-qc0","dag_run_id":"csst-msc-l1-qc0-inttest","dataset":"test-msc-c9-25sqdeg-v3","instrument":"MSC","obs_type":"WIDE","obs_group":"W5","obs_id":"10100131914","detector":"09","batch_id":"inttest", "docker_images":{"csst-msc-l1-qc0":"latest"}}'

with DAG(
# Airflow 需要能够直接在模块的全局命名空间中找到 DAG 对象
# 我们使用 globals() 来确保这一点
dag = create_csst_pipeline_dag(
    dag_id="csst-msc-l1-qc0",
    default_args=default_args,
    description="Run MSC L1 QC0 container",
    start_date=datetime(2023, 1, 1),
    schedule=None,
    catchup=False,
    tags={"csst", "msc", "l1", "qc0"},
) as dag:
    # Create output directory
    create_dir = BashOperator(
        task_id="create_output_dir",
        bash_command="mkdir -p /tmp/csst-msc-l1-qc0-{{ (params.input | fromjson).get('dag_run_id', run_id) }}"
    )
    
    # Run the container with volume mount
    run_qc0 = DockerOperator(
        task_id="run_msc_l1_qc0",
        image="csu-harbor.csst.nao:10443/csst/csst-msc-l1-qc0:{{ (params.input | fromjson).get('docker_images', {}).get('csst-msc-l1-qc0', 'latest') }}",
        command="run {{ params.input }}",
        docker_url="unix://var/run/docker.sock",
        network_mode="bridge",
        auto_remove=True,
        force_pull=True,
        mount_tmp_dir=False,
        volumes=["/tmp/csst-msc-l1-qc0-{{ (params.input | fromjson).get('dag_run_id', run_id) }}:/pipeline/output", "/home/cham/PycharmProjects/csst-airflow/docker-celery-3.0.5/volumes/env:/env"],
        env_file="/env/common.env",
        template_fields=["image", "command", "volumes"],
    )
    
    # Capture and display the log file content
    capture_logs = BashOperator(
        task_id="capture_logs",
        bash_command="echo '=== Pipeline Log Content ===' && cat /tmp/csst-msc-l1-qc0-{{ (params.input | fromjson).get('dag_run_id', run_id) }}/pipeline.log || echo 'Log file not found'"
    tags=["csst", "msc", "l1", "qc0"],
    image_name="csst-msc-l1-qc0",
    default_params_input=PARAMS_INPUT,
    log_tasks=[
        {
            "task_id": "capture_logs",
            "bash_command": "echo '=== Pipeline Log Content ===' && cat {BASE_OUTPUT_DIR}/pipeline.log || echo 'Log file not found'"
        }
    ]
)
    
    create_dir >> run_qc0 >> capture_logs
 No newline at end of file
globals()[dag.dag_id] = dag
+2 −0
Original line number Diff line number Diff line
@@ -84,6 +84,8 @@ x-airflow-common:
    # ===========================
    # Customize airflow.cfg below
    # ===========================
    # Custom Environment Variables for DAGs
    HARBOR: "${HARBOR:-csu-harbor.csst.nao:10443}"
    # DEFAULT_TIMEZONE
    AIRFLOW__CORE__DEFAULT_TIMEZONE: "Asia/Shanghai" # "utc"
    AIRFLOW__CORE__PARALLELISM: 1024  # Max task instances per scheduler
Loading