Planificación basada en datos con Assets

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

El problema del tiempo

 

  • El pipeline A carga datos en bruto desde una API
  • El pipeline B crea un panel con esos datos
  • Si ambos van por cron, B podría ejecutarse antes de que A termine
  • Añadir un margen de tiempo es frágil

El problema del tiempo

Creación de canalizaciones de datos con Airflow

¿Qué es un Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Una referencia a un dato, identificada por un nombre único
  • Se puede adjuntar un URI cuando representa una entidad de datos concreta
  • Puede ser un archivo, una tabla de base de datos o cualquier dato
Creación de canalizaciones de datos con Airflow

Productores y consumidores

  • Las tareas pueden actualizar assets, lo que crea un evento de asset
  • Los DAGs pueden dispararse con actualizaciones mediante un asset schedule

 

Flujo de assets

Creación de canalizaciones de datos con Airflow

Productor: señalizar con outlets

sales_data = Asset("s3://bucket/sales/daily.parquet")

@task(outlets=[sales_data]) def write_sales(): # Write data to S3 ...

 

  • Con éxito, Airflow registra que el asset se actualizó
  • El productor no necesita saber nada de los consumidores
Creación de canalizaciones de datos con Airflow

Consumidor: planificar según un asset

sales_data = Asset("s3://bucket/sales/daily.parquet")

@dag(schedule=[sales_data]) def build_dashboard(): ...

 

  • Sin expresión cron, solo la referencia al asset
  • Airflow lanza el consumidor automáticamente tras la actualización
Creación de canalizaciones de datos con Airflow

Planificación condicional

sales = Asset("s3://bucket/sales.parquet")
inventory = Asset("s3://bucket/inventory.parquet")

# Esperar a que AMBOS assets se actualicen @dag(schedule=(sales & inventory)) def full_report(): ...
# Disparar cuando CUALQUIER asset se actualice @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Creación de canalizaciones de datos con Airflow

Verificar assets con la CLI

$ airflow assets list
name                                  | group    | uri
s3://data-lake/sales/daily.csv        | asset    | s3://data-lake/sales/daily.csv
$ airflow assets details --name "s3://data-lake/sales/daily.csv"
property_name    | property_value
=================+=====================                                         
name             | s3://data-lake/sales/daily.csv
                ...
updated_at       | 2026-04-14T10:00:50.820041Z
  • airflow assets list muestra todos los assets registrados
  • airflow assets details muestra metadatos, incluido updated_at
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...