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

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

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

ท่อรับไฟล์รายเดือน ฉบับ Airflow

เดือนกรกฎาคมของบทที่ 1 รายงานว่าสาขา ก เหลือ 0 ชิ้น ทั้งที่ไฟล์มี 121 ชิ้น · บทนี้ประกอบท่อเจ็ดขั้นให้อ่านไฟล์จริง ได้ 121 ชิ้นครบพร้อมบอกว่ามีแถวซ่อน · และไขปริศนาที่ค้างมาสองข้อ — ตัวอ่านควรทำยังไงกับแถวที่ซ่อน และทำไมห้ามส่งข้อมูลก้อนใหญ่ผ่าน XCom

8.1 แยกตรรกะออกจากการต่อสาย

ท่อนี้มีสองไฟล์: stock_steps.py เก็บตรรกะของทุกขั้นเป็นฟังก์ชันธรรมดาที่ไม่ import อะไรจาก Airflow · stock_monthly.py เป็นไฟล์ DAG ที่แค่เรียกฟังก์ชันพวกนั้นแล้วต่อเส้น · ได้สองอย่าง: ทดสอบงานได้โดยไม่ต้องมี Airflow (บล็อกในหัวข้อ 8.3 รันบน Windows ด้วย Python เปล่า ๆ) และไฟล์ DAG เบา — ไม่มีการอ่าน Excel ที่ชั้นบนสุด · ที่ยังถูกอ่านซ้ำทุกราว 30 วินาทีคือตัว stock_steps.py เอง (มันอยู่ใน dags ด้วย) กับการ import openpyxl ทุกครั้งที่อ่าน stock_monthly.py ราว 0.1 วินาทีในแล็บ — ยังเบา แต่ถ้าเป็นไลบรารีหนักอย่าง pandas หน้า Best Practices แนะนำให้ import ในฟังก์ชัน

8.2 ไฟล์ปลอมของทั้งเล่ม

ตัวสร้างนี้เขียนไฟล์นับสต็อกของร้านสมมติสามสาขา ห้าเดือน (พฤษภาคม–กันยายน 2569) — ทุกไฟล์ในเล่มตั้งแต่บทนี้มาจากมัน · QUIRKS คืออาการแบบไฟล์จริง หนึ่งไฟล์หนึ่งอาการ ให้แต่ละบทลองทีละเรื่อง · seed ตายตัว รันกี่ครั้งก็ได้ข้อมูลเดิมทุกช่อง (ตัวไฟล์ต่างกันแค่เวลาบันทึกที่ฝังอยู่ในไฟล์) · รันเปล่า ๆ เขียนลงโฟลเดอร์ชั่วคราว ใส่ path เขียนลงโฟลเดอร์นั้น

# ~/airflow-lab/make_fake_files.py
"""สร้างไฟล์นับสต็อกรายเดือนของร้านเครื่องเขียนสมมติ 3 สาขา — ข้อมูลปลอมทั้งหมด หน่วยเป็นชิ้น

รันเปล่า ๆ = เขียนลงโฟลเดอร์ชั่วคราวแล้วพิมพ์สรุป · ใส่ path = เขียนลงโฟลเดอร์นั้น
"""

import random
import sys
import tempfile
from datetime import date
from pathlib import Path

from openpyxl import Workbook

SEED = 2569
BRANCHES = {"a": ("สาขา ก", 120), "b": ("สาขา ข", 340), "c": ("สาขา ค", 560)}  # ยอดเดือนแรก
MONTHS = [date(2026, 5, 31), date(2026, 6, 30), date(2026, 7, 31), date(2026, 8, 31), date(2026, 9, 30)]
CATALOG = {
    "ปากกา": ["ปากกาลูกลื่น", "ปากกาเจล", "ปากกาเน้นข้อความ"],
    "ดินสอ": ["ดินสอไม้ 2B", "ดินสอกด", "ยางลบ"],
    "กระดาษ": ["สมุดปกอ่อน", "กระดาษ A4", "กระดาษโน้ต"],
    "แฟ้ม": ["แฟ้มสอด", "แฟ้มห่วง"],
}
PACKS = ["แพ็ก 1", "แพ็ก 6", "แพ็ก 12"]
COLUMNS = ["รหัส", "ชื่อสินค้า", "หมวด", "จำนวน"]

