《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
阅读 —
·
全站 —