《Airflow DAG 定义与调度》

《Airflow DAG 定义与调度》

本文演示如何在 Airflow 中定义一个 DAG 并测试运行。完整官方教程参考:https://airflow.readthedocs.io/en/latest/tutorial.html

1 定义一个 DAG

在 airflow 家目录下新建 dags 目录,编写 example.py:

from datetime import timedelta from airflow import DAG # 镜像版本(1.10.x)与当前官方文档不一致,这里使用 airflow.operators.bash from airflow.operators.bash import BashOperator from airflow.utils.dates import days_ago default_args = { 'owner': 'airflow', 'depends_on_past': False, 'email': ['alert@example.com'], 'email_on_failure': True, 'email_on_retry': True, 'retries': 3, 'retry_delay': timedelta(minutes=1), # 'queue': 'bash_queue', # 'pool': 'backfill', # 'priority_weight': 10, } dag = DAG( 'tutorial', default_args=default_args, description='A simple tutorial DAG', schedule_interval='*/2 * * * *', start_date=days_ago(2), tags=['example'], ) t1 = BashOperator( task_id='print_date', bash_command='date', dag=dag, ) t2 = BashOperator( task_id='sleep', depends_on_past=False, bash_command='sleep 5', retries=3, dag=dag, ) templated_command = """ """ t3 = BashOperator( task_id='templated', depends_on_past=False, bash_command=templated_command, params={'my_param': 'Parameter I passed in'}, dag=dag, ) t1 >> [t2, t3]

要点:

  • schedule_interval 支持 cron 表达式(如 '*/2 * * * *' 每 2 分钟执行)。
  • 通过 t1 >> [t2, t3] 声明依赖:t2、t3 依赖 t1。

2 校验脚本是否 OK

python /usr/local/airflow/dags/example.py

3 列出 DAG 与任务

# 列出当前所有 DAG airflow list_dags # 列出指定 DAG 的所有 task airflow list_tasks tutorial # print_date # sleep # templated # 以树形结果展示依赖 airflow list_tasks tutorial --tree # <Task(BashOperator): print_date> # <Task(BashOperator): templated> # <Task(BashOperator): sleep>

4 手动测试某个任务

airflow test tutorial print_date 2015-06-01 airflow test tutorial sleep 2015-06-01 airflow test tutorial templated 2015-06-01
阅读 — · 全站 —
🎸 我的歌单 0 首