# อาการที่ไฟล์จริงเป็น — หนึ่งไฟล์หนึ่งอาการ ให้แต่ละบทลองทีละเรื่อง
QUIRKS = {
    ("2026-06", "b"): "columns_moved",  # คอลัมน์สลับที่ มีคอลัมน์เพิ่ม และชื่อชีตเปลี่ยน
    ("2026-07", "a"): "autofilter_stuck",  # ตัวกรอง "จำนวน = 0" ค้าง ทุกแถวข้อมูลติดธงซ่อน
    ("2026-07", "b"): "second_sheet",  # ข้อมูลอยู่ชีตที่สอง และมีแถวคำอธิบายใต้หัวตาราง
    ("2026-08", "a"): "bad_date",  # ช่องวันที่นับเป็นข้อความ อ่านเป็นวันที่ไม่ได้
    ("2026-08", "b"): "new_category",  # หมวดใหม่โผล่มาเดือนนี้
    ("2026-08", "c"): "nbsp_blank",  # ช่องจำนวนที่ดูว่าง แต่มี non-breaking space
    ("2026-09", "a"): "missing_column",  # คอลัมน์ "จำนวน" หาย ด่านต้องหยุดท่อ
    ("2026-09", "c"): "late",  # ไฟล์มาช้า — แยกไว้ในโฟลเดอร์ late/
}


def catalog(quirk):
    items = [(cat, f"{name} {pack}") for cat, names in CATALOG.items() for name in names for pack in PACKS]
    if quirk == "new_category":
        items += [("อุปกรณ์ศิลปะ", f"สีไม้ {pack}") for pack in PACKS]
    return [(f"S{i:03d}", cat, name) for i, (cat, name) in enumerate(items, start=1)]


def split(total, n, rng):
    """แบ่ง total ชิ้นให้ n รายการ ทุกรายการได้อย่างน้อย 1 ชิ้น"""
    counts = [1] * n
    for _ in range(total - n):
        counts[rng.randrange(n)] += 1
    return counts


def make_rows(total, quirk, rng):
    items = catalog(quirk)
    counts = split(total, len(items), rng)
    rows = [
        {"รหัส": code, "ชื่อสินค้า": name, "หมวด": cat, "จำนวน": n}
        for (code, cat, name), n in zip(items, counts)
    ]
    if quirk == "nbsp_blank":
        for row in rng.sample(rows, 3):
            row["จำนวน"] = "\u00a0"
    if quirk == "missing_column":
        for row in rows:
            del row["จำนวน"]
    return rows


