Planificación con particiones usando Asset Partitions

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Cuando ds no basta

 

  • {{ ds }} vincula cada ejecución a una fecha; sirve para Dags programados en el tiempo
  • Pero en Dags activados por assets ds vale None; no está disponible
  • El Dag downstream no sabe qué fecha se acaba de actualizar
  • Asset Partitions propaga un partition_key de upstream a downstream

Ejecución programada upstream activa ejecución downstream basada en asset donde ds es None

Creación de canalizaciones de datos con Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    ...

 

  • Cada ejecución programada recibe un partition_key automáticamente
  • Las ejecuciones manuales no se particionan salvo que se proporcione clave
Creación de canalizaciones de datos con Airflow

Acceder a las claves de partición

En tareas de Python:

@task
def process_data(dag_run=None):
    partition_key = dag_run.partition_key
    print(f"Processing: {partition_key}")

En plantillas SQL:

DELETE FROM daily_summary
WHERE order_date = '{{ dag_run.partition_key[:10] }}';

 

  • Disponible en cada tarea dentro de una ejecución particionada
Creación de canalizaciones de datos con Airflow

Eventos de assets particionados

from airflow.sdk import dag, task, Asset, CronPartitionTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    @task(outlets=[daily_sales])
    def load_data(**context):
        ...

 

  • Outlet de la tarea + CronPartitionTimetable = evento de asset particionado
  • El evento lleva la misma clave de partición que la ejecución del Dag
  • Indica qué partición está lista, no solo que hubo actualización
Creación de canalizaciones de datos con Airflow

PartitionedAssetTimetable

from airflow.sdk import dag, task, Asset, PartitionedAssetTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=PartitionedAssetTimetable(assets=daily_sales))
def sales_report():
    @task
    def generate_report(dag_run=None):
        partition_key = dag_run.partition_key
        print(f"Report for: {partition_key}")

 

  • Se activa solo con eventos de assets particionados
  • La ejecución downstream hereda la clave de partición
Creación de canalizaciones de datos con Airflow

StartOfDayMapper

from airflow.sdk import (
    dag, PartitionedAssetTimetable,
    StartOfDayMapper,
)

@dag(schedule=PartitionedAssetTimetable(
    assets=daily_sales,
    partition_mapper_config={
        daily_sales: StartOfDayMapper()
    },
))
def sales_report():
    ...

 

Antes del mapeo: 2026-04-23T00:00:00

Después de StartOfDayMapper: 2026-04-23

$$

  • Claves temporales: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Claves no temporales: AllowedKeyMapper para regiones, departamentos
Creación de canalizaciones de datos con Airflow

El flujo completo

Flujo completo de particiones

  • Upstream: CronPartitionTimetable añade una clave de partición a cada ejecución
  • La tarea outlet emite un evento de asset particionado con esa clave
  • StartOfDayMapper normaliza el timestamp a una fecha
  • Downstream: PartitionedAssetTimetable se activa y hereda la clave mapeada
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...