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

← คู่มือ Airflow — ทำไมงานหลายขั้นต้องมีคนคุมลำดับ

LEVEL 1 · พื้นฐานที่ใช้ทุกวัน

อ่าน Grid ให้เป็น

เช้าวันหนึ่งคุณเปิดหน้าจอของ Airflow แล้วเห็นแถบของทุกรอบเป็นสีเขียว — เมื่อคืนทุกอย่างผ่านไหม · บทนี้สร้างสามรอบที่มีเรื่องต่างกัน แล้วอ่านหน้าจอเทียบกับตารางของบทที่ 4 ทีละช่อง

7.1 Grid คือตารางของบทที่ 4

หน้า UI ในเอกสารทางการบรรยาย Grid ไว้ว่า แต่ละแถวคือขั้น แต่ละคอลัมน์คือรอบ — ตารางของบทที่ 4 ตรงตัว

DAG ข้างล่างคือ stock_first เติมสามอย่าง: สวิตช์จำลองความพังที่ส่งมากับรอบ (params) · ทุกขั้นได้สิทธิ์ลองใหม่ 1 ครั้ง รอ 10 วินาที (default_args) · notify ตั้ง trigger_rule="all_done" ให้ส่งรายงานทุกเช้าแม้ขั้นก่อนหน้าพัง · notify ไม่รับค่าจากใคร เส้นของมันจึงวาดด้วย >> ซึ่งแปลว่า "ซ้ายต้องจบก่อนขวา" โดยไม่ส่งค่า

# ~/airflow-lab/dags/stock_grid.py — DAG ของบทที่ 6 + สวิตช์จำลองความพัง ไว้ดู Grid
from datetime import datetime, timedelta

from airflow.sdk import dag, task

COUNTS = {"a": 120, "b": 340, "c": 560}


@dag(
    schedule=None,
    start_date=datetime(2026, 5, 1),
    catchup=False,
    params={"fail": "none"},  # none · receive_c (พังเฉพาะครั้งแรก) · transform (พังทุกครั้ง)
    default_args={"retries": 1, "retry_delay": timedelta(seconds=10)},
)
def stock_grid():
    @task
    def receive(branch, **context):
        if f"receive_{branch}" == context["params"]["fail"] and context["ti"].try_number == 1:
            raise FileNotFoundError(f"ไฟล์สาขา {branch} ยังมาไม่ถึง")
        return COUNTS[branch]

    @task
    def merge_check(counts):
        return counts

    @task
    def transform(counts, **context):
        if context["params"]["fail"] == "transform":
            raise ValueError("ไม่มีคอลัมน์ จำนวน")
        return sum(counts)

    @task
    def load(total):
        print(f"load: รับ {total} ชิ้น")

    @task(trigger_rule="all_done")  # ← ส่งรายงานทุกเช้า แม้ขั้นก่อนหน้าจะพัง
    def notify():
        print("notify: ส่งรายงานแล้ว")

    counts = [receive.override(task_id=f"receive_{b}")(b) for b in COUNTS]
    load(transform(merge_check(counts))) >> notify()


stock_grid()

สร้างสามรอบ (standalone ต้องเปิดอยู่ ครั้งนี้ให้ scheduler เป็นคนเดินรอบ)

airflow dags reserialize          # ไม่รอรอบสแกน 5 นาที (หัวข้อ 6.5)
airflow dags unpause stock_grid   # DAG ใหม่เริ่มในสถานะหยุด
airflow dags trigger stock_grid -c '{"fail": "none"}'
airflow dags trigger stock_grid -c '{"fail": "receive_c"}'
airflow dags trigger stock_grid -c '{"fail": "transform"}'

แล็บเว้นแต่ละคำสั่ง trigger ราว 50 วินาที ให้รอบก่อนหน้าจบก่อน · -c ส่งค่าที่ทับ params ของรอบนั้น

7.2 อ่านภาพแรก

หน้า Grid ของ stock_grid หลังสามรอบ: แท่งรอบเขียวทั้งสาม รอบที่สาม transform แดง load ส้ม และการ์ด 1 Failed Task กับ 0 Failed Runs

  • สามคอลัมน์ = สามรอบ เก่าไปใหม่จากซ้ายไปขวา · แท่งบนหัวคอลัมน์คือตัวรอบ สูงตามเวลาที่ใช้ สีตามสถานะของรอบ — เขียวทั้งสามแท่ง
  • เจ็ดแถว = เจ็ดขั้น · ช่องเขียวมีเครื่องหมายถูก = success
  • คอลัมน์ที่สาม (รอบ fail = transform): transform กากบาทแดง = failed · load ไอคอนส้ม = upstream_failed · notify เขียว
  • ซีกขวา: รอบล่าสุด (15:27:04 คือรอบที่สาม) มีเครื่องหมายถูกเขียว · การ์ดสองใบเขียนว่า 1 Failed Task กับ 0 Failed Runs

"พังที่ไหน รอบไหน" ตอบได้จากตำแหน่งของช่องเดียว: แถว transform คอลัมน์ที่สาม · แต่การ์ดสองใบขัดกันเอง — ขั้นพังหนึ่งขั้น รอบพังศูนย์รอบ