def write_book(path, title, month_end, rows, quirk):
    wb = Workbook()
    ws = wb.active
    ws.title = "นับสต็อก" if quirk == "columns_moved" else "สต็อก"
    if quirk == "second_sheet":
        ws.title = "หน้าปก"
        ws["A1"] = "ข้อมูลอยู่ชีตถัดไป"
        ws = wb.create_sheet("ข้อมูล")
    columns = list(COLUMNS)
    if quirk == "columns_moved":
        columns = ["ชื่อสินค้า", "รหัส", "จำนวน", "หมวด", "หมายเหตุ"]
    if quirk == "missing_column":
        columns.remove("จำนวน")
    ws["A1"] = f"นับสต็อก {title}"
    ws["A2"] = "วันที่นับ"
    ws["B2"] = "สิ้นเดือน ส.ค." if quirk == "bad_date" else month_end
    for letter, width in zip("ABCDE", [12, 24, 14, 10, 12]):
        ws.column_dimensions[letter].width = width
    header_row = 4  # แถว 3 เว้นว่าง
    for col, name in enumerate(columns, start=1):
        ws.cell(header_row, col, name)
    first = header_row + 1
    if quirk == "second_sheet":
        ws.cell(first, 1, "(หน่วย: ชิ้น · นับ ณ สิ้นเดือน)")
        first += 1
    for r, row in enumerate(rows, start=first):
        for col, name in enumerate(columns, start=1):
            ws.cell(r, col, row.get(name))
    last = first + len(rows) - 1
    if quirk == "autofilter_stuck":
        ws.auto_filter.ref = f"A{header_row}:{ws.cell(header_row, len(columns)).column_letter}{last}"
        ws.auto_filter.add_filter_column(columns.index("จำนวน"), ["0"])
        for r in range(first, last + 1):
            ws.row_dimensions[r].hidden = True
    wb.save(path)
    hidden = sum(1 for r in range(first, last + 1) if ws.row_dimensions[r].hidden)
    return wb.sheetnames, hidden


def main(out_dir):
    rng = random.Random(SEED)
    lines = []
    for i, month_end in enumerate(MONTHS):
        month = f"{month_end:%Y-%m}"
        for code, (title, base) in BRANCHES.items():
            quirk = QUIRKS.get((month, code), "-")
            total = base if i == 0 else round(base * rng.uniform(0.8, 1.2))
            rows = make_rows(total, quirk, rng)
            folder = out_dir / "late" if quirk == "late" else out_dir
            folder.mkdir(parents=True, exist_ok=True)
            name = f"stock_{month}_{code}.xlsx"
            sheets, hidden = write_book(folder / name, title, month_end, rows, quirk)
            quantities = [row.get("จำนวน") for row in rows]
            pieces = sum(n for n in quantities if isinstance(n, int))
            where = "late/" if quirk == "late" else ""
            lines.append(
                f"{where}{name} · ชีต {'/'.join(sheets)} · {len(rows)} แถว · ซ่อน {hidden} · "
                f"{pieces} ชิ้น · {quirk}"
            )
    print("\n".join(lines))
    first = [base for _, base in BRANCHES.values()]
    print(f"เดือนแรก: {' + '.join(map(str, first))} = {sum(first)} ชิ้น")


if __name__ == "__main__":
    if len(sys.argv) > 1:
        main(Path(sys.argv[1]))
    else:
        main(Path(tempfile.mkdtemp()))

ผลรันจริง

stock_2026-05_a.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 120 ชิ้น · -
stock_2026-05_b.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 340 ชิ้น · -
stock_2026-05_c.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 560 ชิ้น · -
stock_2026-06_a.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 140 ชิ้น · -
stock_2026-06_b.xlsx · ชีต นับสต็อก · 33 แถว · ซ่อน 0 · 364 ชิ้น · columns_moved
stock_2026-06_c.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 610 ชิ้น · -
stock_2026-07_a.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 33 · 121 ชิ้น · autofilter_stuck
stock_2026-07_b.xlsx · ชีต หน้าปก/ข้อมูล · 33 แถว · ซ่อน 0 · 382 ชิ้น · second_sheet
stock_2026-07_c.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 448 ชิ้น · -
stock_2026-08_a.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 139 ชิ้น · bad_date
stock_2026-08_b.xlsx · ชีต สต็อก · 36 แถว · ซ่อน 0 · 406 ชิ้น · new_category
stock_2026-08_c.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 550 ชิ้น · nbsp_blank
stock_2026-09_a.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 0 ชิ้น · missing_column
stock_2026-09_b.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 320 ชิ้น · -
late/stock_2026-09_c.xlsx · ชีต สต็อก · 33 แถว · ซ่อน 0 · 559 ชิ้น · late
เดือนแรก: 120 + 340 + 560 = 1020 ชิ้น

