← คู่มือ Airflow — ทำไมงานหลายขั้นต้องมีคนคุมลำดับ
LEVEL 2 · ระดับกลาง
ด่านตรวจก่อนโหลด
ลองเดือนสิงหาคมตามข้อ 2 ของหัวข้อ 8.6 แล้ว transform พังด้วย ValueError: invalid literal for int() with base 10: '\xa0' — ไม่บอกว่าสาขาไหน แถวไหน · ส่วนไฟล์ศูนย์แถวแบบเดือนกรกฎาคมของบทที่ 1 ยังเดินผ่าน check_columns ไปถึงขั้นโหลดได้ · บทนี้สร้างด่านที่หยุดทั้งสองแบบก่อนขั้นโหลด พร้อมเหตุเป็นภาษาคน
ลองเดาก่อน: เดือนสิงหาคม สาขา ข มีหมวดใหม่ "อุปกรณ์ศิลปะ" สามรายการ — ด่านที่ดีควรหยุดไฟล์นี้ไหม · จดไว้ เฉลยอยู่หัวข้อ 9.2
9.1 ตรวจโครง ไม่ตรวจค่าที่สาขาเป็นเจ้าของ
ด่านตรวจว่าไฟล์มีรูปร่างที่ขั้นถัดไปต้องใช้ ไม่ได้ตรวจว่าข้อมูลหน้าตาเหมือนเดือนก่อน
ชื่อหมวด ชื่อสินค้า จำนวนรายการ เป็นของที่สาขาเปลี่ยนเองได้ทุกเดือน ด่านที่ล็อกค่าพวกนี้จะพังทุกครั้งที่ร้านขายของใหม่ · ด่านจึงตรวจแค่สี่เรื่องที่ขั้นถัดไปจะพังหรือเข้าใจผิดถ้าไม่มี
| ตรวจ | เพราะ | ไฟล์ที่โดน |
|---|---|---|
| คอลัมน์ที่ต้องใช้ครบ | transform อ่าน "จำนวน" | 2026-09_a |
| จำนวนแถว 10–200 | ศูนย์แถวคือไฟล์ผิด ไม่ใช่สต็อกหมด · ช่วงกว้างโดยตั้งใจ | บล็อกทดลอง |
| ช่องจำนวนไม่ว่าง | int() ของ transform | 2026-08_c |
| วันที่นับอยู่ในเดือนที่สั่ง | รู้ว่าไฟล์เป็นของเดือนไหน (บทที่ 13) | 2026-08_a |
แถวซ่อนไม่อยู่ในตาราง — นับรวมแล้วเตือนในรายงานตามที่บทที่ 8 ตัดสินไว้ · ช่อง "ว่าง" เช็กด้วย .strip() เปล่า ๆ เพราะช่องที่ใส่ non-breaking space ของสาขา ค ผ่าน == "" กับ .strip(" ") ได้ (บรรทัดแรกของผลรันข้างล่าง)
inspect ตรวจทีละสาขาแล้วคืนรายการปัญหา · check รวมปัญหาทุกสาขาแล้วหยุดครั้งเดียว — หยุดที่ปัญหาแรก = แก้ทีละเรื่อง รันซ้ำทีละรอบ · BadFile คือชนิดของ error ที่แปลว่า "ต้องให้สาขาส่งใหม่" บทที่ 10 จะใช้มันแยกความพังถาวร
# ~/airflow-lab/dags/stock_gate.py — ด่านตรวจโครงก่อนโหลด: ฟังก์ชันธรรมดา ไม่ import อะไรจาก Airflow
from datetime import datetime
from openpyxl import load_workbook
NAMES = {"a": "ก", "b": "ข", "c": "ค"}
ROWS = range(10, 201) # กว้างโดยตั้งใจ: จับไฟล์ผิดตัวหรือถูกตัดหาย ไม่ได้ล็อกขนาดแคตตาล็อกที่สาขาเป็นเจ้าของ
class BadFile(ValueError):
"""ไฟล์ผิดรูป — ต้องให้สาขาส่งใหม่ถึงจะหาย"""
def count_date(path):
"""ค่าในช่องข้าง ๆ คำว่า "วันที่นับ" — หาด้วยคำ ไม่เชื่อตำแหน่ง แบบเดียวกับหัวตาราง"""
for ws in load_workbook(path).worksheets:
for row in ws.iter_rows(max_row=10):
for cell in row:
if cell.value == "วันที่นับ":
return ws.cell(cell.row, cell.column + 1).value
return None
def inspect(path, rows, missing, month):
"""ตรวจโครงของไฟล์หนึ่งสาขา → รายการปัญหา (ว่าง = ผ่าน) · ไม่ตรวจค่าที่สาขาเป็นเจ้าของ เช่นชื่อหมวด"""
problems = []
if missing:
problems.append(f"ไม่มีคอลัมน์ {', '.join(missing)}")
if len(rows) not in ROWS:
problems.append(f"มี {len(rows)} แถว นอกช่วง {ROWS.start}–{ROWS.stop - 1}")
if "จำนวน" not in missing:
# .strip() เปล่า ๆ ตัด non-breaking space ได้ · == "" กับ .strip(" ") ตัดไม่ได้
blank = [row["รหัส"] for row in rows if row["จำนวน"] is None or str(row["จำนวน"]).strip() == ""]
if blank:
problems.append(f"ช่องจำนวนว่าง {len(blank)} แถว (รหัส {', '.join(blank)})")
when = count_date(path)
if not isinstance(when, datetime):
problems.append(f"วันที่นับอ่านไม่ได้: {when!r}")
elif f"{when:%Y-%m}" != month:
problems.append(f"วันที่นับ {when:%Y-%m-%d} ไม่ใช่เดือน {month}")
return problems
def check(staged):
"""รวมปัญหาของทุกสาขาแล้วหยุดครั้งเดียว — หยุดที่ปัญหาแรก = แก้ทีละเรื่อง รันซ้ำทีละรอบ"""
bad = [f"สาขา {NAMES[s['branch']]}: {p}" for s in staged for p in s["problems"]]
if bad:
raise BadFile("ด่านหยุดท่อ — " + " · ".join(bad))
return staged
if __name__ == "__main__": # ทดลองด่านโดยไม่มี Airflow — สามสาขา สาขาละอาการ (แถวแบบที่ read_branch คืน)
import tempfile
from pathlib import Path
from openpyxl import Workbook
nbsp = "\u00a0"
print(f'== "" → {nbsp == ""} · .strip(" ") → {nbsp.strip(" ") == ""} · .strip() → {nbsp.strip() == ""}')
work = Path(tempfile.mkdtemp())
cases = { # สาขา: (วันที่นับ, แถว, คอลัมน์ที่ขาด)
"a": (datetime(2026, 8, 31), [], []), # มีแต่หัวตาราง
"b": (datetime(2026, 8, 31), [{"รหัส": f"S{i:03d}", "หมวด": "สีไม้" if i > 10 else "ปากกา",
"จำนวน": nbsp if i in (3, 7) else 5} for i in range(1, 13)], []),
"c": ("สิ้นเดือน", [{"รหัส": f"S{i:03d}", "หมวด": "ปากกา"} for i in range(1, 13)], ["จำนวน"]),
}
staged = []
for b, (date, rows, missing) in cases.items():
wb = Workbook()
wb.active.append(["วันที่นับ", date])
wb.save(work / f"{b}.xlsx")
staged.append({"branch": b, "problems": inspect(work / f"{b}.xlsx", rows, missing, "2026-08")})
try:
check(staged)
except BadFile as e:
print(type(e).__name__, "·", str(e).replace(" · สาขา", "\n สาขา"))ผลรันจริง
== "" → False · .strip(" ") → False · .strip() → True
BadFile · ด่านหยุดท่อ — สาขา ก: มี 0 แถว นอกช่วง 10–200
สาขา ข: ช่องจำนวนว่าง 2 แถว (รหัส S003, S007)
สาขา ค: ไม่มีคอลัมน์ จำนวน
สาขา ค: วันที่นับอ่านไม่ได้: 'สิ้นเดือน'สาขา ข มีหมวดใหม่ "สีไม้" ด่านไม่ทัก แต่จับช่องว่างแปลกได้พร้อมรหัสสินค้า · สาขา ค มีสองปัญหา ได้ครบทั้งสองในรอบเดียว
9.2 ต่อด่านเข้าท่อ แล้วดูสถานะทุกช่อง
stock_monthly.py ฉบับนี้เขียนทับฉบับของบทที่ 8 · ต่างกันสามจุด: ขั้นใหม่ arrived (หัวข้อ 9.3) · receive แนบผลของ inspect ไปกับใบสรุป · gate เข้าแทน merge_check
# ~/airflow-lab/dags/stock_monthly.py — บทที่ 9: ขั้น arrived (ยังไม่มีไฟล์ = ข้าม) + ด่านตรวจโครงแทน merge_check
from datetime import datetime
from pathlib import Path
import stock_gate
import stock_steps as steps
from airflow.sdk import dag, task
from airflow.sdk.exceptions import AirflowSkipException
LAB = Path.home() / "airflow-lab"
@dag(schedule=None, start_date=datetime(2026, 5, 1), catchup=False, params={"month": "2026-07"})
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"
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'])} ข้อ")
return stock_gate.check(staged)
@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))
received = [receive.override(task_id=f"receive_{b}")(b) for b in "abc"]
arrived() >> received
staged = gate(received)
summary = transform(staged)
notify(staged, summary, load(summary))
stock_monthly()ตั้งแต่บทนี้หลายรอบจะพังโดยตั้งใจ · สคริปต์ข้างล่างพิมพ์คอลัมน์ของรอบล่าสุดเป็นตารางของบทที่ 4 (สถานะกับครั้งที่ลอง เรียงตามจังหวะแบบบทที่ 3) · มันอ่านฐานข้อมูล SQLite ของแล็บตรง ๆ ใช้ได้แค่เครื่องทดลอง ของจริงดูที่ Grid · ตารางที่มันอ่านคือเรื่องของบทที่ 15
# ~/airflow-lab/states.py — พิมพ์คอลัมน์ของรอบล่าสุดจากฐานข้อมูลของแล็บ ในรูปตารางของบทที่ 4
import importlib
import sqlite3
import sys
from pathlib import Path
import airflow # noqa: F401 — ต้องรันด้วย python ของแล็บ เพราะต้อง import ไฟล์ DAG เพื่ออ่านกราฟ
LAB = Path.home() / "airflow-lab"
dag_id = sys.argv[1]
sys.path.insert(0, str(LAB / "dags"))
dag = getattr(importlib.import_module(dag_id), dag_id)() # ประกาศกราฟอย่างเดียว ไม่รันงาน (บทที่ 6)
ups = {t.task_id: t.upstream_task_ids for t in dag.tasks}
def beat(task): # จังหวะแบบบทที่ 3: ขั้นที่ไม่รอใคร = 1 · ขั้นอื่น = 1 + จังหวะที่มากที่สุดของต้นน้ำ
return 1 + max((beat(u) for u in ups[task]), default=0)
db = sqlite3.connect(f"file:{LAB / 'airflow.db'}?mode=ro", uri=True) # อ่านอย่างเดียว ไม่แก้ตาราง
run_id, run_state, conf = db.execute(
"SELECT run_id, state, conf FROM dag_run WHERE dag_id = ? ORDER BY id DESC LIMIT 1", [dag_id]
).fetchone()
cells = {
task: (state, tries)
for task, state, tries in db.execute(
"SELECT task_id, state, try_number FROM task_instance WHERE dag_id = ? AND run_id = ?", [dag_id, run_id]
)
}
print(f"รอบล่าสุดของ {dag_id} {conf} → {run_state}")
for task in sorted(ups, key=lambda t: (beat(t), t)):
state, tries = cells.get(task, (None, 0))
print(f" {task:<14}{state or 'none'} {tries}") # ช่องว่างในฐานข้อมูล = none ของบทที่ 4วาง stock_gate.py กับ stock_monthly.py ฉบับใหม่ไว้ใน ~/airflow-lab/dags/ แล้วรันเดือนสิงหาคม
airflow dags test stock_monthly -c '{"month": "2026-08"}' 2>&1 | grep -E "ด่าน ←|BadFile:"ผลรันจริง
ด่าน ← receive_a: 33 แถว · ปัญหา 1 ข้อ
ด่าน ← receive_b: 36 แถว · ปัญหา 0 ข้อ
ด่าน ← receive_c: 33 แถว · ปัญหา 1 ข้อ
stock_gate.BadFile: ด่านหยุดท่อ — สาขา ก: วันที่นับอ่านไม่ได้: 'สิ้นเดือน ส.ค.' · สาขา ค: ช่องจำนวนว่าง 3 แถว (รหัส S007, S008, S023)python ~/airflow-lab/states.py stock_monthlyผลรันจริง
รอบล่าสุดของ stock_monthly {"month": "2026-08"} → failed
arrived success 1
receive_a success 1
receive_b success 1
receive_c success 1
gate failed 1
transform upstream_failed 0
load upstream_failed 0
notify upstream_failed 0สาขา ข 36 แถวพร้อมหมวดใหม่ผ่านด่าน (เฉลยคำถามเปิดบท: ไม่ควรหยุด หมวดเป็นของสาขา) · '\xa0' ใน transform ของบทที่ 8 กลายเป็นเหตุที่ลงมือได้ทันที: สาขา ค รหัส S007 S008 S023 · transform ถึง notify เป็น upstream_failed 0 ไม่มีอะไรลงคลัง
9.3 ไม่มีไฟล์ ไม่เท่ากับไฟล์ผิด
arrived นับไฟล์ของเดือนที่สั่ง ถ้าไม่มีสักสาขามันยก AirflowSkipException แทน error · ลองเดือนตุลาคมที่ยังไม่มีใครส่ง
airflow dags test stock_monthly -c '{"month": "2026-10"}' 2>&1 | grep -o "reason=.*"ผลรันจริง
reason='ยังไม่มีไฟล์ของเดือน 2026-10'python ~/airflow-lab/states.py stock_monthlyผลรันจริง
รอบล่าสุดของ stock_monthly {"month": "2026-10"} → success
arrived skipped 1
receive_a skipped 0
receive_b skipped 0
receive_c skipped 0
gate skipped 0
transform skipped 0
load skipped 0
notify skipped 0states.py ในหัวข้อ 9.2–9.3 · ไม่มีอะไรให้ทำ — การข้ามไหลลงปลายน้ำเป็น skipped ขั้นปลายสุดที่ข้ามนับเป็นผ่าน รอบจึงเขียว · ไฟล์ผิด — ด่านพัง ปลายน้ำเป็น upstream_failed รอบแดง · ทั้งสองแบบไม่มีอะไรถูกโหลด ต่างกันที่ว่าต้องมีคนลงมือไหม💡 skipped — "ไม่ได้ทำ และไม่มีอะไรผิด" · หน้า Tasks ทางการยกตัวอย่างไว้ตรงตัว: ข้ามเมื่อรู้ว่าไม่มีข้อมูล
การข้ามไหลลงปลายน้ำผ่าน trigger rule ตั้งต้น all_success เหมือนความพัง (หน้า Dags ทางการ) · ต่างที่ปลายทาง: ขั้นปลายสุด skipped นับเป็นผ่าน รอบจึงเขียว ไม่มีใครถูกปลุก
เลือกด้วยคำถามเดียว: ถ้าเรื่องนี้เกิด ต้องมีคนลงมือไหม — ต้อง = failed (ไฟล์ผิดรูป ต้องขอไฟล์ใหม่) · ไม่ต้อง เพราะรอบหน้าจะเก็บเอง = skipped (ยังไม่มีใครส่ง — ตราบที่มีรอบถัดไปมาเก็บ บทที่ 13–14)
⚠️ arrived ข้ามเฉพาะเมื่อไม่มีไฟล์สักสาขา · ขาดสาขาเดียว receive ของสาขานั้นพัง — ข้ามไม่ได้ เพราะรอบเขียวที่ขาดหนึ่งสาขาคือยอดผิดแบบเงียบ ๆ ของบทที่ 1 · ความพังแบบนี้รอแล้วอาจหาย — บทถัดไป
ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)
- ท่อรับไฟล์ CSV ยอดขายรายวันจากหลายร้าน · วันอาทิตย์ร้านปิด ไม่มีไฟล์ · วันหนึ่งไฟล์มาแต่ไม่มีคอลัมน์ยอดขาย · อีกวันมีร้านใหม่โผล่มา — แต่ละกรณีควร
skippedfailedหรือปล่อยผ่าน - เพื่อนแก้
gateให้ยกAirflowSkipExceptionแทน error เมื่อไฟล์ผิด "จะได้ไม่มีรอบแดงรกตา" · รอบเดือนสิงหาคมจะเป็นสีอะไร และเช้านั้นใครจะรู้ว่าสาขา ค มีช่องว่าง - ไฟล์เดือนกรกฎาคมของบทที่ 1 (ซ่อนทั้ง 33 แถว) ผ่านด่านนี้ไหม · ถ้ายังใช้ตัวอ่านเก่าที่ข้ามแถวซ่อน ด่านจะทำอะไร