← คู่มือ 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

หน้าต่างเลือก 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 1receive_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 — ครึ่งแรกจะเป็นยังไง · ทั้งสองเรื่องคือบทถัดไป
ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)
- ถ้าเช้านั้นสาขา ค ต่างหากที่ส่งไฟล์แก้มา ต้อง clear ช่องไหน และหน้าต่าง Clear จะบอกว่ากี่ช่อง
- หลัง Clear
gateเป็นnone 1และได้เพดาน 3 · ถ้าไฟล์ใหม่ของสาขา ก ยังผิดอยู่gateจะจบที่ครั้งที่เท่าไหร่ เพราะอะไร watcherจบเป็นskipped 1· ทำไมไม่ใช่skipped 0แบบช่องที่ถูกข้ามในบทที่ 9