Compartir datos entre tareas

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

¿Qué es XCom?

Flujo push/pull de XCom

 

  • XCom = comunicación entre tareas
  • Valores guardados en la base de datos de metadatos
  • Enfoque clásico: xcom_push() y xcom_pull() explícitos
Creación de canalizaciones de datos con Airflow

XComArgs: el enfoque TaskFlow

 

  • El valor de retorno de @task se convierte en un XComArg
  • XComArg es una referencia perezosa, no el dato real
  • Pasarlo a otra función crea la dependencia
  • En ejecución, Airflow lo resuelve extrayéndolo de XCom

XComArgs en TaskFlow

Creación de canalizaciones de datos con Airflow

XComArgs en código

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # dependencia + datos
Creación de canalizaciones de datos con Airflow

Mezclar clásico y TaskFlow

get_timestamp = BashOperator(
    task_id="get_timestamp",
    bash_command="date +%Y-%m-%d")

@task
def log_timestamp(ts_value):
    print(f"Timestamp: {ts_value}")


log_timestamp(get_timestamp.output)
Creación de canalizaciones de datos con Airflow

Limitaciones de XCom

Límites de tamaño de XCom por backend

 

  • El tamaño máximo depende de tu backend de base de datos
  • Debe ser serializable a JSON (dicts, listas, cadenas, números)
  • Cada valor añade carga a la base de datos
Creación de canalizaciones de datos con Airflow

Backends XCom personalizados

Enrutado de backend XCom personalizado

  • Redirige valores a un almacenamiento de objetos
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: solo guarda XComs en el almacenamiento de objetos si se alcanza el umbral
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...