Операторы Airflow

Введение в Apache Airflow на Python

Mike Metzger

Data Engineer

Операторы

  • Представляют одну задачу в рабочем процессе
  • Выполняются независимо (как правило)
  • Обычно не обмениваются данными
  • Разные операторы для разных задач
@dag(
  dag_id="Example_Dag"
)
def example_dag():
  @task
  def task1():
    return "The result from task1"

  task1()

example_dag()
Введение в Apache Airflow на Python

@task (PythonOperator)

  • Выполняет Python-функцию
  • Любую функцию Python можно декорировать с помощью @task
  • Может передавать данные другим функциям и задачам
from airflow.sdk import task

@task def printme(): print("This goes in the logs!")
printme()
Введение в Apache Airflow на Python

Аргументы @task

  • Аргументы передаются в задачу так же, как в обычную функцию Python
@task
def printme(name: str):
    print(f"Hi {name} - This goes in the logs!")

printme(name='DataCamp')
# Adds: # Hi DataCamp - This goes in the logs! # to the Airflow logs
Введение в Apache Airflow на Python

@task.bash (BashOperator)

@task.bash
def bash_example():
  return "echo 'Example!'"

bash_example()
@task.bash
def run_cleanup():
  return "runcleanup.sh"

run_cleanup()
  • Выполняет команду или скрипт Bash
  • Запускает команду во временной директории
  • Позволяет задавать переменные окружения для команды
Введение в Apache Airflow на Python

Зависимости задач

  • Каждый DAG содержит набор задач для выполнения
  • Зависимости задач определяют порядок их запуска
  • Существуют разные способы задать зависимости

 

DAG Airflow с цепочкой задач, отражающей порядок зависимостей между ними

Введение в Apache Airflow на Python

Синтаксис битового сдвига

>> и <<

task1() >> task2()

# task1 completes before task2 starts
task1() >> task2() >> task3()
# task1 completes before task2 and task2 completes before task3
task1() >> task3() task2() >> task3()
# task1 and task2 can run together but both must complete before task3 runs
Введение в Apache Airflow на Python

Пример синтаксиса битового сдвига

# Download sales data before reconciling
download_sales_data() >> reconcile()

# Download inventory data before reconciling
download_inventory_data() >> reconcile()

 

  • Обе задачи должны завершиться перед сверкой, но порядок между ними не задан
Введение в Apache Airflow на Python

Давайте потренируемся!

Введение в Apache Airflow на Python

Preparing Video For Download...