ดูความสูงของแท่งแรกด้วย: รอบที่ไม่มีอะไรพังใช้ 10.7 วินาที ทั้งที่ตั้งแต่ merge_check ลงไปแต่ละขั้นใช้ราว 0.1 วินาที — บันทึกเวลาในฐานข้อมูลของแล็บบอกว่าขั้นถัดไปเริ่มหลังขั้นก่อนหน้าจบราว 1 วินาทีทุกครั้ง · ช่วงว่างนั้นมาจากไหน ในเมื่อแบบจำลองของบทที่ 4 เดินต่อทันทีในจังหวะถัดไป คือปริศนาของบทที่ 15

7.3 รอบเขียวไม่ได้แปลว่าทุกขั้นผ่าน

นี่คือคำตอบของปริศนาข้อ 2 ในหัวข้อ 4.7 · สถานะของรอบดูจากขั้นปลายสุดเท่านั้น ซึ่งท่อนี้คือ notify · all_done แปลว่า "รันเมื่อต้นน้ำจบ ไม่ว่าจบแบบไหน" — load จบแบบ upstream_failed ก็นับว่าจบ notify จึงรันและผ่าน ⇒ รอบ success · หน้า DAG Runs ในเอกสารทางการเตือนกรณีนี้ตรงตัว

⚠️ แถบเขียวตอบได้แค่ว่าขั้นปลายสุดผ่านไหม ไม่ได้ตอบว่าทุกขั้นผ่านไหม · ท่อที่มีขั้นรายงานแบบ all_done ต้องอ่านทีละช่อง หรือใช้การ์ด Failed Task เป็นสัญญาณแรก · ทางแก้ที่ทำให้รอบแดงจริงเรียกว่า watcher pattern — บทที่ 11

7.4 ครั้งที่เท่าไหร่ และเพราะอะไร

กดที่ช่อง transform ของคอลัมน์ที่สาม

หน้าช่อง transform ของรอบที่สาม: ป้าย Failed · Try Number 2 · Task Tries ปุ่ม 1 และ 2 กากบาทแดงทั้งคู่ · บันทึก ValueError ไม่มีคอลัมน์ จำนวน · แถบบนสุดของรอบมีเครื่องหมายถูกเขียว

  • แถบบนสุด DAG → รอบ → ขั้น · ตัวรอบมีเครื่องหมายถูกเขียว — ข้อ 7.3 อีกครั้ง
  • Failed · Try Number 2 — retries=1 จึงลองได้สองครั้ง แบบ transform ของเดือนสิงหาคมในหัวข้อ 4.3
  • Task Tries 1 · 2 — ปุ่มละครั้ง เปิดบันทึกแยกกัน ทั้งสองครั้งกากบาทแดง
  • บันทึก: ValueError: ไม่มีคอลัมน์ จำนวน — ตอบว่าเพราะอะไร

ทีนี้กดช่อง receive_c ของคอลัมน์ที่สอง ซึ่งเขียวเหมือนช่องรอบข้าง · หน้าจอบอก Success · Try Number 2 · Task Tries ปุ่ม 1 กากบาทแดง ปุ่ม 2 ถูกเขียว · บันทึกครั้งที่ 2 มีบรรทัด Done. Returned value was: 560 — ค่าที่ถูกเขียนลง XCom · สีของช่องบอกสถานะสุดท้าย ไม่บอกประวัติ

7.5 ลำดับการอ่าน Grid ตอนเช้า

  1. หาคอลัมน์ของคืนนั้น — ขวาสุดคือรอบใหม่สุด
  2. ไล่หาช่องที่ไม่ใช่ถูกเขียวทีละแถว — อย่าเชื่อสีของแท่งรอบ
  3. หาต้นเหตุ — ช่องแดงที่อยู่ต้นน้ำที่สุดคือจุดที่พัง ช่องส้มใต้มันคือผลที่ตามมา
  4. เปิดช่องต้นเหตุ — อ่าน Try Number แล้วหาบรรทัด ERROR ในบันทึกของครั้งสุดท้าย
  5. เผื่อใจกับช่องเขียว — Try Number มากกว่า 1 คือความพังชั่วคราวที่ retry กลบไว้ ถ้าเกิดทุกคืนก็ควรรู้

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

  1. เพื่อนส่ง Grid มาให้ดู คอลัมน์ขวาสุด: receive ทั้งสามเขียว · merge_check แดง · transform load notify ส้มทั้งสามช่อง (ท่อของเพื่อนไม่ได้ตั้ง all_done) · แท่งของรอบจะเป็นสีอะไร พังที่ไหน และช่องไหนที่ควรกดเปิดก่อน
  2. ในรอบที่สามของ stock_grid ทำไม load เป็นส้ม แต่ notify เป็นเขียว ทั้งที่ทั้งคู่อยู่ปลายน้ำของ transform
  3. ถ้าเปลี่ยน retries ของ stock_grid เป็น 3 แล้วสั่งรอบ fail = transform ใหม่ ช่อง transform จะมีปุ่มใน Task Tries กี่ปุ่ม และรอบนั้นจะใช้เวลานานขึ้นอย่างน้อยกี่วินาที

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