Operatori Airflow

Introducere în Apache Airflow în Python

Mike Metzger

Data Engineer

Operatori

  • Reprezintă o singură sarcină dintr-un flux de lucru
  • Rulează independent (de obicei)
  • În general nu partajează informații
  • Diferiți operatori pentru sarcini diferite
@dag(
  dag_id="Example_Dag"
)
def example_dag():
  @task
  def task1():
    return "The result from task1"

  task1()

example_dag()
Introducere în Apache Airflow în Python

@task (PythonOperator)

  • Execută o funcție Python
  • Orice funcție Python poate fi decorată cu @task
  • Poate transmite date altor funcții / sarcini
from airflow.sdk import task

@task def printme(): print("This goes in the logs!")
printme()
Introducere în Apache Airflow în Python

Argumente @task

  • Permite transmiterea de argumente ca oricărei funcții Python obișnuite
@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
Introducere în Apache Airflow în Python

@task.bash (BashOperator)

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

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

run_cleanup()
  • Execută o comandă sau un script Bash dat
  • Rulează comanda într-un director temporar
  • Permite specificarea variabilelor de mediu
Introducere în Apache Airflow în Python

Dependențe între sarcini

  • Fiecare DAG are un set de sarcini de executat
  • Dependențele specifică ordinea de rulare
  • Există mai multe metode pentru a defini dependențele

 

DAG Airflow cu sarcini conectate, indicând ordinea dependențelor

Introducere în Apache Airflow în Python

Sintaxa 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
Introducere în Apache Airflow în Python

Exemplu de sintaxă bitshift

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

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

 

  • Ambele trebuie finalizate înainte de reconciliere, ordinea nu este specificată
Introducere în Apache Airflow în Python

Să exersăm!

Introducere în Apache Airflow în Python

Preparing Video For Download...