บรรทัดที่เจ็ดคือไฟล์ของบทที่ 1: stock_2026-07_a.xlsx 33 แถว ซ่อนทั้ง 33 แถว 121 ชิ้น · บรรทัดสุดท้ายคือยอดเดือนแรกของบทที่ 1 กับ 6 · "\u00a0" ใน make_rows คือช่องว่างที่ไม่ใช่ช่องว่างธรรมดา (non-breaking space) — ตาเห็นว่าว่าง โค้ดเห็นว่ามีของ ด่านตรวจของบทที่ 9 ต้องจับให้ได้

8.3 ตัวอ่านที่ไม่เชื่อตำแหน่ง และไม่ทิ้งแถวซ่อน

⚠️ ความเข้าใจผิด: "แถวที่ซ่อนอยู่ไม่ต้องนับ — กรองทิ้งตั้งแต่ตอนอ่านได้" — ผู้เขียนเจอทั้งสองทางในงานจริง: แถวซ่อนที่เป็นของเลิกใช้ไหลขึ้นยอดรวม และแถวซ่อนที่เป็นข้อมูลจริงถูกกรองทิ้งเงียบ ๆ แบบเดือนกรกฎาคม · ผิดเพราะธงซ่อนบอกแค่ว่ามีคนซ่อนแถวนั้น ไม่ได้บอกว่าทำไม · กลไกที่ถูกคืออ่านทุกแถว → ติดป้ายว่าแถวไหนซ่อน → ตัดสินทีละเรื่อง · พิสูจน์ได้ที่ผลรันข้างล่าง

read_branch หาหัวตารางด้วยชื่อ ไม่ใช่ตำแหน่ง (ชีตไหนก็ได้ แถวไหนก็ได้ที่มีทั้ง "รหัส" และ "ชื่อสินค้า") และนับเป็นข้อมูลเฉพาะแถวที่มีทั้งสองช่อง แถวคำอธิบายใต้หัวตารางจึงหลุดไปเอง · stage เขียนแถวทั้งหมดลงไฟล์ แล้วคืนใบสรุปเล็ก ๆ · ท้ายบล็อกลองทุกขั้นกับไฟล์เล็กสองไฟล์ที่มีอาการแบบของจริง โดยไม่มี Airflow

# ~/airflow-lab/dags/stock_steps.py — ตรรกะของทุกขั้นเป็นฟังก์ชันธรรมดา ไม่ import อะไรจาก Airflow
import json
import sqlite3
from pathlib import Path

from openpyxl import load_workbook

REQUIRED = ["รหัส", "ชื่อสินค้า", "หมวด", "จำนวน"]
NAMES = {"a": "ก", "b": "ข", "c": "ค"}


def find_header(book):
    """หาชีตกับแถวหัวตารางด้วยชื่อคอลัมน์ ไม่เชื่อตำแหน่ง → (ชีต, เลขแถว, {ชื่อคอลัมน์: เลขคอลัมน์})"""
    for ws in book.worksheets:
        for r in range(1, min(ws.max_row, 20) + 1):
            cols = {ws.cell(r, c).value: c for c in range(1, ws.max_column + 1)}
            if "รหัส" in cols and "ชื่อสินค้า" in cols:
                return ws, r, {name: c for name, c in cols.items() if name}
    raise ValueError("หาแถวหัวตารางไม่เจอในทุกชีต")


def read_branch(path):
    """อ่านทุกแถวใต้หัวตาราง แล้วติดป้ายว่าแถวไหนซ่อน — ไม่กรองแถวซ่อนทิ้งตอนอ่าน"""
    ws, header, cols = find_header(load_workbook(path))
    rows = []
    for r in range(header + 1, ws.max_row + 1):
        row = {name: ws.cell(r, c).value for name, c in cols.items()}
        if row["รหัส"] is None or row["ชื่อสินค้า"] is None:
            continue  # แถวคำอธิบายใต้หัวตาราง หรือแถวว่าง
        row["ซ่อน"] = bool(ws.row_dimensions[r].hidden)
        rows.append(row)
    return rows, [name for name in REQUIRED if name not in cols]


