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

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

LEVEL 2 · ระดับกลาง

พังกลางทาง: ล้างสถานะเฉพาะขั้นแล้วรันต่อ

เช้าวันที่ด่านหยุดเดือนกันยายน สาขา ก ส่งไฟล์แก้มา · สาขา ข กับ ค อ่านผ่านไปแล้ว — ไม่ต้องรันทั้งรอบ: ล้างช่องที่ต้องทำใหม่กับปลายน้ำของมัน แล้วให้วงรอบเดินต่อ (หัวข้อ 4.5)

11.1 รอบแดงจริง: watcher

ฉบับนี้เติมสองขั้น · morning ส่งข้อความทุกเช้าแม้ท่อพัง จึงใช้ all_done — ซึ่งทำให้รอบเขียวทั้งที่ข้างในพัง (บทที่ 7) · watcher คือทางแก้ที่หน้า Best Practices ทางการเรียกว่า watcher pattern: ต่อจากทุกขั้น ตั้ง one_failed และยก error เสมอ · มีขั้นพัง มันรันแล้วพัง รอบจึงแดง · ไม่มี มันถูกข้าม

# ~/airflow-lab/dags/stock_monthly.py — บทที่ 11: เพิ่มข้อความเช้า (all_done) กับ watcher ที่ทำให้รอบแดงจริง
from datetime import datetime, timedelta
from pathlib import Path

import stock_gate
import stock_steps as steps
from airflow.sdk import dag, task
from airflow.sdk.exceptions import AirflowFailException, AirflowSkipException

LAB = Path.home() / "airflow-lab"


@dag(
    schedule=None,
    start_date=datetime(2026, 5, 1),
    catchup=False,
    params={"month": "2026-07"},
    default_args={"retries": 2, "retry_delay": timedelta(seconds=20)},  # แล็บรอ 20 วินาที ของจริงหลายนาที
)
def stock_monthly():
    @task
    def arrived(**context):
        month = context["params"]["month"]
        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 = context["params"]["month"]
        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" / context["params"]["month"] / "all.json")

    @task
    def load(summary, **context):
        return steps.load(summary, LAB / "warehouse.db", context["params"]["month"])

    @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()

11.2 Clear ช่องที่อ่านของที่เปลี่ยน

ฉากนี้ใช้ scheduler จริงกับปุ่มบน Grid จึงเป็นบันทึกจากแล็บ · คำสั่งข้างล่างจำลองไฟล์แก้ของสาขา ก: สร้างไฟล์ปลอมใหม่โดยถอดอาการคอลัมน์หาย (seed เดิม ไฟล์อื่นเท่าเดิม)

cd ~/airflow-lab && python -c "
import make_fake_files as m
from pathlib import Path
del m.QUIRKS[('2026-09', 'a')]
m.main(Path('fixed'))" && cp fixed/stock_2026-09_a.xlsx inbox/

clear ช่องไหน · gate พัง แต่มันตรวจใบสรุปของ receive_a จากเมื่อคืน ซึ่งยังบอกว่าไม่มีคอลัมน์ — clear gate อย่างเดียวพังซ้ำด้วยเหตุเดิม · ต้องเป็น ช่องที่อ่านของที่เปลี่ยน: กดช่อง receive_a ใน Grid แล้วกด Clear

หน้าต่าง Clear ของช่อง receive_a เลือก Downstream · Affected Tasks 7 แถว: receive_a Success · gate Failed · transform load notify Upstream Failed · morning Success · watcher Failed

หน้าต่างเลือก Downstream ไว้ให้ และบอกก่อนกดว่าจะล้าง 7 ช่อง รวม morning ที่ผ่านไปแล้ว · arrived receive_b receive_c ไม่ถูกแตะ

บันทึกจากแล็บ — ผลของ states.py สามครั้งวางเคียงกัน

