Operátory Airflow

Úvod do Apache Airflow v Pythonu

Mike Metzger

Data Engineer

Operátory

  • Představují jednu úlohu v pracovním postupu
  • Spouštějí se nezávisle (obvykle)
  • Obecně si nesdílí informace
  • Různé operátory pro různé úlohy
@dag(
  dag_id="Example_Dag"
)
def example_dag():
  @task
  def task1():
    return "The result from task1"

  task1()

example_dag()
Úvod do Apache Airflow v Pythonu

@task (PythonOperator)

  • Spustí funkci v Pythonu
  • Libovolnou funkci lze dekorovat pomocí @task
  • Umožňuje předávat data jiným funkcím / úlohám
from airflow.sdk import task

@task def printme(): print("This goes in the logs!")
printme()
Úvod do Apache Airflow v Pythonu

Argumenty @task

  • Argumenty lze předávat stejně jako běžné funkci v Pythonu
@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
Úvod do Apache Airflow v Pythonu

@task.bash (BashOperator)

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

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

run_cleanup()
  • Spustí zadaný příkaz nebo skript Bash
  • Spouští příkaz v dočasném adresáři
  • Umožňuje nastavit proměnné prostředí pro příkaz
Úvod do Apache Airflow v Pythonu

Závislosti úloh

  • Každý DAG obsahuje sadu úloh k dokončení
  • Závislosti úloh určují pořadí spuštění
  • Různé způsoby definice závislostí

 

DAG Airflow se spojenými úlohami zobrazující pořadí závislostí

Úvod do Apache Airflow v Pythonu

Syntaxe bitshift

>> a <<

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
Úvod do Apache Airflow v Pythonu

Příklad syntaxe bitshift

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

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

 

  • Obě musí být dokončeny před spuštěním reconcile, ale vzájemné pořadí není určeno
Úvod do Apache Airflow v Pythonu

Pojďme si procvičit!

Úvod do Apache Airflow v Pythonu

Preparing Video For Download...