ข้ามไปยังเนื้อหา
Tayakorn

← คู่มือ 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 สาขา" ×1

up_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 แถว
ไฟล์สาขา ค มาถึง inboxนาฬิกา (บทที่ 13)sensor แบบ rescheduleAsset แบ่งพาร์ทิชันปลุกตามเวลาreceive_cfailed 3failedไม่มีใครปลุกท่ออีก — ต้อง clear เอง (บทที่ 11)2/32/32/32/3ถามทุก 10 วินาที · ยังไม่ครบ = ปล่อยช่องคืน (up_for_reschedule)3/3success 1all_filesabcประกาศ a · b ยังไม่ปลุกท่อ · สั่ง c ก่อนไฟล์มา = ผู้ผลิตพังเองcsuccess 1stock_monthly · partition=2026-09
FIG 14.1 — เดือนกันยายนเดียวกันสามแบบ คัดจากบันทึกในหัวข้อ 13.3 · 14.1 · 14.2 · นาฬิกาปลุกก่อนไฟล์มาแล้วพัง หลังไฟล์มาไม่มีใครปลุกอีก · sensor ถามซ้ำจนครบ (ช้ากว่าไฟล์มาไม่เกินรอบถามหนึ่งรอบ) · Asset ถูกปลุกโดยคนส่งไฟล์ และรู้เดือนจาก partition_key โดยไม่ต้องเดาเวลาหรือความถี่

สองประกาศแรกไม่ปลุกท่อ — เดือนกันยายนยังขาดสาขา ค · การสั่งสาขา ค ก่อนไฟล์มาพังที่ผู้ผลิตเอง ไม่มีประกาศหลุดไปถึงผู้ใช้ · พอไฟล์มาจริง ท่อถูกปลุกภายในไม่กี่วินาที และรู้เดือนจากพาร์ทิชันโดยไม่มีใครพิมพ์ -c · ไม่มีตรงไหนต้องเดาว่าไฟล์จะมาเมื่อไหร่

14.3 เลือกให้ตรงกับแหล่งข้อมูล

แหล่งข้อมูลใช้เพราะ
มาตรงเวลาแน่นอน เช่นระบบอื่นส่งออกทุกตีหนึ่งนาฬิกา (บทที่ 13)ง่ายสุด ไม่มีอะไรต้องรอ
มาไม่ตรงเวลา และไม่มีใครบอกได้ว่ามาแล้วsensor แบบ rescheduleต้องถามเอง แต่ไม่ถือ worker ระหว่างรอ
มาไม่ตรงเวลา และคนส่งบอกได้Asset (แบ่งพาร์ทิชันถ้าต้องรอหลายชิ้นของงวดเดียวกัน)ไม่ต้องเดาทั้งเวลาและความถี่

ทั้งสามแบบยังต้องการบทที่ 12 · ผู้ผลิตอาจประกาศซ้ำ คนอาจสั่งรอบเดือนเดิมซ้ำ — ท่อที่ถูกปลุกด้วยอะไรก็ตามยังต้องรันซ้ำได้

ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)

  1. ระบบบัญชีส่งไฟล์ยอดขายรายวันเข้าโฟลเดอร์ทุกตี 2 ตรงเวลาเสมอ · อีกระบบหนึ่งส่งไฟล์ราคาสินค้าเมื่อมีการเปลี่ยนราคา ซึ่งเดาไม่ได้ว่าเมื่อไหร่ แต่ระบบนั้นเรียก API ได้ · แต่ละแหล่งควรใช้นาฬิกา sensor หรือ Asset
  2. ถ้า stock_files.py ไม่ยันว่าไฟล์มีจริงก่อนประกาศ แล้วมีคนสั่งสาขา ค เดือนกันยายนตอนไฟล์ยังอยู่ใน late/ ตารางของบทที่ 4 ของ stock_monthly จะเกิดอะไรขึ้นทีละขั้น — ใช้บทที่ 9 กับ 10 ตอบ
  3. สาขา ข ส่งไฟล์เดือนกันยายนฉบับแก้มา แล้วระบบรับไฟล์สั่ง stock_file_b ประกาศเดือน 2026-09 อีกครั้ง · ถ้าผู้ใช้ถูกปลุกอีกรอบ คลังจะเป็นยังไง และทำไมไม่ต้องกังวล

อ่านแบบเต็มเล่ม