ก่อน Clear                      ทันทีหลัง Confirm               จบรอบ
รอบ → failed                    รอบ → running                    รอบ → success
  arrived       success 1         arrived       success 1          arrived       success 1
  receive_a     success 1         receive_a     running 2          receive_a     success 2
  receive_b     success 1         receive_b     success 1          receive_b     success 1
  receive_c     success 1         receive_c     success 1          receive_c     success 1
  gate          failed 1          gate          none 1             gate          success 2
  transform     upstream_failed 0 transform     none 0             transform     success 1
  load          upstream_failed 0 load          none 0             load          success 1
  notify        upstream_failed 0 notify        none 0             notify        success 1
  morning       success 1         morning       none 1             morning       success 2
  watcher       failed 1          watcher       none 1             watcher       skipped 1
ก่อน Clearsuccess 1success 1success 1success 1failed 1upstream_failed 0upstream_failed 0upstream_failed 0success 1failed 1failedทันทีหลัง Confirmsuccess 1running 2success 1success 1none 1none 0none 0none 0none 1none 1runningจบรอบsuccess 1success 2success 1success 1success 2success 1success 1success 1success 2skipped 1successarrivedreceive_areceive_breceive_cgatetransformloadnotifymorningwatcherสถานะของรอบClear ล้าง 7 ช่อง: receive_a · gate · transform · load · notify · morning · watcher
FIG 11.1 — Clear ช่อง receive_a พร้อม Downstream คัดจากบันทึกในหัวข้อ 11.2 · ช่องที่ถูกล้างกลับเป็น noneแต่ตัวนับครั้งที่ลองไม่หาย (gate none 1 แล้วจบที่ครั้งที่ 2) · ช่องจาง — arrived receive_b receive_c — ไม่ถูกแตะตลอดฉาก · morning ได้รันครั้งที่ 2 = ข้อความเช้าฉบับที่สอง (บทที่ 12)
  • ก่อน Clear รอบแดงเพราะ watcher — ไม่มีมัน ขั้นปลายสุดที่เหลือคือ morning ที่ผ่าน รอบจะเขียวแบบบทที่ 7
  • ล้างแล้วตัวนับไม่หาย — gate เป็น none 1 ได้เพดานใหม่ 3 = ลองไปแล้ว 1 + retries 2 · ขั้นที่พังครบสามครั้งแล้ว clear จะรันเป็นครั้งที่ 4 (ยันในแล็บ)
  • จบรอบ รายงาน 1,000 ชิ้น 99 แถว · watcher ถูกข้าม รอบเขียวจริง · ประวัติครั้งเก่ายังอยู่ใน Task Tries

คำสั่งเทียบเท่า airflow tasks clear stock_monthly -t receive_a -d มีกับดักสองข้อ · -t หาคำย่อย ไม่ใช่ regex แบบที่ --help บอก — -t receive โดนทั้งสามสาขา · และไม่ใส่ช่วงวันที่ = ทุกรอบของ DAG (แล็บ: 6 รอบ ถามยืนยัน 24 ช่อง · -y ข้ามคำถามนั้น) · ปุ่มบน Grid ล้างรอบเดียวเสมอ

บันทึกนี้ซ่อนสองเรื่อง · morning รันสองครั้ง ข้อความเช้าถูกส่งสองฉบับ · และถ้าเมื่อคืนพังที่ load ตอนเขียนไปครึ่งหนึ่ง แล้วคุณ clear load — ครึ่งแรกจะเป็นยังไง · ทั้งสองเรื่องคือบทถัดไป

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

  1. ถ้าเช้านั้นสาขา ค ต่างหากที่ส่งไฟล์แก้มา ต้อง clear ช่องไหน และหน้าต่าง Clear จะบอกว่ากี่ช่อง
  2. หลัง Clear gate เป็น none 1 และได้เพดาน 3 · ถ้าไฟล์ใหม่ของสาขา ก ยังผิดอยู่ gate จะจบที่ครั้งที่เท่าไหร่ เพราะอะไร
  3. watcher จบเป็น skipped 1 · ทำไมไม่ใช่ skipped 0 แบบช่องที่ถูกข้ามในบทที่ 9

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