def stage(branch, rows, missing, path):
    """เขียนแถวทั้งหมดลงไฟล์ แล้วคืนแค่ใบสรุปเล็ก ๆ — ใบนี้คือสิ่งที่เดินทางผ่าน XCom"""
    path.parent.mkdir(parents=True, exist_ok=True)
    path.write_text(json.dumps(rows, ensure_ascii=False), encoding="utf-8")
    hidden = sum(row["ซ่อน"] for row in rows)
    return {"branch": branch, "path": str(path), "rows": len(rows), "hidden": hidden, "missing": missing}


def check_columns(staged):
    for s in staged:
        if s["missing"]:
            raise ValueError(f"สาขา {NAMES[s['branch']]} ไม่มีคอลัมน์ {', '.join(s['missing'])}")
    return staged


def transform(staged, path):
    rows = []
    for s in staged:
        for row in json.loads(Path(s["path"]).read_text(encoding="utf-8")):
            rows.append([s["branch"], row["รหัส"], int(row["จำนวน"]), row["ซ่อน"]])
    path.write_text(json.dumps(rows, ensure_ascii=False), encoding="utf-8")
    pieces = {s["branch"]: sum(row[2] for row in rows if row[0] == s["branch"]) for s in staged}
    return {"path": str(path), "rows": len(rows), "pieces": pieces}


def load(summary, db_path, month):
    """⚠️ ต่อท้ายตาราง — รันซ้ำแล้วจะเกิดอะไรขึ้น คือเรื่องของบทที่ 12"""
    rows = json.loads(Path(summary["path"]).read_text(encoding="utf-8"))
    db = sqlite3.connect(db_path)
    db.execute("CREATE TABLE IF NOT EXISTS stock (month, branch, code, qty, hidden)")
    db.executemany("INSERT INTO stock VALUES (?, ?, ?, ?, ?)", [[month, *row] for row in rows])
    db.commit()
    return db.execute("SELECT COUNT(*) FROM stock WHERE month = ?", [month]).fetchone()[0]


def report(staged, summary, in_db):
    each = " · ".join(f"{NAMES[b]} {n:,}" for b, n in summary["pieces"].items())
    total = sum(summary["pieces"].values())
    lines = [f"สต็อกรวม {total:,} ชิ้น ({each}) จาก {summary['rows']} แถว · ในคลังเดือนนี้ {in_db} แถว"]
    lines += [f"⚠️ สาขา {NAMES[s['branch']]} ซ่อนอยู่ {s['hidden']} จาก {s['rows']} แถว — เปิดไฟล์ดูก่อนเชื่อยอด"
              for s in staged if s["hidden"]]
    return "\n".join(lines)


if __name__ == "__main__":  # ทดลองทุกขั้นโดยไม่มี Airflow — ไฟล์เล็กสองไฟล์ที่มีอาการแบบของจริง
    import tempfile

    from openpyxl import Workbook

    work = Path(tempfile.mkdtemp())
    book = Workbook()  # สาขา ก: ตัวกรองค้าง ทุกแถวข้อมูลติดธงซ่อน
    ws = book.active
    ws.append(["นับสต็อก สาขา ก"])
    ws.append([])
    ws.append([])
    ws.append(["รหัส", "ชื่อสินค้า", "หมวด", "จำนวน"])
    for i in range(1, 4):
        ws.append([f"S{i:03d}", f"ปากกา {i}", "ปากกา", 10 * i])
        ws.row_dimensions[ws.max_row].hidden = True
    book.save(work / "a.xlsx")
    book = Workbook()  # สาขา ข: ข้อมูลอยู่ชีตที่สอง คอลัมน์สลับที่ มีแถวคำอธิบายใต้หัวตาราง
    book.active.title = "หน้าปก"
    ws = book.create_sheet("ข้อมูล")
    ws.append(["ชื่อสินค้า", "รหัส", "จำนวน", "หมวด"])
    ws.append(["(หน่วย: ชิ้น)"])
    ws.append(["ดินสอ 1", "S004", 7, "ดินสอ"])
    ws.append(["ดินสอ 2", "S005", 5, "ดินสอ"])
    book.save(work / "b.xlsx")

    staged = [stage(b, *read_branch(work / f"{b}.xlsx"), work / f"{b}.json") for b in "ab"]
    for s in staged:
        print({k: v for k, v in s.items() if k != "path"})
    summary = transform(check_columns(staged), work / "all.json")
    print(report(staged, summary, load(summary, work / "warehouse.db", "2026-07")))

