การกำหนดเวลาแบบ Data-aware ด้วย Assets

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

Volker Janz

Senior Developer Advocate at Astronomer

ปัญหาด้านเวลา

 

  • Pipeline A โหลดข้อมูลดิบจาก API
  • Pipeline B สร้างแดชบอร์ดจากข้อมูลนั้น
  • หากทั้งคู่ใช้ cron B อาจรันก่อนที่ A จะเสร็จ
  • การเพิ่มช่วงเวลาบัฟเฟอร์เป็นวิธีที่ เปราะบาง

ปัญหาด้านเวลา

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

Asset คืออะไร?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • การอ้างอิงถึงข้อมูลชิ้นหนึ่ง โดยระบุด้วยชื่อที่ไม่ซ้ำกัน
  • สามารถแนบ URI เข้ากับ asset ได้เมื่อ asset นั้นแทนข้อมูลจริง
  • อาจเป็นไฟล์ ตารางฐานข้อมูล หรือข้อมูลประเภทใดก็ได้
การสร้าง Data Pipeline ด้วย Airflow

Producer และ Consumer

  • Task สามารถอัปเดต asset ซึ่งจะสร้าง asset event ขึ้น
  • DAG สามารถถูก trigger เมื่อ asset อัปเดตโดยใช้ asset schedule

 

Asset flow

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

Producer: ส่งสัญญาณด้วย outlets

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

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

 

  • เมื่อ สำเร็จ Airflow จะบันทึกว่า asset ถูกอัปเดตแล้ว
  • Producer ไม่จำเป็นต้องรู้ ว่ามี consumer ใดอยู่
การสร้าง Data Pipeline ด้วย Airflow

Consumer: กำหนดเวลาตาม asset

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

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

 

  • ไม่ต้องใช้ cron expression เพียงระบุ asset reference เท่านั้น
  • Airflow จะ trigger consumer โดยอัตโนมัติ หลังการอัปเดต
การสร้าง Data Pipeline ด้วย Airflow

การกำหนดเวลาแบบมีเงื่อนไข

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

# Wait for BOTH assets to update @dag(schedule=(sales & inventory)) def full_report(): ...
# Trigger when ANY asset updates @dag(schedule=(sales | inventory)) def quick_refresh(): ...
การสร้าง Data Pipeline ด้วย Airflow

ตรวจสอบ assets ด้วย 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 แสดง asset ทั้งหมดที่ลงทะเบียนไว้
  • airflow assets details แสดง metadata รวมถึง updated_at
การสร้าง Data Pipeline ด้วย Airflow

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

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

Preparing Video For Download...