การกำหนดเวลาแบบ Partition-aware ด้วย Asset Partitions

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

เมื่อ ds ไม่เพียงพอ

 

  • {{ ds }} ผูกแต่ละ run เข้ากับวันที่ ซึ่งเหมาะกับ Dag ที่ กำหนดเวลา (scheduled)
  • แต่ Dag ที่ ถูกทริกเกอร์ด้วย asset จะมี ds เป็น None ไม่สามารถใช้งานได้
  • Dag ปลายทางไม่มีทางรู้ได้ว่า วันที่ใด ถูกอัปเดตไป
  • Asset Partitions ส่งต่อ partition_key จาก upstream ไปยัง downstream

Upstream ที่กำหนดเวลาทริกเกอร์ downstream ที่กำหนดเวลาด้วย asset โดย ds เป็น None

การสร้าง Data Pipeline ด้วย Airflow

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • ทุก scheduled run จะได้รับ partition_key โดยอัตโนมัติ
  • การรันแบบ manual จะ ไม่ ถูก partition ยกเว้นระบุ key ไว้
การสร้าง Data Pipeline ด้วย Airflow

การเข้าถึง partition key

ใน Python task:

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

ใน SQL template:

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

 

  • ใช้ได้ใน ทุก task ภายใน partitioned run
การสร้าง Data Pipeline ด้วย Airflow

Partitioned asset event

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 ของ task + CronPartitionTimetable = partitioned asset event
  • event นั้นจะพา partition key เดียวกัน กับ Dag run ไปด้วย
  • บ่งบอกว่า partition ใด พร้อมใช้งาน ไม่ใช่แค่ว่าข้อมูลอัปเดตแล้ว
การสร้าง Data Pipeline ด้วย 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}")

 

  • ทริกเกอร์ เฉพาะ เมื่อมี partitioned asset event
  • Downstream run จะ รับช่วง partition key ต่อโดยอัตโนมัติ
การสร้าง Data Pipeline ด้วย Airflow

StartOfDayMapper

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

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

 

ก่อน mapping: 2026-04-23T00:00:00

หลัง StartOfDayMapper: 2026-04-23

$$

  • key แบบเวลา: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • key แบบไม่ใช่เวลา: AllowedKeyMapper สำหรับภูมิภาค, แผนก
การสร้าง Data Pipeline ด้วย Airflow

ขั้นตอนทั้งหมด

ขั้นตอน partition flow ทั้งหมด

  • Upstream: CronPartitionTimetable แนบ partition key กับแต่ละ run
  • Outlet task ปล่อย partitioned asset event พร้อม key นั้น
  • StartOfDayMapper แปลง timestamp ให้เป็น วันที่
  • Downstream: PartitionedAssetTimetable ทริกเกอร์และ รับช่วง mapped key
การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...