ผลรันจริง

{'branch': 'a', 'rows': 3, 'hidden': 3, 'missing': []}
{'branch': 'b', 'rows': 2, 'hidden': 0, 'missing': []}
สต็อกรวม 72 ชิ้น (ก 60 · ข 12) จาก 5 แถว · ในคลังเดือนนี้ 5 แถว
⚠️ สาขา ก ซ่อนอยู่ 3 จาก 3 แถว — เปิดไฟล์ดูก่อนเชื่อยอด

สาขา ก ซ่อนทั้งสามแถวแต่ยังนับครบ 60 ชิ้นพร้อมบรรทัดเตือน · สาขา ข อยู่ชีตที่สอง คอลัมน์สลับที่ มีแถวคำอธิบาย — อ่านได้สองแถวถูกต้อง · ท่อนี้ตัดสินว่า ยอดรวมนับทุกแถว เพราะซ่อนไม่ใช่ลบ แต่รายงานแยกบอกว่าซ่อนกี่แถว ให้คนที่รู้จักสาขาตัดสินว่าเป็นของเลิกขายหรือตัวกรองค้าง

8.4 ต่อสาย แล้วรันกับไฟล์เดือนกรกฎาคม

# ~/airflow-lab/dags/stock_monthly.py — ต่อสายอย่างเดียว ตรรกะทั้งหมดอยู่ใน stock_steps.py ข้าง ๆ
from datetime import datetime
from pathlib import Path

import stock_steps as steps
from airflow.sdk import dag, task

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 receive(branch, **context):
        month = context["params"]["month"]
        rows, missing = steps.read_branch(LAB / "inbox" / f"stock_{month}_{branch}.xlsx")
        return steps.stage(branch, rows, missing, LAB / "staging" / month / f"{branch}.json")

    @task
    def merge_check(staged):
        for s in staged:  # สิ่งที่มาทาง XCom มีแค่ใบสรุป ไม่มีแถวข้อมูลสักแถว
            print(f"merge_check ← receive_{s['branch']}: {Path(s['path']).name} · {s['rows']} แถว · ซ่อน {s['hidden']}")
        return steps.check_columns(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))

    staged = merge_check([receive.override(task_id=f"receive_{b}")(b) for b in "abc"])
    summary = transform(staged)
    notify(staged, summary, load(summary))


stock_monthly()

โครงเดียวกับ stock_first เนื้อในแต่ละ @task เหลือแค่การเรียก steps · ต่างกันเส้นเดียว: notify รอสามขั้น เพราะรายงานใช้ทั้งใบสรุป ยอด และจำนวนแถวในคลัง · เดือนส่งมาเป็น params (อ่านจาก context ที่ฟังก์ชันรับเป็น **context) ค่าตั้งต้น "2026-07" — ทำไมไม่ใช้วันที่ของรอบบอกเดือน ซับซ้อนกว่าที่คิด บทที่ 13 จะกลับมา

เตรียมแล็บแล้วรัน — venv ของแล็บยังไม่มี openpyxl ลงด้วยไฟล์ constraint ชุดเดิม (ได้ 3.1.5 รุ่นเดียวกับบน Windows)

