Operatory Airflow

Wprowadzenie do Apache Airflow w Pythonie

Mike Metzger

Data Engineer

Operatory

  • Reprezentują pojedyncze zadanie w przepływie pracy
  • Działają niezależnie (zazwyczaj)
  • Zazwyczaj nie współdzielą informacji
  • Różne operatory do różnych zadań
@dag(
  dag_id="Example_Dag"
)
def example_dag():
  @task
  def task1():
    return "The result from task1"

  task1()

example_dag()
Wprowadzenie do Apache Airflow w Pythonie

@task (PythonOperator)

  • Wykonuje funkcję Pythona
  • Dowolna funkcja Pythona może być udekorowana przez @task
  • Umożliwia przekazywanie danych między funkcjami/zadaniami
from airflow.sdk import task

@task def printme(): print("This goes in the logs!")
printme()
Wprowadzenie do Apache Airflow w Pythonie

Argumenty @task

  • Umożliwia przekazywanie argumentów do zadania/funkcji jak do zwykłej funkcji Pythona
@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
Wprowadzenie do Apache Airflow w Pythonie

@task.bash (BashOperator)

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

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

run_cleanup()
  • Wykonuje podane polecenie lub skrypt Bash
  • Uruchamia polecenie w katalogu tymczasowym
  • Umożliwia określenie zmiennych środowiskowych
Wprowadzenie do Apache Airflow w Pythonie

Zależności zadań

  • Każdy DAG zawiera zestaw zadań do wykonania
  • Zależności zadań określają kolejność uruchamiania
  • Istnieje kilka metod definiowania zależności

 

DAG Airflow z połączonymi zadaniami pokazujący kolejność zależności między nimi

Wprowadzenie do Apache Airflow w Pythonie

Składnia bitshift

>> and <<

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
Wprowadzenie do Apache Airflow w Pythonie

Przykład składni bitshift

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

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

 

  • Oba muszą zakończyć się przed uzgodnieniem, ale kolejność nie jest określona
Wprowadzenie do Apache Airflow w Pythonie

Czas na ćwiczenia!

Wprowadzenie do Apache Airflow w Pythonie

Preparing Video For Download...