Un ejemplo práctico de cómo utilizar Apache Airflow
Supongamos que queremos crear una tarea programada que se ejecute todos los días a las 9AM para enviar un correo electrónico a los miembros de un equipo con el resumen de las tareas que deben completar ese día. Para esto, podemos crear un flujo de trabajo en Apache Airflow, que se encargue de lo siguiente:
- Comprobar si hay nuevas tareas que se hayan agregado al sistema.
- Generar el resumen de tareas que deben completar los miembros del equipo.
- Enviar un correo electrónico con el resumen de tareas a los miembros del equipo.
Aquí te dejo un código que ilustra cómo sería la definición de tareas dentro de un flujo de trabajo en Apache Airflow:
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from airflow.operators.email_operator import EmailOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'team',
'start_date': datetime(2020, 1, 1),
'email': ['team@mycompany.com'],
'email_on_failure': False,
'email_on_retry': False,
'retries': 1,
'retry_delay': timedelta(minutes=5)
}
dag = DAG('daily_task_summary', default_args=default_args, schedule_interval='0 9 * * *')
task_1 = BashOperator(
task_id='check_for_new_tasks',
bash_command='python /path/to/check_for_new_tasks.py',
dag=dag
)
task_2 = BashOperator(
task_id='generate_task_summary',
bash_command='python /path/to/generate_task_summary.py',
dag=dag
)
task_3 = EmailOperator(
task_id='send_task_summary_email',
to='{{ ti.xcom_pull(task_ids="generate_task_summary", key="email_list") }}',
subject='Daily Task Summary',
html_content='{{ ti.xcom_pull(task_ids="generate_task_summary", key="task_summary") }}',
dag=dag
)
task_1 >> task_2 >> task_3
En este código, hemos definido tres tareas. La primera tarea se encarga de comprobar si hay nuevas tareas que se hayan agregado al sistema. La segunda tarea genera el resumen de tareas que deben completar los miembros del equipo y la tercera tarea envía un correo electrónico con el resumen de tareas a los miembros del equipo. La tarea 1 y la tarea 2 están conectadas mediante el operador de dependencia >>, lo que indica que la tarea 2 depende de la tarea 1 y no puede comenzar hasta que la tarea 1 se haya completado. La tarea 3 depende de la tarea 2 y se utiliza el operador xcom para pasar datos entre tareas. La dirección de correo electrónico que se utiliza para enviar el correo electrónico y el resumen de tareas se pasan a través de xcom desde la tarea 2 a la tarea 3.