uv pip install openpyxl --constraint "https://raw.githubusercontent.com/apache/airflow/constraints-3.3.2/constraints-3.12.txt"
python ~/airflow-lab/make_fake_files.py ~/airflow-lab/inbox   # ผลเดียวกับหัวข้อ 8.2
# วาง stock_steps.py กับ stock_monthly.py ไว้ใน ~/airflow-lab/dags/
airflow dags test stock_monthly 2>&1 | grep -E "ชิ้น|แถว"

ผลรันจริง

merge_check ← receive_a: a.json · 33 แถว · ซ่อน 33
merge_check ← receive_b: b.json · 33 แถว · ซ่อน 0
merge_check ← receive_c: c.json · 33 แถว · ซ่อน 0
สต็อกรวม 951 ชิ้น (ก 121 · ข 382 · ค 448) จาก 99 แถว · ในคลังเดือนนี้ 99 แถว
⚠️ สาขา ก ซ่อนอยู่ 33 จาก 33 แถว — เปิดไฟล์ดูก่อนเชื่อยอด

สาขา ก ได้ 121 ชิ้นครบ — ไฟล์เดียวกับที่ทำให้รายงานของบทที่ 1 บอก 0 ชิ้น · สาขา ข อยู่ชีตที่สองพร้อมแถวคำอธิบาย อ่านได้ 33 แถวถูกต้อง · และบรรทัดเตือนบอกคนอ่านรายงานว่ามีอะไรผิดปกติ แทนที่จะเงียบ

8.5 อะไรเดินทางผ่าน XCom อะไรอยู่บนดิสก์

สามบรรทัดแรกของผลรันพิมพ์จากใบสรุปที่ merge_check ได้รับทาง XCom — ใบละห้าช่อง (สาขา · path · จำนวนแถว · แถวซ่อน · คอลัมน์ที่ขาด) ไม่มีแถวข้อมูลสักแถว · ทั้ง 99 แถวอยู่บนดิสก์ตลอดทาง

ขั้นผ่าน XCom — ใบสรุปเล็ก ๆบนดิสก์ใต้ ~/airflow-lab — ข้อมูลเต็มreceive_a · b · cใบสรุป a.json · 33 แถว · ซ่อน 33สาขาละใบ — ข กับ ค ซ่อน 0อ่านinbox/stock_2026-07_a.xlsxเขียนstaging/2026-07/a.json33 แถวต่อสาขาmerge_checkใบสรุปสามใบเดิม ส่งต่อไม่แตะดิสก์transformใบสรุป all.json · 99 แถวก 121 · ข 382 · ค 448 ชิ้นอ่านstaging/2026-07/a·b·c.jsonเขียนstaging/2026-07/all.json99 แถวload99แถวในคลังเดือนนี้อ่านstaging/2026-07/all.jsonเขียนwarehouse.dbต่อท้าย 99 แถวnotifyไม่คืนค่าไม่แตะดิสก์ · พิมพ์รายงาน
FIG 8.1 — ท่อ stock_monthly ทีละขั้น คัดจากผลรันในหัวข้อ 8.4 · สิ่งที่ผ่าน XComมีแค่ใบสรุป — ชื่อไฟล์กับตัวเลขไม่กี่ตัว ซึ่งถูกเขียนลงฐานข้อมูลของ Airflow · แถวข้อมูลทั้งหมดอยู่บนดิสก์ที่คุณเป็นเจ้าของ ขั้นถัดไปอ่านจาก path ในใบสรุปเอง · merge_check กับ notify ไม่แตะดิสก์เลย

