Skip to content
Tayakorn

← คู่มือ 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() ของ transform2026-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 0
เดือนตุลาคม — ยังไม่มีไฟล์arrived ยก AirflowSkipExceptionเดือนสิงหาคม — ไฟล์ผิดgate ยก BadFilearrivedreceive_areceive_breceive_cgatetransformloadnotifyskipped 1skipped 0skipped 0skipped 0skipped 0skipped 0skipped 0skipped 0success 1success 1success 1success 1failed 1upstream_failed 0upstream_failed 0upstream_failed 0สถานะของรอบsuccessfailedไม่มีใครถูกปลุกต้องมีคนไปขอไฟล์ใหม่ไหลลงปลายน้ำแบบเดียวกันทั้งคู่
FIG 9.1 — ท่อเดียวกัน สองรอบ คัดจากผลของ states.py ในหัวข้อ 9.2–9.3 · ไม่มีอะไรให้ทำ — การข้ามไหลลงปลายน้ำเป็น skipped ขั้นปลายสุดที่ข้ามนับเป็นผ่าน รอบจึงเขียว · ไฟล์ผิด — ด่านพัง ปลายน้ำเป็น upstream_failed รอบแดง · ทั้งสองแบบไม่มีอะไรถูกโหลด ต่างกันที่ว่าต้องมีคนลงมือไหม

💡 skipped — "ไม่ได้ทำ และไม่มีอะไรผิด" · หน้า Tasks ทางการยกตัวอย่างไว้ตรงตัว: ข้ามเมื่อรู้ว่าไม่มีข้อมูล

การข้ามไหลลงปลายน้ำผ่าน trigger rule ตั้งต้น all_success เหมือนความพัง (หน้า Dags ทางการ) · ต่างที่ปลายทาง: ขั้นปลายสุด skipped นับเป็นผ่าน รอบจึงเขียว ไม่มีใครถูกปลุก

เลือกด้วยคำถามเดียว: ถ้าเรื่องนี้เกิด ต้องมีคนลงมือไหม — ต้อง = failed (ไฟล์ผิดรูป ต้องขอไฟล์ใหม่) · ไม่ต้อง เพราะรอบหน้าจะเก็บเอง = skipped (ยังไม่มีใครส่ง — ตราบที่มีรอบถัดไปมาเก็บ บทที่ 13–14)

⚠️ arrived ข้ามเฉพาะเมื่อไม่มีไฟล์สักสาขา · ขาดสาขาเดียว receive ของสาขานั้นพัง — ข้ามไม่ได้ เพราะรอบเขียวที่ขาดหนึ่งสาขาคือยอดผิดแบบเงียบ ๆ ของบทที่ 1 · ความพังแบบนี้รอแล้วอาจหาย — บทถัดไป

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

  1. ท่อรับไฟล์ CSV ยอดขายรายวันจากหลายร้าน · วันอาทิตย์ร้านปิด ไม่มีไฟล์ · วันหนึ่งไฟล์มาแต่ไม่มีคอลัมน์ยอดขาย · อีกวันมีร้านใหม่โผล่มา — แต่ละกรณีควร skipped failed หรือปล่อยผ่าน
  2. เพื่อนแก้ gate ให้ยก AirflowSkipException แทน error เมื่อไฟล์ผิด "จะได้ไม่มีรอบแดงรกตา" · รอบเดือนสิงหาคมจะเป็นสีอะไร และเช้านั้นใครจะรู้ว่าสาขา ค มีช่องว่าง
  3. ไฟล์เดือนกรกฎาคมของบทที่ 1 (ซ่อนทั้ง 33 แถว) ผ่านด่านนี้ไหม · ถ้ายังใช้ตัวอ่านเก่าที่ข้ามแถวซ่อน ด่านจะทำอะไร

Read the full book