← คู่มือ Airflow — ทำไมงานหลายขั้นต้องมีคนคุมลำดับ
LEVEL 3 · ระดับสูง
ข้อมูลมาเมื่อไหร่ค่อยรัน: Asset แทนนาฬิกา
รอบเดือนกันยายนในหัวข้อ 13.3 พังเพราะไฟล์สาขา ค ยังไม่มา · นาฬิกาปลุกท่อตอนเที่ยงคืนวันที่ 1 ตุลาคมตรงเวลา แต่ไฟล์ไม่ได้มาตามนาฬิกา · ทางที่มีอยู่แล้วไม่พอทั้งคู่: retries ของบทที่ 10 รอได้แค่ในโควตา และ skipped ของบทที่ 9 ใช้ได้เฉพาะเมื่อไม่มีไฟล์สักสาขา · บทนี้เปลี่ยนคำถามจาก "ถึงเวลาหรือยัง" เป็น "ข้อมูลมาหรือยัง"
ลองเดาก่อน: ถ้าให้ท่อวนถามทุก 10 วินาทีว่าไฟล์มาหรือยัง ระหว่างที่รอสามวัน เครื่องเสียอะไรไป
ตัวคุมที่รอนาฬิกา ต้องเดาว่าข้อมูลจะมาเมื่อไหร่ · ตัวคุมที่รอข้อมูล ไม่ต้องเดา
14.1 วนถามเอง: sensor
💡 sensor — ขั้นที่ไม่ได้ทำงานอะไร นอกจากถามซ้ำ ๆ ว่าเงื่อนไขเป็นจริงหรือยัง · @task.sensor เปลี่ยนฟังก์ชันที่คืน True/False ให้เป็นขั้นแบบนี้
stock_wait.py ถามทุก 10 วินาทีว่าไฟล์ของเดือนครบสามสาขาหรือยัง · ค่าตั้งต้นของเดือนคือกรกฎาคม (ครบแล้ว) ตัวตรวจรันได้ทันที
# ~/airflow-lab/dags/stock_wait.py — sensor: วนถามว่าไฟล์ครบสามสาขาหรือยัง คืนช่องให้ระหว่างรอ
from datetime import datetime
from pathlib import Path
from airflow.sdk import dag, task
INBOX = Path.home() / "airflow-lab" / "inbox"
@dag(schedule=None, start_date=datetime(2026, 5, 1), catchup=False, params={"month": "2026-07"})
def stock_wait():
@task.sensor(poke_interval=10, timeout=3 * 24 * 3600, mode="reschedule") # แล็บถามทุก 10 วินาที
def all_files(**context):
month = context["params"]["month"]
found = [b for b in "abc" if (INBOX / f"stock_{month}_{b}.xlsx").exists()]
print(f"เดือน {month}: มาแล้ว {len(found)} จาก 3 สาขา")
return len(found) == 3
@task
def ready(**context):
print(f"ไฟล์เดือน {context['params']['month']} ครบแล้ว")
all_files() >> ready()
stock_wait()สั่งเดือนกันยายนขณะที่ไฟล์สาขา ค ยังอยู่ใน late/ แล้วย้ายไฟล์เข้าหลัง 40 วินาที
บันทึกจากแล็บ
10:14:59 all_files up_for_reschedule 1 · ready none 0
10:15:09 all_files up_for_reschedule 1 · ready none 0
10:15:19 all_files up_for_reschedule 1 · ready none 0
10:15:29 all_files up_for_reschedule 1 · ready none 0 ← ย้ายไฟล์สาขา ค เข้า inbox
10:15:39 all_files success 1 · ready success 1
บันทึกของ all_files: "มาแล้ว 2 จาก 3 สาขา" ×4 · "มาแล้ว 3 จาก 3 สาขา" ×1up_for_reschedule คือสถานะใหม่ของเล่ม: ถามแล้วยังไม่จริง → ปล่อยช่องคืน รอ 10 วินาทีแล้วถามใหม่ · ครั้งที่ลองยังเป็น 1 ตลอด เพราะไม่ได้พัง · หน้า Sensors ทางการแยกสองโหมด: poke (ค่าตั้งต้น) ถือช่องของ worker ไว้ตลอดเวลาที่รอ · reschedule ถือเฉพาะตอนถาม · แล็บรันทั้งสองโหมดคู่กัน แบบ poke เป็น running ค้างตลอดการรอ — ถ้ารอสามวัน worker หนึ่งตัวในสามสิบสองตัวหายไปสามวัน (บทที่ 5)
sensor ยังเป็นการเดาอยู่ดี — เดาว่าควรถามถี่แค่ไหน และนานเท่าไหร่ถึงยอมแพ้ (timeout) · และคนที่รู้ว่าไฟล์มาแล้วคือคนส่งไฟล์ ไม่ใช่ sensor
14.2 ให้คนส่งบอก: Asset
💡 Asset — ชื่อของข้อมูลชิ้นหนึ่งที่ DAG หนึ่งผลิต และ DAG อื่นรอใช้ · ผู้ผลิตจบงานแล้ว Airflow บันทึก "ของชิ้นนี้อัปเดตแล้ว" ลงฐานข้อมูล · ผู้ใช้ที่ตั้ง schedule เป็น asset ตัวนั้นถูกปลุกเอง ไม่มีนาฬิกาเกี่ยว
แต่ไฟล์สามสาขาของเดือนเดียวกันต้องมาครบก่อน · ไฟล์สาขา ก ของตุลาคมกับสาขา ข ของกันยายนรวมกันไม่ได้ ⇒ ใช้ asset แบบแบ่งพาร์ทิชัน (เพิ่มในรุ่น 3.2.0 · ป้ายรุ่นตามหน้า Assets ทางการ): เหตุการณ์ของ asset พกกุญแจพาร์ทิชันไปด้วย ในท่อนี้กุญแจคือเดือน
ผู้ผลิตมีสาขาละหนึ่ง asset · ระบบรับไฟล์ของสาขาสั่งรันผู้ผลิตของสาขานั้นพร้อมเดือน (ในแล็บสั่งด้วย airflow dags trigger · ของจริงยิง REST API) · ผู้ผลิตยันว่าไฟล์อยู่ในกล่องจริงก่อนประกาศ
# ~/airflow-lab/dags/stock_files.py — ผู้ผลิต: สาขาละหนึ่ง asset · ถูกสั่งเมื่อไฟล์ของสาขามาถึง แล้วประกาศเดือนของไฟล์นั้น
from pathlib import Path
from airflow.sdk import PartitionedAtRuntime, asset
INBOX = Path.home() / "airflow-lab" / "inbox"
def announce(branch, context, outlet_events, self):
month = context["dag_run"].conf["month"]
path = INBOX / f"stock_{month}_{branch}.xlsx"
if not path.exists(): # สั่งประกาศทั้งที่ไฟล์ยังไม่มา = สั่งผิด ไม่ใช่ไฟล์มาช้า
raise FileNotFoundError(f"ยังไม่มี {path.name} ในกล่องรับไฟล์")
outlet_events[self].add_partitions(month) # ประกาศว่า asset ของสาขานี้ มีของเดือนนี้แล้ว
print(f"ประกาศ: ไฟล์สาขา {branch} เดือน {month} มาแล้ว")
@asset(uri="file://inbox/stock-a", schedule=PartitionedAtRuntime())
def stock_file_a(self, context, outlet_events):
announce("a", context, outlet_events, self)
@asset(uri="file://inbox/stock-b", schedule=PartitionedAtRuntime())
def stock_file_b(self, context, outlet_events):
announce("b", context, outlet_events, self)
@asset(uri="file://inbox/stock-c", schedule=PartitionedAtRuntime())
def stock_file_c(self, context, outlet_events):
announce("c", context, outlet_events, self)ผู้ใช้คือ stock_monthly ฉบับใหม่ที่เปลี่ยนแค่ schedule กับ month_of · PartitionedAssetTimetable ที่ต่อสาม asset ด้วย & ปลุกท่อเมื่อทั้งสามประกาศกุญแจเดียวกัน และเดือนมาถึงงานทาง dag_run.partition_key
# ~/airflow-lab/dags/stock_monthly.py — บทที่ 14: รันเมื่อไฟล์ครบสามสาขาของเดือนเดียวกัน ไม่ใช่เมื่อถึงเวลา
from datetime import datetime, timedelta
from pathlib import Path
import stock_gate
import stock_load
import stock_steps as steps
from airflow.sdk import Asset, PartitionedAssetTimetable, dag, task
from airflow.sdk.exceptions import AirflowFailException, AirflowSkipException
LAB = Path.home() / "airflow-lab"
def month_of(context):
"""เดือนที่รอบนี้ประมวลผล: -c ก่อน · แล้วพาร์ทิชันของ asset · แล้วต้นช่วงข้อมูล · ไม่มีเลย = หยุด"""
if context["params"]["month"]:
return context["params"]["month"]
if context["dag_run"].partition_key: # รอบที่ asset ปลุก — เดือนมากับพาร์ทิชัน
return context["dag_run"].partition_key
start = context.get("data_interval_start") # รอบที่กดรันเองไม่มีคีย์นี้เลย — ["..."] จะยก KeyError
if start is None:
raise AirflowFailException("รอบนี้ไม่มีช่วงข้อมูล — บอกเดือนด้วย -c หรือย้อนรันด้วย backfill")
return f"{start:%Y-%m}"
@dag(
schedule=PartitionedAssetTimetable( # รอจนทั้งสาม asset ประกาศเดือนเดียวกัน
assets=Asset.ref(name="stock_file_a") & Asset.ref(name="stock_file_b") & Asset.ref(name="stock_file_c")
),
start_date=datetime(2026, 5, 1),
catchup=False,
params={"month": ""},
default_args={"retries": 2, "retry_delay": timedelta(seconds=20)}, # แล็บรอ 20 วินาที ของจริงหลายนาที
)
def stock_monthly():
@task
def arrived(**context):
month = month_of(context)
files = sorted((LAB / "inbox").glob(f"stock_{month}_*.xlsx"))
if not files: # ยังไม่มีไฟล์ของเดือนนี้สักสาขา — ไม่มีอะไรให้ทำ ไม่ใช่ความผิดของใคร
raise AirflowSkipException(f"ยังไม่มีไฟล์ของเดือน {month}")
return len(files)
@task
def receive(branch, **context):
month = month_of(context)
path = LAB / "inbox" / f"stock_{month}_{branch}.xlsx"
if not path.exists(): # ไฟล์มาช้า — ชั่วคราว ปล่อยให้ retries ลองใหม่
raise FileNotFoundError(f"ไฟล์สาขา {stock_gate.NAMES[branch]} เดือน {month} ยังมาไม่ถึง")
rows, missing = steps.read_branch(path)
staged = steps.stage(branch, rows, missing, LAB / "staging" / month / f"{branch}.json")
return {**staged, "problems": stock_gate.inspect(path, rows, missing, month)}
@task
def gate(staged):
for s in staged:
print(f"ด่าน ← receive_{s['branch']}: {s['rows']} แถว · ปัญหา {len(s['problems'])} ข้อ")
try:
return stock_gate.check(staged)
except stock_gate.BadFile as e: # ไฟล์ผิดรูป ลองกี่ครั้งก็ไม่หาย — หยุดเลย ไม่ใช้ retries ที่เหลือ
raise AirflowFailException(str(e)) from e
@task
def transform(staged, **context):
return steps.transform(staged, LAB / "staging" / month_of(context) / "all.json")
@task
def load(summary, **context):
return stock_load.load_month(summary, LAB / "warehouse.db", month_of(context)) # ← เดิม steps.load
@task
def notify(staged, summary, in_db):
print(steps.report(staged, summary, in_db))
@task(trigger_rule="all_done") # ส่งทุกเช้า ไม่ว่าท่อจะผ่านหรือพัง — แบบ notify ของบทที่ 7
def morning():
print("morning: ส่งข้อความเช้าแล้ว")
@task(trigger_rule="one_failed") # watcher: มีขั้นไหนพัง ขั้นนี้ได้รันแล้วพังเสมอ
def watcher():
raise AirflowFailException("มีขั้นที่พังในรอบนี้")
found = arrived()
received = [receive.override(task_id=f"receive_{b}")(b) for b in "abc"]
found >> received
staged = gate(received)
summary = transform(staged)
loaded = load(summary)
done = notify(staged, summary, loaded)
note = morning()
done >> note
[found, *received, staged, summary, loaded, done, note] >> watcher() # ต่อจากทุกขั้น
stock_monthly()สั่งเองด้วย -c ยังได้เหมือนเดิม
airflow dags test stock_monthly -c '{"month": "2026-07"}' 2>&1 | grep "ในคลัง"ผลรันจริง
สต็อกรวม 951 ชิ้น (ก 121 · ข 382 · ค 448) จาก 99 แถว · ในคลังเดือนนี้ 99 แถวฉากจริงใช้ scheduler (ไฟล์สาขา ก เดือนกันยายนเป็นฉบับแก้ของบทที่ 11)
บันทึกจากแล็บ
$ airflow dags trigger stock_file_a -c '{"month": "2026-09"}' ← ประกาศ: ไฟล์สาขา a เดือน 2026-09 มาแล้ว
$ airflow dags trigger stock_file_b -c '{"month": "2026-09"}' ← ประกาศ: ไฟล์สาขา b เดือน 2026-09 มาแล้ว
stock_monthly: (ยังไม่มีรอบ)
$ airflow dags trigger stock_file_c -c '{"month": "2026-09"}' ← ไฟล์ยังอยู่ใน late/
stock_file_c failed · FileNotFoundError: ยังไม่มี stock_2026-09_c.xlsx ในกล่องรับไฟล์
stock_monthly: (ยังไม่มีรอบ)
$ mv ~/airflow-lab/inbox/late/stock_2026-09_c.xlsx ~/airflow-lab/inbox/
$ airflow dags trigger stock_file_c -c '{"month": "2026-09"}' ← ประกาศ: ไฟล์สาขา c เดือน 2026-09 มาแล้ว
stock_monthly: asset_triggered · partition=2026-09 · success
notify: สต็อกรวม 1,000 ชิ้น (ก 121 · ข 320 · ค 559) จาก 99 แถว · ในคลังเดือนนี้ 99 แถวpartition_key โดยไม่ต้องเดาเวลาหรือความถี่สองประกาศแรกไม่ปลุกท่อ — เดือนกันยายนยังขาดสาขา ค · การสั่งสาขา ค ก่อนไฟล์มาพังที่ผู้ผลิตเอง ไม่มีประกาศหลุดไปถึงผู้ใช้ · พอไฟล์มาจริง ท่อถูกปลุกภายในไม่กี่วินาที และรู้เดือนจากพาร์ทิชันโดยไม่มีใครพิมพ์ -c · ไม่มีตรงไหนต้องเดาว่าไฟล์จะมาเมื่อไหร่
14.3 เลือกให้ตรงกับแหล่งข้อมูล
| แหล่งข้อมูล | ใช้ | เพราะ |
|---|---|---|
| มาตรงเวลาแน่นอน เช่นระบบอื่นส่งออกทุกตีหนึ่ง | นาฬิกา (บทที่ 13) | ง่ายสุด ไม่มีอะไรต้องรอ |
| มาไม่ตรงเวลา และไม่มีใครบอกได้ว่ามาแล้ว | sensor แบบ reschedule | ต้องถามเอง แต่ไม่ถือ worker ระหว่างรอ |
| มาไม่ตรงเวลา และคนส่งบอกได้ | Asset (แบ่งพาร์ทิชันถ้าต้องรอหลายชิ้นของงวดเดียวกัน) | ไม่ต้องเดาทั้งเวลาและความถี่ |
ทั้งสามแบบยังต้องการบทที่ 12 · ผู้ผลิตอาจประกาศซ้ำ คนอาจสั่งรอบเดือนเดิมซ้ำ — ท่อที่ถูกปลุกด้วยอะไรก็ตามยังต้องรันซ้ำได้
ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)
- ระบบบัญชีส่งไฟล์ยอดขายรายวันเข้าโฟลเดอร์ทุกตี 2 ตรงเวลาเสมอ · อีกระบบหนึ่งส่งไฟล์ราคาสินค้าเมื่อมีการเปลี่ยนราคา ซึ่งเดาไม่ได้ว่าเมื่อไหร่ แต่ระบบนั้นเรียก API ได้ · แต่ละแหล่งควรใช้นาฬิกา sensor หรือ Asset
- ถ้า
stock_files.pyไม่ยันว่าไฟล์มีจริงก่อนประกาศ แล้วมีคนสั่งสาขา ค เดือนกันยายนตอนไฟล์ยังอยู่ในlate/ตารางของบทที่ 4 ของstock_monthlyจะเกิดอะไรขึ้นทีละขั้น — ใช้บทที่ 9 กับ 10 ตอบ - สาขา ข ส่งไฟล์เดือนกันยายนฉบับแก้มา แล้วระบบรับไฟล์สั่ง
stock_file_bประกาศเดือน 2026-09 อีกครั้ง · ถ้าผู้ใช้ถูกปลุกอีกรอบ คลังจะเป็นยังไง และทำไมไม่ต้องกังวล