⚠️ ความเข้าใจผิด: "ส่งข้อมูลก้อนใหญ่ระหว่างขั้นผ่าน XCom ได้ return ทั้ง DataFrame ไปเลย" — หน้า Best Practices ทางการให้ใช้ XCom ส่ง "ข้อความเล็ก ๆ" ระหว่างขั้น · ผิดเพราะหน้า XComs ระบุว่ามันออกแบบมาสำหรับข้อมูลปริมาณน้อย อย่าใช้ส่งค่าใหญ่อย่าง dataframe และที่เก็บตั้งต้นคือฐานข้อมูลของ Airflow — ที่เดียวกับสถานะของทุกช่อง · กลไกที่ถูกคือเขียนข้อมูลลงที่เก็บของตัวเอง แล้วส่ง path ผ่าน XCom แบบที่ Best Practices แนะนำ · พิสูจน์เอง: สั่ง airflow dags test stock_monthly 2>&1 | grep "Returned value" — ทุกขั้นพิมพ์ค่าที่มันคืนลง XCom บรรทัดละขั้น ของ receive ทั้งสามคือใบสรุปเท่านั้น (รอบที่รันด้วย dags test ไม่มีไฟล์บันทึกให้กดดูใน Grid ต้องดูจากหน้าต่างที่สั่ง)

บทที่ 2 บอกว่าผลของแต่ละขั้นต้องอยู่นอกโปรเซส · บทนี้เติมครึ่งหลัง: ข้อมูลอยู่ในที่เก็บของคุณ ใบสั่งให้ไปหยิบอยู่ใน XCom

8.6 สิ่งที่ท่อนี้ยังทำไม่ได้

  1. load ต่อท้ายตาราง — รันซ้ำแล้ว "ในคลังเดือนนี้" จะเป็นเท่าไหร่ เดาก่อนเปิดเฉลยข้อ 3 → บทที่ 12
  2. ด่านยังหยาบ — check_columns หยุดท่อได้แค่เมื่อคอลัมน์หาย · ไฟล์ศูนย์แถวกับหมวดใหม่ผ่านด่านไปได้ทั้งคู่ · ช่องว่างแปลก ๆ ก็ผ่านด่าน แล้วไปพังใน transform แบบไม่บอกว่าสาขาไหนแถวไหน (ลอง -c '{"month": "2026-08"}') → บทที่ 9
  3. ไฟล์สาขา ค เดือนกันยายนมาช้า (อยู่ใน late/) — ตั้งเวลาไว้ตีสองแล้วไฟล์ยังไม่มา ท่อนี้ทำได้แค่พัง · ลอง -c '{"month": "2026-09"}' จะเห็นว่า receive_c พังเพราะหาไฟล์ไม่เจอ ก่อนที่ด่านจะได้เห็นว่าไฟล์สาขา ก เดือนนั้นไม่มีคอลัมน์ "จำนวน" → บทที่ 14
  4. เดือนถูกส่งด้วยมือ — ทุกเดือนต้องมีคนพิมพ์ -c → บทที่ 13

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

  1. ระบบหนึ่งวางไฟล์ CSV ยอดขายรายวันไว้ในโฟลเดอร์ทุกเช้า · ร่างท่อสามขั้น (รับไฟล์ · ตรวจ · โหลด) แบบบทนี้: แต่ละขั้นเป็นฟังก์ชันอะไร คืนอะไร และอะไรที่ไม่ควรเดินทางผ่าน XCom
  2. ถ้าแก้ receive ให้คืน rows ทั้งก้อนแทนใบสรุป ท่อยังรันผ่าน — แล้วเสียอะไร · และทำไมการทำแบบนั้นยัง "ถูก" ตามบทเรียนของบทที่ 2 แต่ผิดตามบทนี้
  3. สั่ง airflow dags test stock_monthly 2>&1 | grep -E "ชิ้น|แถว" ซ้ำอีกครั้งทันทีโดยไม่ล้างอะไร บรรทัดที่สี่จะเปลี่ยนไปยังไง และอาการนี้คืออาการเดียวกับเช้าวันไหนในบทที่ 1

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