Передача данных между задачами с помощью XCom

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

Mike Metzger

Data Engineer

Что такое XCom?

  • «Cross Communication»
    • Позволяет задачам обмениваться данными
  • Хранится в базе метаданных Airflow
    • Передача небольших объёмов данных
    • Имена файлов, URI, количество строк

Иллюстрация передачи небольших данных между задачами Airflow через XCom

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

Что не стоит передавать через XCom

  • Большие файлы
  • DataFrames
  • Полные базы данных
  • Большие изображения

Иллюстрация типов данных, которые не следует передавать через XCom: большие файлы и DataFrames

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

Реализация XCom

  • Существует множество способов использовать XCom
  • Сосредоточимся на TaskFlow API
  • Расширение того, что уже сделано с помощью @tasks
Введение в Apache Airflow на Python

Пример XCom

@dag(dag_id='Example_XCom')
def example_xcom():

@task def get_data(): return data
@task(multiple_outputs=True) def clean_data(sourcedata): return clean(sourcedata) # Example, not implemented
clean_data(get_data()) example_xcom()
Введение в Apache Airflow на Python

Зависимости XCom

  • XCom автоматически определяет порядок зависимостей
  • Пример
    clean_data(get_data())
    
  • Концептуально то же самое, что get_data() >> clean_data()
  • Ещё один пример
     result = clean_data(get_data())
     result >> alert_when_complete()
    
Введение в Apache Airflow на Python

Просмотр данных XCom

Страница XCom в Airflow со списком сохранённых значений по ключу, DAG и задаче

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

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

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

Preparing Video For Download...