ข้ามไปยังเนื้อหา
Tayakorn
← คู่มือทั้งหมด
dataอัปเดต 2026-10-06อ่านฟรี

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

เข้าใจว่าตัวคุมลำดับงานจำอะไรไว้แทนคุณ ตั้งแต่สคริปต์ Python ที่พังกลางทางจนถึง DAG ของ Airflow 3 อ่านจบแล้วเขียนท่อรับไฟล์รายเดือนที่พังแล้วรันต่อจากจุดที่พังได้ และรันซ้ำแล้วข้อมูลไม่เบิ้ล

#airflow#orchestration#data-pipeline#python#idempotency

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

อ่านจบแล้วคุณจะ

  1. ดูงานหลายขั้นตรงหน้าแล้วตัดสินได้ว่าต้องมีตัวคุมไหม หรือ Task Scheduler กับ cron ก็พอ
  2. เขียน DAG ของ Airflow 3 สำหรับท่อรับไฟล์รายเดือนได้เอง ที่มีด่านตรวจก่อนโหลด ลองใหม่เฉพาะความพังชั่วคราว และรันซ้ำแล้วข้อมูลไม่เบิ้ล
  3. เช้าวันที่มีขั้นพังตอนกลางคืน เปิดหน้า Grid แล้วบอกได้ว่าพังที่ไหน ครั้งที่เท่าไหร่ เพราะอะไร และสั่งรันเฉพาะส่วนที่ต้องรัน
  4. อ่านบทเรียนออนไลน์แล้วรู้ว่าเป็นโค้ดรุ่น 2 และแปลงเป็นรุ่น 3 ได้
  • ระดับ: เริ่มจากศูนย์ฝั่งตัวคุมลำดับงาน — ต้องเขียน Python ได้แล้ว (อ่านไฟล์ วนลูป เขียนฟังก์ชัน)
  • เวลาอ่าน: ~75 นาทีสำหรับบทที่ 1–4 (เล่มนี้ทยอยเขียน บทที่ 5–16 ตามมา)
  • ภาษา: ไทย พร้อมโค้ดที่รันได้จริง · บทที่ 1–4 เป็น Python ล้วน รันบน Windows ได้ทันที · บทที่ใช้ Airflow ปักรุ่น 3.3.2
  • ของจริง: ผลรันทุกบล็อกมาจากการรันจริงด้วย Python 3.13.5 บน Windows (บล็อกที่อ่านไฟล์ Excel ใช้ openpyxl 3.1.5) · ข้อมูลตัวอย่างทั้งเล่มเป็นของปลอม — ร้านเครื่องเขียนสมมติสามสาขา หน่วยเป็นชิ้น

เคยเขียนตัวรันที่จำได้ว่าทำถึงไหนแล้วมาก่อน? ข้ามบทที่ 2 กับ 3 ไปบทที่ 4 ได้เลย แต่อย่าข้ามบทที่ 4 — ทุกบทหลังจากนั้นอ้างตารางในบทนั้น

LEVEL 0 · ปูพื้นจากศูนย์

คืนที่สคริปต์พังตรงกลาง#

ทุกต้นเดือน ร้านเครื่องเขียนสามสาขาส่งไฟล์ Excel นับสต็อกมาที่คุณ สาขาละไฟล์ แล้วสคริปต์ Python ตัวหนึ่งที่คุณเขียนไว้ก็ทำงานต่อจนจบ: รับไฟล์สามสาขา → รวมกันและตรวจว่าคอลัมน์ครบ → แปลงให้อยู่ในรูปเดียวกัน → โหลดเข้าฐานข้อมูล → ส่งรายงานยอดสต็อกรวม · Task Scheduler ปลุกมันตอนตีสองของวันที่ 1 ทุกเดือน และมันก็ทำแบบนั้นมาหลายเดือนโดยไม่มีใครต้องนึกถึงมัน

ร้านนี้ไม่มีอยู่จริง ชื่อสินค้าและตัวเลขทุกตัวในเล่มตั้งขึ้นใหม่ แต่รูปร่างของปัญหาทุกข้อที่จะเจอเป็นของจริง

1.1 เช้าวันที่ไฟล์ "ว่าง"

เช้าวันที่ 1 สิงหาคม รายงานยอดสต็อกเดือนกรกฎาคมถูกส่งออกไปแล้วตามเวลา สคริปต์จบโดยไม่มี error สักบรรทัด และรายงานบอกว่า สาขา ก เหลือของ 0 ชิ้น

คุณเปิดไฟล์ของสาขา ก ใน Excel แล้วเห็นแบบนี้

หน้าต่าง Excel เปิดไฟล์นับสต็อกสาขา ก เดือนกรกฎาคม เห็นแถว 1 ถึง 4 คือชื่อสาขา วันที่นับ 2026-07-31 แถวว่าง และหัวตาราง รหัส ชื่อสินค้า หมวด จำนวน จากนั้นเลขแถวกระโดดไปที่ 38 ปุ่มลูกศรที่หัวคอลัมน์จำนวนเป็นรูปกรวย และแถบสถานะด้านล่างขึ้นคำว่า Filter Mode

ชื่อสาขา วันที่นับ แล้วก็หัวตาราง ไม่มีสินค้าสักแถว ข้อสรุปที่ตามมาแทบจะอัตโนมัติคือสาขาส่งไฟล์เปล่ามา แล้วคุณก็พิมพ์ข้อความไปขอไฟล์ใหม่

ก่อนกดส่ง ลองดูเลขแถวทางซ้ายอีกครั้ง: 1 · 2 · 3 · 4 แล้วกระโดดไป 38 — แถว 5 ถึง 37 หายไปไหน · ปุ่มที่หัวคอลัมน์จำนวนไม่ใช่ลูกศรธรรมดาแต่เป็นรูปกรวย และแถบสถานะด้านล่างเขียนว่า Filter Mode

ไฟล์นี้ไม่ได้ว่าง มันถูกเซฟทั้งที่ตัวกรองยังค้างอยู่ คนที่สาขาเปิดตัวกรองหาสินค้าที่ของหมด (จำนวน = 0) แล้วเซฟปิดไปโดยไม่ได้ล้างตัวกรอง · เดือนนี้ไม่มีสินค้าไหนหมดเลยสักรายการ ทุกแถวจึงไม่ผ่านตัวกรองและถูกซ่อนหมด — ซ่อน ไม่ใช่ลบ ข้อมูลยังอยู่ในไฟล์ครบทุกแถว ขนาดไฟล์ก็ปกติ

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

บล็อกนี้สร้างไฟล์หน้าตาเดียวกับของสาขา ก (33 รายการ ตัวกรอง "จำนวน = 0" ค้างอยู่) แล้วเปิดกลับมาอ่านสองแบบ — อ่านทุกแถว กับอ่านเฉพาะแถวที่มองเห็น ซึ่งคือสิ่งที่ตาคุณเห็นบนจอ · ต้องมี openpyxl ก่อน (pip install openpyxl)

import tempfile
from pathlib import Path

from openpyxl import Workbook, load_workbook

# ── สร้างไฟล์หน้าตาเดียวกับของสาขา ก เดือนกรกฎาคม: 33 รายการ ตัวกรอง "จำนวน = 0" ค้างอยู่ ──
path = Path(tempfile.mkdtemp()) / "stock_2026-07_a.xlsx"
wb = Workbook()
ws = wb.active
ws.title = "สต็อก"
ws["A1"] = "นับสต็อก สาขา ก"
ws["A2"], ws["B2"] = "วันที่นับ", "2026-07-31"
for col, name in enumerate(["รหัส", "ชื่อสินค้า", "หมวด", "จำนวน"], start=1):
    ws.cell(4, col, name)  # แถว 4 = หัวตาราง
for i in range(1, 34):
    ws.append([f"S{i:03d}", f"สินค้า {i}", "เครื่องเขียน", 3 + (i % 3 != 0)])
ws.auto_filter.ref = "A4:D37"
ws.auto_filter.add_filter_column(3, ["0"])  # คนก่อนหน้ากรองหาของที่หมด แล้วเซฟทั้งที่ตัวกรองยังค้าง
for r in range(5, 38):
    ws.row_dimensions[r].hidden = True  # เดือนนี้ไม่มีของหมดสักรายการ ทุกแถวจึงไม่ผ่านตัวกรอง
wb.save(path)

# ── เปิดไฟล์กลับมาอ่านสองแบบ ──
ws = load_workbook(path).active
rows = range(5, ws.max_row + 1)
hidden = [r for r in rows if ws.row_dimensions[r].hidden]
visible = [r for r in rows if not ws.row_dimensions[r].hidden]
print(f"ตัวกรองค้างอยู่ที่ช่วง {ws.auto_filter.ref}")
print(f"แถวข้อมูลในไฟล์ {len(rows)} แถว · ติดธงซ่อน {len(hidden)} แถว")
print(f"อ่านทุกแถว → {len(rows)} แถว · {sum(ws.cell(r, 4).value for r in rows)} ชิ้น")
print(f"อ่านเฉพาะแถวที่มองเห็น → {len(visible)} แถว · {sum(ws.cell(r, 4).value for r in visible)} ชิ้น")

ผลรันจริง

ตัวกรองค้างอยู่ที่ช่วง A4:D37
แถวข้อมูลในไฟล์ 33 แถว · ติดธงซ่อน 33 แถว
อ่านทุกแถว → 33 แถว · 121 ชิ้น
อ่านเฉพาะแถวที่มองเห็น → 0 แถว · 0 ชิ้น

ข้อมูล 121 ชิ้นอยู่ในไฟล์ครบ แต่ตัวอ่านที่เคารพธงซ่อนได้ศูนย์

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

เดือนนี้สคริปต์ไม่ได้พัง ไม่มี error · รหัสจบการทำงานเป็น 0 · รายงานออกตรงเวลา และผิด · ในระบบไม่มีอะไรรู้ว่ามันผิด จนกว่าคนอ่านรายงานจะทักมา

💡 ตัวอ่านควรทำยังไงกับแถวที่ซ่อน — ข้ามทิ้ง หรือนับรวม — ไม่มีคำตอบเดียว เพราะแถวซ่อนตัดได้สองทาง บทที่ 8 จะกลับมาที่ลำดับที่ถูก: อ่านทุกแถว ติดป้ายว่าแถวไหนซ่อน แล้วค่อยตัดสินทีละเรื่อง · ส่วนการกันไม่ให้ไฟล์ที่ได้ศูนย์แถวเดินผ่านไปถึงขั้นโหลด คือด่านตรวจของบทที่ 9

1.2 คืนที่มันพังจริง

เดือนกรกฎาคมพังแบบเงียบ แต่สคริปต์ตัวนี้เคยพังแบบเสียงดังมาก่อน และคืนนั้นสอนได้มากกว่า

ตีสองของวันที่ 1 มิถุนายน คืนปิดยอดเดือนพฤษภาคม ฐานข้อมูลปลายทางตัดการเชื่อมต่อระหว่างที่สคริปต์กำลังโหลดของสาขาที่สาม

บล็อกข้างล่างมีสองส่วน: ฟังก์ชัน monthly() คือสคริปต์ทั้งก้อน ทุกขั้นเรียงกันตามบรรทัดในฟังก์ชันเดียว · ส่วนล่างเรียกมันเป็นอีกโปรเซสแบบที่ Task Scheduler เรียก แล้วดูว่ามันทิ้งอะไรไว้ให้บ้าง — ครั้งแรกตอนฐานข้อมูลล่ม ครั้งที่สองตอนเช้าที่ฐานข้อมูลกลับมาแล้ว และคุณสั่งรันใหม่ทั้งก้อน

import sqlite3
import subprocess
import sys
import tempfile
from pathlib import Path

BRANCHES = {"ก": 120, "ข": 340, "ค": 560}  # ยอดที่แต่ละสาขานับได้เดือนพฤษภาคม (ชิ้น)


def monthly(work):
    """สคริปต์ปิดยอดรายเดือนแบบก้อนเดียว — ทุกขั้นอยู่ในฟังก์ชันนี้ เรียงกันตามบรรทัด"""
    files = {}
    for branch, total in BRANCHES.items():  # รับไฟล์ 3 สาขา (สาขาละ 33 รายการ)
        q, extra = divmod(total, 33)
        files[branch] = [q + (i < extra) for i in range(33)]
    rows = [(b, i + 1, n) for b, counts in files.items() for i, n in enumerate(counts)]  # รวม
    db = sqlite3.connect(work / "stock.db")
    db.execute("CREATE TABLE IF NOT EXISTS stock (branch, item, qty)")
    for branch in BRANCHES:  # โหลดทีละสาขา
        if branch == "ค" and not (work / "db-fixed").exists():
            raise ConnectionError("ฐานข้อมูลปลายทางตัดการเชื่อมต่อ")
        db.executemany("INSERT INTO stock VALUES (?, ?, ?)", [r for r in rows if r[0] == branch])
        db.commit()
    total = db.execute("SELECT SUM(qty) FROM stock").fetchone()[0]
    print(f"ส่งรายงาน: สต็อกรวม {total:,} ชิ้น")  # แจ้งผล


if len(sys.argv) == 2:  # ← ตัวลูก: คือสคริปต์ที่ Task Scheduler เรียกตอนตีสอง
    monthly(Path(sys.argv[1]))
    sys.exit()

# ── ตัวแม่: เรียกสคริปต์ข้างบนเป็นอีกโปรเซส แล้วดูว่าเหลืออะไรไว้ให้เราบ้าง ──
work = Path(tempfile.mkdtemp())


def run(label):
    r = subprocess.run([sys.executable, "-X", "utf8", __file__, str(work)],
                       capture_output=True, text=True, encoding="utf-8")
    print(f"── {label} ──")
    print(f"รหัสจบการทำงาน: {r.returncode}")
    print("ข้อความสุดท้าย:", (r.stdout or r.stderr).strip().splitlines()[-1])
    db = sqlite3.connect(work / "stock.db")
    for branch, n, qty in db.execute("SELECT branch, COUNT(*), SUM(qty) FROM stock GROUP BY branch ORDER BY branch"):
        print(f"  ในฐานข้อมูล สาขา {branch}: {n} แถว · {qty:,} ชิ้น")
    db.close()


run("คืนวันที่ 1 มิถุนายน 02:00")
(work / "db-fixed").touch()  # เช้ามาฐานข้อมูลกลับมาทำงานแล้ว
run("เช้าวันเดียวกัน — รันใหม่ทั้งก้อน")

ผลรันจริง

── คืนวันที่ 1 มิถุนายน 02:00 ──
รหัสจบการทำงาน: 1
ข้อความสุดท้าย: ConnectionError: ฐานข้อมูลปลายทางตัดการเชื่อมต่อ
  ในฐานข้อมูล สาขา ก: 33 แถว · 120 ชิ้น
  ในฐานข้อมูล สาขา ข: 33 แถว · 340 ชิ้น
── เช้าวันเดียวกัน — รันใหม่ทั้งก้อน ──
รหัสจบการทำงาน: 0
ข้อความสุดท้าย: ส่งรายงาน: สต็อกรวม 1,480 ชิ้น
  ในฐานข้อมูล สาขา ก: 66 แถว · 240 ชิ้น
  ในฐานข้อมูล สาขา ข: 66 แถว · 680 ชิ้น
  ในฐานข้อมูล สาขา ค: 33 แถว · 560 ชิ้น

ไล่ดูทีละบรรทัด

  • รหัสจบการทำงาน 1 กับบรรทัดสุดท้ายของ error คือทั้งหมดที่เหลือจากคืนนั้น — สิ่งที่ตัวตั้งเวลาอย่าง Task Scheduler หรือ cron ได้กลับมาจากโปรเซส มีแค่รหัสจบการทำงานกับข้อความที่มันพิมพ์ออกมา ไม่มีส่วนไหนบอกว่าข้างในทำไปถึงไหน
  • จะรู้ว่าโหลดไปแล้วสองสาขา ต้องเข้าไปเปิดฐานข้อมูลนับเอง — และต้องรู้ด้วยว่าควรนับอะไร
  • เช้ามารันใหม่ทั้งก้อน สคริปต์ทำทุกขั้นตั้งแต่ต้น รวมทั้งโหลดสาขา ก กับ ข ซ้ำ · รายงานบอก 1,480 ชิ้น ทั้งที่สามสาขานับรวมกันได้ 1,020 (120 + 340 + 560) · รอบนี้รหัสจบการทำงานเป็น 0 และรายงานผิด — อาการเดียวกับเดือนกรกฎาคม ต่างแค่ต้นเหตุ

💡 รันบล็อกในเล่มบน Windows แล้วส่งผลลงไฟล์ (python x.py > out.txt) ภาษาไทยกับเครื่องหมายบางตัวอาจทำให้ Python หยุดด้วย UnicodeEncodeError เพราะตอนส่งลงไฟล์ Python ใช้รหัสอักขระของระบบ ไม่ใช่ UTF-8 · ตั้งตัวแปร PYTHONUTF8=1 ก่อนรันแล้วหาย ส่วนการรันให้ผลขึ้นจอตรง ๆ ไม่มีปัญหานี้ · ส่วน -X utf8 ในบล็อกข้างบนคือสวิตช์เดียวกัน ใส่ให้โปรเซสลูกที่ถูกเรียกขึ้นมา

1.3 สามคำถามที่สคริปต์ตอบเองไม่ได้

สองเดือนนั้นทิ้งคำถามชุดเดียวกันไว้

คำถามเดือนพฤษภาคม (พังเสียงดัง)เดือนกรกฎาคม (พังเงียบ)
ทำถึงไหนแล้วอยู่ในตัวแปรของโปรเซสที่ตายไปแล้ว เหลือแค่ร่องรอยในฐานข้อมูลให้ไปนับเองสคริปต์คิดว่าทำครบทุกขั้น
ขั้นไหนรอขั้นไหนอยากรันแค่ "โหลดสาขา ค" แต่ข้อมูลที่แปลงแล้วอยู่ในหน่วยความจำ หายไปกับโปรเซส ต้องเริ่มใหม่ตั้งแต่รับไฟล์ขั้นโหลดเชื่อขั้นอ่านไฟล์ทุกอย่าง ไม่มีใครถามว่าศูนย์แถวสมเหตุสมผลไหม
ใครรู้ว่ามันพังรหัสจบการทำงาน 1 ในหน้าต่างที่ไม่มีใครเปิดดูตอนตีสองไม่มีใครเลย จนคนอ่านรายงานทัก

ทั้งสามคำถามมีคำตอบอยู่ครู่หนึ่ง — ขณะที่สคริปต์ยังรันอยู่ สคริปต์รู้ว่ากำลังทำขั้นไหน รู้ว่าขั้นก่อนหน้าเสร็จแล้ว รู้ว่าเพิ่งเจอ error · ปัญหาคือมันเก็บคำตอบทั้งหมดไว้ในตัวเอง: ลำดับขั้นอยู่ในลำดับบรรทัด · ความคืบหน้าอยู่ในตัวแปร · ความล้มเหลวอยู่ใน traceback ที่พิมพ์ออกจอ · พอโปรเซสจบหรือตาย คำตอบก็หายไปพร้อมกัน

ความรู้ว่างานทำถึงไหนแล้ว ต้องอยู่นอกตัวงาน — ไม่อย่างนั้นมันตายไปพร้อมกับงาน

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

1.4 ตัวคุมลำดับงานคืออะไร

💡 ตัวคุมลำดับงาน (workflow orchestrator) — โปรแกรมที่ถือคำตอบของสามคำถามนั้นไว้นอกตัวงาน: จดว่าแต่ละขั้นของแต่ละรอบอยู่ในสถานะไหน · รู้ว่าขั้นไหนต้องรอขั้นไหน · และเป็นคนสั่งรัน สั่งลองใหม่ จดบันทึก และบอกคนเมื่อมีอะไรพัง · Apache Airflow คือตัวหนึ่งในกลุ่มนี้ และเป็นตัวที่เล่มนี้ใช้ตั้งแต่บทที่ 5

สิ่งที่ตัวคุมทำไม่ได้สำคัญพอ ๆ กัน และบทนี้มีตัวอย่างให้ครบทั้งสองข้อแล้ว

  • มันไม่รู้ว่าข้อมูลผิด — เดือนกรกฎาคมทุกขั้นผ่านหมด ตัวคุมจะรายงานว่าสำเร็จเหมือนกัน · กันเรื่องนี้ต้องมีด่านตรวจก่อนโหลด ซึ่งคุณต้องเขียนเอง (บทที่ 9)
  • มันรันงานซ้ำแน่ — ลองใหม่ รันต่อจากจุดพัง ย้อนรันเดือนเก่า ล้วนเป็นการรันซ้ำ · ถ้าขั้นโหลดเขียนแบบต่อท้าย การรันซ้ำก็เบิ้ลแบบเช้าวันที่ 1 มิถุนายนทุกครั้ง · ทำให้แต่ละขั้นรันซ้ำแล้วได้ผลเท่าเดิมเป็นงานของคนเขียนขั้น ตัวคุมรับประกันให้ไม่ได้ (บทที่ 12)

1.5 คำถามที่บทนี้ทิ้งไว้

  1. รันใหม่ทั้งก้อนแล้วยอดเบิ้ล — ทำยังไงให้รันซ้ำแล้วได้ผลเท่าเดิม → บทที่ 12 ซึ่งเป็นบทพีคของเล่ม
  2. ถ้างานของคุณมีขั้นเดียว รันใหม่ทั้งก้อนได้โดยไม่เสียอะไร และไม่มีใครรอผล ยังต้องมีตัวคุมไหม → บทที่ 16
  3. ตัวอ่านควรทำยังไงกับแถวที่ซ่อน → บทที่ 8

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

  1. งานรายสัปดาห์ของคุณ: ดึงรายการสั่งซื้อจากระบบ → รวมเป็นไฟล์สรุป → ส่งอีเมลหาลูกค้า 200 คน · คืนหนึ่งมันพังตอนส่งไปได้ 80 ฉบับ ถ้าคุณรันใหม่ทั้งก้อน จะเกิดอะไรขึ้น และความรู้ข้อไหนที่หายไปกับโปรเซส
  2. ทำไมรหัสจบการทำงาน 0 ของเดือนกรกฎาคม ถึงไม่ได้แปลว่างานถูก
  3. เปิดไฟล์ Excel แล้วเห็นแค่หัวตาราง มีอะไรบนหน้าจอบ้างที่ต้องเช็กก่อนสรุปว่าไฟล์ว่าง

ลองคุมลำดับเองด้วย Python ล้วน#

บทที่แล้วจบด้วยประโยคว่าความรู้ว่างานทำถึงไหนแล้วต้องอยู่นอกตัวงาน บทนี้ลองทำตามประโยคนั้นด้วยมือ — เขียนตัวรันเล็ก ๆ ด้วย Python ล้วน แล้วดูว่าต้องเพิ่มชิ้นอะไรบ้าง กว่าจะรอดคืนแบบวันที่ 1 มิถุนายนได้

ไม่ต้องติดตั้งอะไรเลยในบทนี้ และชิ้นส่วนที่โผล่มาระหว่างทางคือชิ้นเดียวกับที่ตัวคุมลำดับงานทุกตัวมี

2.1 แยกงานเป็นขั้น แล้วจดทุกครั้งที่ขั้นหนึ่งเสร็จ

เปลี่ยนจากสคริปต์ก้อนเดียวสองจุด

  1. แยกงานเป็นขั้นที่มีชื่อ — หนึ่งขั้นหนึ่งฟังก์ชัน แล้วเก็บไว้ในรายการที่บอกลำดับ · ชื่อขั้นเป็นภาษาอังกฤษเพราะจะใช้ชื่อชุดนี้ไปจนจบเล่ม: receive_a receive_b receive_c (รับไฟล์สามสาขา) · merge_check (รวมและตรวจคอลัมน์) · transform (แปลง) · load (โหลด) · notify (ส่งรายงาน)
  2. จดลงไฟล์ทันทีที่แต่ละขั้นเสร็จ — ไฟล์ state.json เก็บว่าขั้นไหนสำเร็จแล้ว · ตอนเริ่มทุกครั้ง ตัวรันอ่านไฟล์นี้ก่อน แล้วข้ามขั้นที่สำเร็จไปแล้ว

ขั้นในตัวอย่างนี้ยังไม่ทำงานจริง (เป็นฟังก์ชันเปล่า) ยกเว้น load ที่จำลองว่าฐานข้อมูลล่มจนกว่าจะมีไฟล์ db-fixed — เพราะสิ่งที่บทนี้ดูคือตัวรัน ไม่ใช่ตัวงาน

import json
import tempfile
from pathlib import Path

WORK = Path(tempfile.mkdtemp())
STATE = WORK / "state.json"  # ← ความจำของตัวรัน อยู่บนดิสก์ ไม่ได้อยู่ในตัวแปร


def load():
    if not (WORK / "db-fixed").exists():
        raise ConnectionError("ฐานข้อมูลปลายทางตัดการเชื่อมต่อ")


STEPS = [  # ลำดับที่ต้องทำ — ขั้นละหนึ่งฟังก์ชัน
    ("receive_a", lambda: None),
    ("receive_b", lambda: None),
    ("receive_c", lambda: None),
    ("merge_check", lambda: None),
    ("transform", lambda: None),
    ("load", load),
    ("notify", lambda: None),
]


def run():
    state = json.loads(STATE.read_text()) if STATE.exists() else {}
    for name, step in STEPS:
        if state.get(name) == "success":
            print(f"  {name:<12} ข้าม — สำเร็จไปแล้ว")
            continue
        try:
            step()
        except Exception as e:
            state[name] = "failed"
            STATE.write_text(json.dumps(state))
            print(f"  {name:<12} พัง: {e} → หยุดรอบนี้")
            return
        state[name] = "success"
        STATE.write_text(json.dumps(state))  # จดทันทีที่แต่ละขั้นเสร็จ ไม่รอจบรอบ
        print(f"  {name:<12} สำเร็จ")


print("รอบที่ 1")
run()
(WORK / "db-fixed").touch()  # แก้ต้นเหตุแล้ว
print("รอบที่ 2 — เรียก run() ใหม่ ซึ่งรู้เรื่องรอบแรกจากไฟล์อย่างเดียว")
run()

ผลรันจริง

รอบที่ 1
  receive_a    สำเร็จ
  receive_b    สำเร็จ
  receive_c    สำเร็จ
  merge_check  สำเร็จ
  transform    สำเร็จ
  load         พัง: ฐานข้อมูลปลายทางตัดการเชื่อมต่อ → หยุดรอบนี้
รอบที่ 2 — เรียก run() ใหม่ ซึ่งรู้เรื่องรอบแรกจากไฟล์อย่างเดียว
  receive_a    ข้าม — สำเร็จไปแล้ว
  receive_b    ข้าม — สำเร็จไปแล้ว
  receive_c    ข้าม — สำเร็จไปแล้ว
  merge_check  ข้าม — สำเร็จไปแล้ว
  transform    ข้าม — สำเร็จไปแล้ว
  load         สำเร็จ
  notify       สำเร็จ

รอบที่สองข้ามห้าขั้นแรก แล้วเริ่มที่ load ซึ่งเป็นขั้นที่พัง — ไม่มีอะไรถูกทำซ้ำ · run() ในรอบที่สองไม่ได้รับตัวแปรอะไรจากรอบแรกเลยนอกจากไฟล์ จะปิดโปรแกรมแล้วเปิดใหม่ก็ได้ผลเดียวกัน (บทที่ 4 จะฆ่ามันทิ้งกลางงานให้ดูจริง ๆ)

สังเกตว่าตัวรันจด state.json ทันทีที่แต่ละขั้นเสร็จ ไม่ได้รอจดทีเดียวตอนจบรอบ — ถ้ารอจดตอนจบ คืนที่พังกลางทางจะไม่มีอะไรถูกจดเลย แล้วรอบที่สองก็ต้องเริ่มใหม่ทั้งหมดเหมือนเดิม

⚠️ ตัวรันนี้โกงอยู่ข้อหนึ่ง: ขั้นในตัวอย่างไม่ได้ส่งข้อมูลให้กัน · ของจริง load ต้องใช้ผลของ transform — ถ้าผลนั้นอยู่ในตัวแปร รอบที่สองจะไม่มีอะไรให้โหลด ⇒ ผลของแต่ละขั้นต้องถูกเขียนไว้นอกโปรเซสด้วย ไม่ใช่แค่สถานะ · Airflow มีกลไกส่งค่าเล็ก ๆ ระหว่างขั้นชื่อ XCom ซึ่งจะเจอในบทที่ 6 และบทที่ 8 จะบอกว่าทำไมห้ามส่งของก้อนใหญ่ผ่านมัน

2.2 พังแล้วลองใหม่ และจดว่าลองไปกี่ครั้ง

ความพังบางแบบหายเองถ้ารอสักพัก — ไฟล์ของสาขายังอัปโหลดไม่เสร็จ ฐานข้อมูลรีสตาร์ตอยู่ · ถ้าตัวรันลองใหม่เองหลังรอเวลาหนึ่ง คืนแบบนั้นจะไม่ต้องมีใครตื่นมาสั่ง

ตัวรันข้างล่างเพิ่มสามอย่างจากรุ่นแรก: จำนวนครั้งที่ยอมลองใหม่ (RETRIES) · เวลารอก่อนลองใหม่ (RETRY_DELAY) · และบันทึก (runner.log) ที่จดทุกครั้งที่ลอง ไม่ใช่แค่ผลสุดท้าย · ส่วน state.json ตอนนี้เก็บจำนวนครั้งที่ลองของแต่ละขั้นด้วย

import json
import tempfile
import time
from pathlib import Path

WORK = Path(tempfile.mkdtemp())
STATE = WORK / "state.json"  # สถานะ + จำนวนครั้งที่ลอง ของทุกขั้น
LOG = WORK / "runner.log"  # ← ชิ้นใหม่: บันทึกว่าเกิดอะไรขึ้นทีละครั้ง
RETRIES = 2  # ← ชิ้นใหม่: พังแล้วลองใหม่ได้อีกกี่ครั้ง
RETRY_DELAY = 0.2  # วินาที · ของจริงตั้งเป็นนาที


def receive_c():
    arrived = WORK / "stock_2026-05_c.xlsx"
    if not arrived.exists():
        arrived.touch()  # จำลองว่าสาขา ค อัปโหลดเสร็จระหว่างที่ตัวรันรอ
        raise FileNotFoundError("ไฟล์สาขา ค ยังมาไม่ถึง")


STEPS = [
    ("receive_a", lambda: None),
    ("receive_b", lambda: None),
    ("receive_c", receive_c),
    ("merge_check", lambda: None),
    ("transform", lambda: None),
    ("load", lambda: None),
    ("notify", lambda: None),
]


def log(line):
    with LOG.open("a", encoding="utf-8") as f:
        f.write(line + "\n")


def run():
    state = json.loads(STATE.read_text()) if STATE.exists() else {}
    for name, step in STEPS:
        cell = state.setdefault(name, {"state": "none", "tries": 0})
        while cell["state"] != "success":
            cell["tries"] += 1
            try:
                step()
                cell["state"] = "success"
                log(f"{name} ครั้งที่ {cell['tries']}: สำเร็จ")
            except Exception as e:
                cell["state"] = "failed"
                log(f"{name} ครั้งที่ {cell['tries']}: พัง — {e}")
                if cell["tries"] > RETRIES:
                    STATE.write_text(json.dumps(state))
                    return
                time.sleep(RETRY_DELAY)
            STATE.write_text(json.dumps(state))


run()
print(LOG.read_text(encoding="utf-8"), end="")
print("state.json → receive_c =", json.loads(STATE.read_text())["receive_c"])

ผลรันจริง

receive_a ครั้งที่ 1: สำเร็จ
receive_b ครั้งที่ 1: สำเร็จ
receive_c ครั้งที่ 1: พัง — ไฟล์สาขา ค ยังมาไม่ถึง
receive_c ครั้งที่ 2: สำเร็จ
merge_check ครั้งที่ 1: สำเร็จ
transform ครั้งที่ 1: สำเร็จ
load ครั้งที่ 1: สำเร็จ
notify ครั้งที่ 1: สำเร็จ
state.json → receive_c = {'state': 'success', 'tries': 2}

receive_c พังครั้งแรกเพราะไฟล์ยังมาไม่ถึง ตัวรันรอแล้วลองครั้งที่ 2 ซึ่งผ่าน · บันทึกเก็บทั้งสองครั้ง และ state.json จำว่าขั้นนี้ใช้ไปสองครั้ง — ตัวเลขนี้จะสำคัญมากเมื่อถึงบทที่ 11

บันทึกในตัวอย่างไม่มีเวลากำกับ เพราะผลรันในเล่มต้องรันซ้ำแล้วได้เหมือนเดิมทุกตัวอักษร · บันทึกของจริงต้องมีเวลาเสมอ

⚠️ ลองใหม่ช่วยได้เฉพาะความพังแบบชั่วคราว · ถ้าไฟล์สาขา ค ไม่มีคอลัมน์จำนวนเลย ลองอีกกี่ครั้งก็พังเหมือนเดิม แถมเสียเวลารอทุกครั้ง · การแยกความพังสองแบบนี้ออกจากกันคือเรื่องของบทที่ 10

2.3 ชิ้นส่วนที่โผล่มาเอง

ไล่ดูว่าตัวรันสองรุ่นนี้มีอะไรบ้าง แล้วจับคู่กับคำถามสามข้อของบทที่ 1

ชิ้นตอบคำถามข้อไหนในตัวรันของเรา
ไฟล์จำสถานะทำถึงไหนแล้วstate.json จดทุกครั้งที่ขั้นหนึ่งเปลี่ยนสถานะ
ลำดับขั้นขั้นไหนรอขั้นไหน (ครึ่งเดียว — หัวข้อ 2.4)รายการ STEPS
วงลองใหม่ไม่ต้องตื่นมาสั่งเองเมื่อพังชั่วคราวRETRIES กับ RETRY_DELAY
บันทึกใครรู้ว่ามันพัง และพังยังไงrunner.log
ตารางเวลาใครเป็นคนเริ่มรอบยังไม่มี — Task Scheduler เรียก python runner.py แทน

ห้าชิ้นนี้ไม่ได้มาจากการออกแบบล่วงหน้า มันโผล่มาเองทีละชิ้นเมื่อพยายามตอบคำถามสามข้อนั้น — นี่คือเหตุผลที่ตัวคุมลำดับงานทุกตัวหน้าตาคล้ายกัน · ของที่ยังขาดชัด ๆ คือการบอกคน: ตอนนี้ถ้าลองครบแล้วยังพัง ตัวรันแค่หยุดเงียบ ๆ แล้วรอให้มีคนมาเปิดบันทึกเอง

2.4 สิ่งที่ตัวรันนี้ยังทำไม่ได้

  1. มันเดินตามรายการทีละขั้น — รับไฟล์สามสาขาไม่ได้รอกันเลย แต่ต้องต่อคิวกัน · และถ้า receive_b พังถาวร receive_c จะไม่ถูกลองเลยสักครั้ง ทั้งที่ไม่เกี่ยวกัน → บทที่ 3
  2. มันมีรอบเดียว — state.json ไม่รู้ว่าเป็นของเดือนไหน · ถ้าเอาตัวรันนี้ไปปิดยอดเดือนมิถุนายนโดยไม่ลบไฟล์ มันจะเห็นว่าทุกขั้นสำเร็จแล้ว แล้วข้ามหมด เดือนมิถุนายนไม่ถูกทำเลยทั้งที่ทุกอย่างดูเรียบร้อย → บทที่ 4
  3. ถ้าตัวรันเองตายระหว่างขั้น — state.json จดแค่ขั้นที่เสร็จแล้ว ไม่ได้จดว่ากำลังทำอะไรอยู่ → บทที่ 4

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

  1. เพื่อนส่งตัวรันมาให้ดู มันจด state.json ครั้งเดียวตอนจบรอบ · ถ้าพังที่ขั้นที่ 6 แล้วรันใหม่ จะเกิดอะไรขึ้น และขาดชิ้นไหน
  2. ถ้าใช้ตัวรันในบทนี้ปิดยอดเดือนมิถุนายนต่อจากเดือนพฤษภาคม โดยไม่ลบ state.json จะเกิดอะไรขึ้น
  3. ลองเขียนเอง: เพิ่มให้ตัวรันในหัวข้อ 2.2 เขียนไฟล์ ALERT.txt บอกชื่อขั้นกับข้อความ error ทันทีที่ขั้นไหนลองครบแล้วยังพัง

ลำดับงานคือกราฟ ไม่ใช่รายการ#

ก่อนอ่านต่อ ลองเดาก่อน: คืนหนึ่งไฟล์สาขา ข เสียจนอ่านไม่ได้ ตัวรันของบทที่ 2 จะทำอะไรกับ receive_c — และควรทำอะไร · จดคำตอบไว้

ตัวรันหยุดที่ receive_b และไม่แตะ receive_c เลย ทั้งที่การรับไฟล์สาขา ค ไม่เกี่ยวอะไรกับไฟล์สาขา ข สักนิด · ถ้าคุณตอบว่า "ก็ต้องหยุดหมด เพราะยังไงก็ต้องรอครบสามสาขาก่อนรวม" — ครึ่งหลังถูก ครึ่งแรกผิด: ขั้นรวมต้องรอครบสามสาขาจริง แต่ receive_c ไม่ได้รออะไรเลย ปล่อยมันรับไฟล์ไว้ก่อนได้ พอแก้ไฟล์สาขา ข เสร็จ จะได้เหลือรันแค่สาขาเดียว

3.1 รายการบอกมากกว่าที่เป็นจริง

รายการเจ็ดขั้นของบทที่ 2 บอกความจริงสองอย่างปนกัน — อะไรต้องเสร็จก่อนอะไร ซึ่งจริง กับลำดับที่บังเอิญเขียน ซึ่งไม่จริง · receive_a อยู่บรรทัดบน receive_b เพราะเขียนก่อน ไม่ใช่เพราะต้องเสร็จก่อน · พอทุกอย่างอยู่ในรายการเดียว ตัวรันแยกสองอย่างนี้ไม่ออก มันจึงเชื่อทุกบรรทัดว่าเป็นข้อบังคับ

ทางแก้คือเลิกเขียน "ทำอะไรก่อน" แล้วเขียนแค่ "อะไรรออะไร" ให้ครบ ที่เหลือให้เครื่องคิดเอง

3.2 เขียนแค่ว่าอะไรรออะไร แล้วให้เครื่องหาลำดับเอง

💡 กราฟ (graph) — จุดกับเส้นที่เชื่อมจุด · ในเล่มนี้จุดคือขั้น และเส้นชี้ว่าขั้นไหนต้องเสร็จก่อนขั้นไหน · ขั้นที่ต้องเสร็จก่อนเรียกว่า ต้นน้ำ (upstream) ขั้นที่รออยู่เรียกว่า ปลายน้ำ (downstream)

Python มีเครื่องมือหาลำดับจากกราฟมาในตัว ชื่อ graphlib (ไลบรารีมาตรฐานตั้งแต่ 3.9 ไม่ต้องติดตั้ง) · บล็อกข้างล่างเขียนท่อเดิมเป็น dict ที่บอกแค่ว่าแต่ละขั้นรอใคร แล้วถาม graphlib ซ้ำ ๆ ว่า "ตอนนี้ขั้นไหนพร้อม"

from graphlib import TopologicalSorter

# ขั้น: ต้นน้ำที่ต้องสำเร็จก่อน — เขียนว่า "อะไรรออะไร" ไม่ได้เขียนว่า "ทำอะไรก่อน"
NEEDS = {
    "receive_a": [],
    "receive_b": [],
    "receive_c": [],
    "merge_check": ["receive_a", "receive_b", "receive_c"],
    "transform": ["merge_check"],
    "load": ["transform"],
    "notify": ["load"],
}

graph = TopologicalSorter(NEEDS)
graph.prepare()
wave = 1
while graph.is_active():
    ready = sorted(graph.get_ready())  # เรียงก่อนพิมพ์ — ในจังหวะเดียวกันไม่มีใครมาก่อนใคร
    print(f"จังหวะที่ {wave}: {', '.join(ready)}")
    graph.done(*ready)  # สมมติว่าทุกขั้นในจังหวะนี้สำเร็จ
    wave += 1

ผลรันจริง

จังหวะที่ 1: receive_a, receive_b, receive_c
จังหวะที่ 2: merge_check
จังหวะที่ 3: transform
จังหวะที่ 4: load
จังหวะที่ 5: notify

จังหวะที่ 1 มีสามขั้นพร้อมกัน — นี่คือข้อมูลที่รายการไม่มีวันบอกได้ · สามขั้นนี้รันพร้อมกันได้ เจ็ดขั้นจึงจบในห้าจังหวะแทนที่จะเป็นเจ็ด

ขั้นที่พร้อมรัน คือขั้นที่ต้นน้ำสำเร็จครบแล้ว — ไม่ใช่ขั้นถัดไปในรายการ

รายการเก็บลำดับที่เขียน1receive_a2receive_b3receive_c4merge_check5transform6load7notifyกราฟเก็บแค่ว่าอะไรรออะไรจังหวะ 1receive_areceive_breceive_cจังหวะ 2merge_checkจังหวะ 3transformจังหวะ 4loadจังหวะ 5notify
FIG 3.1 — ขั้นเจ็ดขั้นเดียวกัน วาดสองแบบ · รายการเก็บลำดับที่บังเอิญเขียน จึงทำได้ทีละขั้นเจ็ดจังหวะ และบอกไม่ได้ว่า receive ทั้งสามไม่ได้รอกัน · กราฟเก็บแค่ “อะไรรออะไร” แล้วลำดับที่ทำได้จริงก็โผล่มาเอง: ห้าจังหวะ จังหวะแรกทำสามขั้นพร้อมกัน — ตรงกับผลรันของ graphlib ในหัวข้อ 3.2

sorted() ในบล็อกมีเหตุผล · เอกสารของ graphlib บอกว่าขั้นที่อยู่ระดับเดียวกันออกมาตามลำดับที่ใส่เข้ากราฟ — ซึ่งก็คือลำดับที่บังเอิญเขียน dict อีกนั่นเอง ไม่ใช่ข้อบังคับของงาน · เรียงตามชื่อก่อนพิมพ์ทำให้ผลไม่ขึ้นกับลำดับที่เขียน และย้ำความจริงของกราฟ: สามขั้นในจังหวะเดียวกันไม่มีใครมาก่อนใคร

3.3 พังหนึ่งขั้น เสียแค่ปลายน้ำของมัน

ทีนี้ลองให้ receive_b พังจริง แล้วเดินท่อสองแบบเทียบกัน — ตามรายการแบบบทที่ 2 กับตามกราฟ · ทุกอย่างเหมือนกันหมด ต่างแค่วิธีเดิน

from graphlib import TopologicalSorter

NEEDS = {
    "receive_a": [],
    "receive_b": [],
    "receive_c": [],
    "merge_check": ["receive_a", "receive_b", "receive_c"],
    "transform": ["merge_check"],
    "load": ["transform"],
    "notify": ["load"],
}
BROKEN = {"receive_b"}  # ไฟล์สาขา ข เสีย อ่านไม่ได้


def by_list():
    """แบบบทที่ 2: เดินตามรายการ เจอขั้นพังแล้วหยุด"""
    ran = []
    for step in NEEDS:  # dict จำลำดับที่เขียน = รายการเดิม
        if step in BROKEN:
            return ran, [step]
        ran.append(step)
    return ran, []


def by_graph():
    """เดินตามกราฟ: ขั้นที่พร้อมคือขั้นที่ต้นน้ำสำเร็จครบ — ขั้นที่พังจะไม่ถูกบอกว่า done"""
    graph = TopologicalSorter(NEEDS)
    graph.prepare()
    ran, failed = [], []
    while ready := sorted(graph.get_ready()):
        for step in ready:
            if step in BROKEN:
                failed.append(step)
            else:
                ran.append(step)
                graph.done(step)
    return ran, failed


for name, walk in [("ตามรายการ", by_list), ("ตามกราฟ", by_graph)]:
    ran, failed = walk()
    untouched = [s for s in NEEDS if s not in ran + failed]
    print(f"{name}")
    print(f"  สำเร็จ: {', '.join(ran) or '-'}")
    print(f"  พัง: {', '.join(failed)}")
    print(f"  ไม่ได้รัน: {', '.join(untouched)}")

ผลรันจริง

ตามรายการ
  สำเร็จ: receive_a
  พัง: receive_b
  ไม่ได้รัน: receive_c, merge_check, transform, load, notify
ตามกราฟ
  สำเร็จ: receive_a, receive_c
  พัง: receive_b
  ไม่ได้รัน: merge_check, transform, load, notify

ตามรายการ ความพังของ receive_b ลามไปถึง receive_c ซึ่งไม่เกี่ยวกัน · ตามกราฟ receive_c รับไฟล์เสร็จไปแล้ว ส่วนที่ไม่ได้รันมีแค่ขั้นที่รอ receive_b อยู่จริง — ความพังไหลตามเส้นของกราฟลงไปหาปลายน้ำเท่านั้น ไม่ไหลข้างไปหาขั้นที่ไม่เกี่ยวกัน

สี่ขั้นที่ไม่ได้รันในแบบกราฟไม่ได้พัง มันรอต้นน้ำที่ไม่มีวันสำเร็จ — สถานะนี้มีชื่อของมันเอง และบทที่ 4 จะตั้งชื่อให้

3.4 กราฟต้องไม่มีวง

สมมติวันหนึ่งมีคนขอ: "ก่อนรวมไฟล์ ช่วยส่งรายงานเบื้องต้นให้หัวหน้าสาขาดูก่อน" แล้วคนแก้ทำแบบง่ายที่สุด คือให้ merge_check รอ notify ด้วย — ทั้งที่ notify ก็รอ load ซึ่งรอ transform ซึ่งรอ merge_check อยู่แล้ว

from graphlib import CycleError, TopologicalSorter

NEEDS = {
    "receive_a": [],
    "receive_b": [],
    "receive_c": [],
    "merge_check": ["receive_a", "receive_b", "receive_c"],
    "transform": ["merge_check"],
    "load": ["transform"],
    "notify": ["load"],
}
# มีคนขอ: "ก่อนรวมไฟล์ ให้ส่งรายงานเบื้องต้นไปให้หัวหน้าสาขาดูก่อน"
# แล้วแก้ด้วยการให้ merge_check รอ notify — ทั้งที่ notify ก็รอ load ซึ่งรอ merge_check อยู่แล้ว
NEEDS["merge_check"] = NEEDS["merge_check"] + ["notify"]

try:
    TopologicalSorter(NEEDS).prepare()
except CycleError as e:
    cycle = e.args[1]
    print("หาจุดเริ่มไม่ได้ — มีวง:")
    print("  " + " → ".join(cycle))

ผลรันจริง

หาจุดเริ่มไม่ได้ — มีวง:
  merge_check → transform → load → notify → merge_check

ทุกขั้นในวงรอกันเองเป็นทอด ๆ จนวนกลับมาที่ตัวเอง ไม่มีขั้นไหนในวงเริ่มได้ก่อน · กราฟที่มีวงจึงไม่มีลำดับที่ทำครบทุกขั้นได้เลย ไม่ใช่แค่หายาก · prepare() ฟ้องตั้งแต่ตอนประกาศกราฟ ก่อนจะมีขั้นไหนได้รัน — receive ทั้งสามไม่อยู่ในวงก็จริง แต่ถ้ารันพวกมันไปก่อน คุณจะไปรู้ตัวเอาตอนท่อค้างอยู่กลางทาง

คำขอนั้นแก้ได้ถูกด้วยการเพิ่มขั้นใหม่ เช่น notify_preview ที่รอแค่ receive ทั้งสาม · "ส่งรายงานสองครั้ง" คือสองขั้น ไม่ใช่ขั้นเดียวที่ถูกเรียกสองหน

💡 DAG (directed acyclic graph) — กราฟที่เส้นมีทิศ (ชี้จากต้นน้ำไปปลายน้ำ) และไม่มีวง · Airflow เรียกท่อแต่ละท่อว่า DAG ด้วยเหตุผลนี้ตรงตัว — ไฟล์แรกที่คุณจะเขียนในบทที่ 6 คือการประกาศกราฟแบบ NEEDS ข้างบนนี้ เพียงแต่เขียนด้วยไวยากรณ์ของ Airflow

นี่คือความจริงข้อที่สองของเล่ม: ลำดับที่แท้ของงานคือ "อะไรต้องเสร็จก่อนอะไร" ไม่ใช่บรรทัดที่เรียงกัน และกราฟนั้นต้องไม่มีวง

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

  1. งานรายงานประจำสัปดาห์มีเจ็ดขั้น: ดึงรายการสั่งซื้อ · ดึงรายชื่อสินค้า · ดึงรายชื่อสาขา · จับคู่รายการสั่งซื้อกับสินค้า · จับคู่ผลนั้นกับสาขา · สรุปรายสาขา · ส่งรายงาน — เขียน NEEDS ให้มัน แล้วบอกว่าจังหวะแรกมีขั้นไหนบ้าง และทั้งท่อจบในกี่จังหวะ
  2. ใน NEEDS ของบทนี้ ถ้า notify พัง มีขั้นไหนบ้างที่ไม่ได้รันเพราะมัน
  3. ทำไมการเจอวงตั้งแต่ตอนประกาศกราฟ ถึงดีกว่าปล่อยให้ท่อรันขั้นที่ไม่อยู่ในวงไปก่อนแล้วค่อยรู้ทีหลัง

แบบจำลองความ⁠คิด: ตารางสถานะ + วงรอบที่ถามว่าใครเดินได้#

บทนี้สำคัญที่สุดในเล่ม · ทุกคำสั่งของ Airflow ที่จะเจอตั้งแต่บทที่ 5 อธิบายได้เป็น "ตารางก่อน → ตารางหลัง" บนตารางในบทนี้ ถ้าอ่านได้บทเดียว อ่านบทนี้

ลองเดาก่อน: ตัวรันกำลังทำ transform อยู่ แล้วเครื่องดับ · เปิดขึ้นมาใหม่ ตัวรันควรทำอะไรกับ transform — ทำต่อ ข้ามไป หรือเริ่มทั้งรอบใหม่ · จดคำตอบไว้ เฉลยอยู่หัวข้อ 4.4

4.1 ประโยคแก่น

ตัวคุมลำดับงานคือตารางสถานะที่อยู่นอกตัวงาน กับวงรอบหนึ่งวงที่อ่านตารางแล้วถามว่าช่องไหนเดินได้แล้ว

ประโยคนี้มีสองครึ่ง และทุกอย่างในสามบทที่ผ่านมาตกอยู่ในครึ่งใดครึ่งหนึ่ง

  • ตาราง คือความจำ — ตอบว่าทำถึงไหนแล้ว และพังหรือยัง · state.json ของบทที่ 2 คือตารางที่มีคอลัมน์เดียว
  • วงรอบ คือการตัดสินใจ — ใช้กราฟของบทที่ 3 ตอบว่าช่องไหนรอใคร แล้วเป็นคนลงมือ ลองใหม่ และจดผลกลับลงตาราง

4.2 ตาราง: แถวคือขั้น คอลัมน์คือรอบ

💡 รอบ (run) — การทำงานทั้งท่อหนึ่งครั้งสำหรับข้อมูลชุดหนึ่ง เช่นรอบปิดยอดเดือนกรกฎาคม · ช่อง — ขั้นหนึ่งของรอบหนึ่ง เช่น transform ของรอบเดือนกรกฎาคม (Airflow เรียกช่องนี้ว่า task instance)

ตารางของบทนี้วางขั้นเป็นแถว รอบเป็นคอลัมน์ และแต่ละช่องเก็บสองอย่าง: สถานะ กับจำนวนครั้งที่ลอง · ปัญหาข้อ 2 ของหัวข้อ 2.4 หายไปเอง — เดือนมิถุนายนคือคอลัมน์ใหม่ ไม่ได้ไปทับสถานะของเดือนพฤษภาคม

ชื่อสถานะใช้ชื่อจริงของ Airflow ทุกตัว ไม่แปล เพราะคุณจะเจอคำเดียวกันนี้บนหน้าจอของมันตั้งแต่บทที่ 7

สถานะแปลว่าเปลี่ยนต่อเป็นอะไรได้
noneยังไม่มีสถานะ — ต้นน้ำยังไม่ครบrunning · upstream_failed
runningกำลังทำอยู่success · failed · up_for_retry
successทำเสร็จโดยไม่มี errorจบ
failedพัง และไม่เหลือสิทธิ์ลองใหม่จบ
up_for_retryพัง แต่ยังมีสิทธิ์ลองใหม่ — รอเวลาแล้วจะถูกรันอีกrunning
upstream_failedไม่ได้รัน เพราะต้นน้ำพังจบ

ของจริงมีสถานะ 14 ตัว (หน้า Tasks ในเอกสารทางการ) เล่มนี้ใช้หกตัวนี้ก่อน ตัวที่เหลือมาเมื่อมีเรื่องให้ใช้ — skipped ในบทที่ 9 · ส่วน scheduled กับ queued เป็นสถานะผ่านทางระหว่าง none กับ running ซึ่งจะเห็นในบทที่ 15 ว่าใครเป็นคนเปลี่ยนมัน

4.3 วงรอบ: หนึ่งจังหวะ ถามทุกช่องคำถามเดียว

วงรอบทำงานเป็นจังหวะ ทุกจังหวะอ่านตารางหนึ่งครั้ง แล้วไล่ถามทุกช่องด้วยกติกาสามข้อ

  1. ช่องที่เป็น none หรือ up_for_retry และต้นน้ำ success ครบ → ลงมือทำ · จด running ลงตารางก่อนลงมือ แล้วค่อยทำ · ผลออกมาเป็น success · หรือ up_for_retry ถ้ายังมีสิทธิ์ลองใหม่ · หรือ failed ถ้าหมดสิทธิ์
  2. ช่องที่ยังไม่ได้รัน แต่มีต้นน้ำ failed หรือ upstream_failed → upstream_failed ทันที โดยไม่ได้รัน
  3. ช่องอื่นทั้งหมด → ไม่แตะ รอจังหวะหน้า

และอีกหนึ่งคำสั่งที่อยู่นอกวงรอบ: เริ่มรอบใหม่ = เพิ่มหนึ่งคอลัมน์ ทุกช่องเป็น none

บล็อกข้างล่างคือแบบจำลองทั้งตัว ตารางอยู่ใน table.json บนดิสก์ วงรอบอ่านและเขียนมันทุกจังหวะ · รอบเดือนกรกฎาคม receive_c พังชั่วคราวครั้งแรก · ก่อนจังหวะที่ 3 มีการเพิ่มคอลัมน์เดือนสิงหาคม ซึ่ง transform พังถาวร · ทุกขั้นได้สิทธิ์ลองใหม่ 1 ครั้ง (RETRIES = 1)

import json
import tempfile
from pathlib import Path

NEEDS = {
    "receive_a": [],
    "receive_b": [],
    "receive_c": [],
    "merge_check": ["receive_a", "receive_b", "receive_c"],
    "transform": ["merge_check"],
    "load": ["transform"],
    "notify": ["load"],
}
RETRIES = 1
TABLE = Path(tempfile.mkdtemp()) / "table.json"  # ← ตารางอยู่บนดิสก์ วงรอบอ่านและเขียนมันทุกจังหวะ

# งานจำลอง: รอบไหน ขั้นไหน พังแบบไหน
TRANSIENT = {("2026-07", "receive_c")}  # ไฟล์มาช้า ลองใหม่แล้วผ่าน
PERMANENT = {("2026-08", "transform")}  # คอลัมน์หาย ลองกี่ครั้งก็พัง


def execute(run, task, attempt):
    if (run, task) in PERMANENT or ((run, task) in TRANSIENT and attempt == 1):
        raise RuntimeError(f"{task} พัง")


def read():
    return json.loads(TABLE.read_text()) if TABLE.exists() else {}


def write(table):
    TABLE.write_text(json.dumps(table))


def trigger(run):
    """เริ่มรอบใหม่ = เพิ่มหนึ่งคอลัมน์ ทุกช่องเป็น none"""
    table = read()
    table[run] = {task: {"state": "none", "try": 0} for task in NEEDS}
    write(table)


def tick():
    """หนึ่งรอบของวงรอบ: อ่านตาราง → หาช่องที่เดินได้ → ลงมือ → เขียนผลกลับ"""
    table, changes = read(), []
    for run, column in table.items():
        seen = {task: cell["state"] for task, cell in column.items()}  # ทั้งจังหวะตัดสินจากภาพเดียวกัน
        for task, ups in NEEDS.items():
            cell = column[task]
            if seen[task] not in ("none", "up_for_retry"):
                continue
            if any(seen[u] in ("failed", "upstream_failed") for u in ups):
                cell["state"] = "upstream_failed"
            elif all(seen[u] == "success" for u in ups):
                cell["try"] += 1
                cell["state"] = "running"
                write(table)  # จดก่อนลงมือ — ถ้าตายกลางงาน ตารางยังบอกได้ว่าค้างที่ช่องไหน
                try:
                    execute(run, task, cell["try"])
                    cell["state"] = "success"
                except Exception:
                    cell["state"] = "up_for_retry" if cell["try"] <= RETRIES else "failed"
            else:
                continue  # ต้นน้ำยังไม่จบ — รอจังหวะหน้า
            write(table)
            changes.append(f"{run} {task:<12} {seen[task]} → {cell['state']}")
    return changes


trigger("2026-07")
n = 0
while True:
    n += 1
    if n == 3:
        trigger("2026-08")
        print("+ เพิ่มคอลัมน์ 2026-08")
    changes = tick()
    if not changes:
        break
    print(f"จังหวะ {n}")
    for line in changes:
        print("  " + line)

table = read()
print()
print((" " * 16 + "".join(f"{run:<19}" for run in table)).rstrip())
for task in NEEDS:
    cells = [f"{table[run][task]['state']} {table[run][task]['try']}" for run in table]
    print((f"{task:<16}" + "".join(f"{c:<19}" for c in cells)).rstrip())
for run, column in table.items():
    leaf = column["notify"]["state"]  # ขั้นปลายสุด — ไม่มีใครรอมัน
    print(f"สถานะของรอบ {run}: {'success' if leaf == 'success' else 'failed'} (ดูจาก notify = {leaf})")

ผลรันจริง

จังหวะ 1
  2026-07 receive_a    none → success
  2026-07 receive_b    none → success
  2026-07 receive_c    none → up_for_retry
จังหวะ 2
  2026-07 receive_c    up_for_retry → success
+ เพิ่มคอลัมน์ 2026-08
จังหวะ 3
  2026-07 merge_check  none → success
  2026-08 receive_a    none → success
  2026-08 receive_b    none → success
  2026-08 receive_c    none → success
จังหวะ 4
  2026-07 transform    none → success
  2026-08 merge_check  none → success
จังหวะ 5
  2026-07 load         none → success
  2026-08 transform    none → up_for_retry
จังหวะ 6
  2026-07 notify       none → success
  2026-08 transform    up_for_retry → failed
จังหวะ 7
  2026-08 load         none → upstream_failed
จังหวะ 8
  2026-08 notify       none → upstream_failed

                2026-07            2026-08
receive_a       success 1          success 1
receive_b       success 1          success 1
receive_c       success 2          success 1
merge_check     success 1          success 1
transform       success 1          failed 2
load            success 1          upstream_failed 0
notify          success 1          upstream_failed 0
สถานะของรอบ 2026-07: success (ดูจาก notify = success)
สถานะของรอบ 2026-08: failed (ดูจาก notify = upstream_failed)

ไล่ดูทีละจุด

  • จังหวะที่ 3 สองคอลัมน์เดินพร้อมกัน — คอลัมน์สิงหาคมเริ่มรับไฟล์ระหว่างที่กรกฎาคมยังไม่จบ วงรอบไม่สนว่าช่องอยู่คอลัมน์ไหน ถามคำถามเดียวกันกับทุกช่อง
  • receive_c ของกรกฎาคม none → up_for_retry → success จบที่ครั้งที่ 2 · ส่วน merge_check รอจนจังหวะที่ 3 เพราะจังหวะที่ 2 ต้นน้ำยังไม่ครบ
  • transform ของสิงหาคม พังครั้งแรกได้ up_for_retry พังครั้งที่สองได้ failed — สิทธิ์ลองใหม่ 1 ครั้ง บวกครั้งแรก เท่ากับลองได้ 2 ครั้งพอดี
  • load กับ notify ของสิงหาคม เป็น upstream_failed ครั้งที่ 0 — ไม่เคยรันเลยสักครั้ง ความพังไหลลงปลายน้ำทีละจังหวะ แบบเดียวกับบทที่ 3
ก่อนจังหวะที่ 6หลังจังหวะที่ 62026-072026-082026-072026-08receive_areceive_breceive_cmerge_checktransformloadnotifysuccess 1success 1success 2success 1success 1success 1none 0success 1success 1success 1success 1up_for_retry 1none 0none 0วงรอบsuccess 1success 1success 2success 1success 1success 1success 1success 1success 1success 1success 1failed 2none 0none 0สถานะของรอบยังไม่จบยังไม่จบsuccessยังไม่จบnonerunningsuccessfailedup_for_retryupstream_failed
FIG 4.1 — หนึ่งจังหวะของวงรอบ คัดจากผลรันในหัวข้อ 4.3 · วงรอบถามทุกช่องคำถามเดียว แล้วหยิบได้สองช่อง (กรอบหนา): notify ของกรกฎาคมที่ต้นน้ำสำเร็จครบ กับ transform ของสิงหาคมที่ยังมีสิทธิ์ลองใหม่ · ลองแล้วพังซ้ำ สิทธิ์หมด จึงเป็น failed ครั้งที่ 2 · ช่อง load กับ notify ของสิงหาคมยังเป็น none — จังหวะหน้าจะกลายเป็น upstream_failed โดยไม่ได้รัน · สถานะของรอบดูจากขั้นปลายสุด: notify ของกรกฎาคมสำเร็จ รอบนั้นจึงจบเป็น success

ภาพนี้คือจังหวะที่ 6 ในผลรัน · วงรอบหยิบได้สองช่องในจังหวะเดียว: notify ของกรกฎาคมที่ต้นน้ำสำเร็จครบ กับ transform ของสิงหาคมที่ยังมีสิทธิ์ลองใหม่ · ช่องอื่นไม่ถูกแตะ — load ของสิงหาคมยังเป็น none ในจังหวะนี้ เพราะตอนต้นจังหวะ transform ยังไม่ได้ failed · ทุกช่องในจังหวะเดียวกันตัดสินจากภาพตารางตอนต้นจังหวะ ไม่ได้ดูผลของช่องที่เพิ่งเปลี่ยนในจังหวะเดียวกัน

สถานะของรอบ — รอบจะมีสถานะของตัวเองก็ต่อเมื่อทุกช่องในคอลัมน์จบแล้ว และตัดสินจากขั้นปลายสุดเท่านั้น คือขั้นที่ไม่มีใครรอมัน ในท่อนี้คือ notify · ขั้นปลายสุดทุกขั้นสำเร็จ (หรือ skipped) = รอบ success · มีขั้นปลายสุดตัวไหน failed หรือ upstream_failed = รอบ failed · นี่คือกติกาจริงของ Airflow (หน้า DAG Runs ในเอกสารทางการ) — ฟังดูไม่มีอะไร แต่บทที่ 7 กับ 11 จะโชว์ว่ามันทำให้รอบเขียวทั้งที่ข้างในพังได้ยังไง

4.4 พิสูจน์ด้วยตาเอง — ฆ่าตัวรันทิ้งกลางงาน

ทีนี้ตอบคำถามเปิดบทด้วยการทำจริง · บล็อกข้างล่างใช้วงรอบเดียวกับบล็อกที่แล้วทุกบรรทัด แต่ให้ตัวรันเป็นโปรเซสลูก · ตัวรันตัวที่ 1 จะค้างอยู่ใน transform นาน ๆ แล้วตัวแม่ฆ่ามันทิ้งตรงนั้น — ไม่ได้ขอให้หยุดดี ๆ แต่ฆ่าเหมือนเครื่องดับ · จากนั้นเปิดตัวรันตัวที่ 2 เป็นโปรเซสใหม่ ซึ่งไม่ได้รับอะไรจากตัวแรกเลยนอกจากไฟล์ในโฟลเดอร์

ชิ้นใหม่ชิ้นเดียวคือ recover() ที่ตัวรันเรียกทุกครั้งที่เปิดขึ้นมา

import json
import subprocess
import sys
import tempfile
import time
from pathlib import Path

NEEDS = {
    "receive_a": [],
    "receive_b": [],
    "receive_c": [],
    "merge_check": ["receive_a", "receive_b", "receive_c"],
    "transform": ["merge_check"],
    "load": ["transform"],
    "notify": ["load"],
}
RETRIES = 1


def execute(run, task, attempt):
    if RUNNER == "1" and task == "transform":
        (WORK / "transform-started").touch()
        time.sleep(60)  # งานยาว — ตัวรันตัวที่ 1 จะถูกฆ่าระหว่างบรรทัดนี้


def read():
    return json.loads(TABLE.read_text()) if TABLE.exists() else {}


def write(table):
    TABLE.write_text(json.dumps(table))


def log(line):
    with (WORK / "runner.log").open("a", encoding="utf-8") as f:
        f.write(f"ตัวรัน {RUNNER}: {line}\n")


def tick():
    """หนึ่งรอบของวงรอบ — โค้ดเดียวกับบล็อกที่แล้วทุกบรรทัด"""
    table, changes = read(), []
    for run, column in table.items():
        seen = {task: cell["state"] for task, cell in column.items()}
        for task, ups in NEEDS.items():
            cell = column[task]
            if seen[task] not in ("none", "up_for_retry"):
                continue
            if any(seen[u] in ("failed", "upstream_failed") for u in ups):
                cell["state"] = "upstream_failed"
            elif all(seen[u] == "success" for u in ups):
                cell["try"] += 1
                cell["state"] = "running"
                write(table)
                try:
                    execute(run, task, cell["try"])
                    cell["state"] = "success"
                except Exception:
                    cell["state"] = "up_for_retry" if cell["try"] <= RETRIES else "failed"
            else:
                continue
            write(table)
            changes.append(f"{run} {task:<12} {seen[task]} → {cell['state']}")
    return changes


def recover():
    """← ชิ้นใหม่: เปิดขึ้นมาเจอช่องที่ค้าง running = ตัวรันก่อนหน้าตายกลางงาน คืนช่องนั้นให้รันใหม่"""
    table = read()
    for run, column in table.items():
        for task, cell in column.items():
            if cell["state"] == "running":
                cell["state"] = "none"  # ตัวนับครั้งไม่ลด — ครั้งที่ตายไปเกิดขึ้นจริง
                log(f"{run} {task:<12} running → none (คนรันเดิมตายกลางงาน)")
    write(table)


if len(sys.argv) == 3:  # ← ตัวลูก = ตัวรันหนึ่งตัว
    WORK, RUNNER = Path(sys.argv[1]), sys.argv[2]
    TABLE = WORK / "table.json"
    recover()
    while changes := tick():
        for line in changes:
            log(line)
    sys.exit()

# ── ตัวแม่: เปิดตัวรัน ฆ่ามันกลางงาน แล้วเปิดตัวใหม่ ──
WORK = Path(tempfile.mkdtemp())
TABLE = WORK / "table.json"
write({"2026-07": {task: {"state": "none", "try": 0} for task in NEEDS}})


def show(title):
    print(title)
    for task, cell in read()["2026-07"].items():
        print(f"  {task:<12} {cell['state']} {cell['try']}")


first = subprocess.Popen([sys.executable, __file__, str(WORK), "1"])
while not (WORK / "transform-started").exists():
    time.sleep(0.05)
first.kill()  # ไม่ได้ขอให้หยุด — ฆ่าทิ้งเหมือนเครื่องดับ
first.wait()
show("ตารางบนดิสก์หลังตัวรันตัวที่ 1 ถูกฆ่า")
subprocess.run([sys.executable, __file__, str(WORK), "2"], check=True)
show("ตารางบนดิสก์หลังตัวรันตัวที่ 2 (โปรเซสใหม่) ทำงานจนจบ")
print("บันทึก")
print((WORK / "runner.log").read_text(encoding="utf-8"), end="")

ผลรันจริง

ตารางบนดิสก์หลังตัวรันตัวที่ 1 ถูกฆ่า
  receive_a    success 1
  receive_b    success 1
  receive_c    success 1
  merge_check  success 1
  transform    running 1
  load         none 0
  notify       none 0
ตารางบนดิสก์หลังตัวรันตัวที่ 2 (โปรเซสใหม่) ทำงานจนจบ
  receive_a    success 1
  receive_b    success 1
  receive_c    success 1
  merge_check  success 1
  transform    success 2
  load         success 1
  notify       success 1
บันทึก
ตัวรัน 1: 2026-07 receive_a    none → success
ตัวรัน 1: 2026-07 receive_b    none → success
ตัวรัน 1: 2026-07 receive_c    none → success
ตัวรัน 1: 2026-07 merge_check  none → success
ตัวรัน 2: 2026-07 transform    running → none (คนรันเดิมตายกลางงาน)
ตัวรัน 2: 2026-07 transform    none → success
ตัวรัน 2: 2026-07 load         none → success
ตัวรัน 2: 2026-07 notify       none → success

ตารางหลังตัวรันตัวแรกตายบอกความจริงอยู่ครึ่งหนึ่ง · ขั้นที่สำเร็จแล้วถูกจดครบ แต่ transform ค้างเป็น running ทั้งที่ไม่มีใครทำมันอยู่แล้ว — ตารางพูดไม่จริงได้ เมื่อคนเขียนตารางตายกลางทาง · นี่คือเหตุผลที่วงรอบจด running ก่อนลงมือ: อย่างน้อยตารางยังบอกได้ว่าค้างอยู่ที่ช่องไหน

ตัวรันตัวที่ 2 จึงต้องมีกติกาเพิ่มหนึ่งข้อ เจอช่องที่ค้าง running ตอนเปิดขึ้นมา = คนรันเดิมตายไปแล้ว คืนช่องนั้นให้รันใหม่ · ส่วนช่องที่ success แล้วไม่ถูกแตะ — receive ทั้งสามกับ merge_check ยังเป็นครั้งที่ 1

เฉลยคำถามเปิดบท: ทำ transform ใหม่ (ครั้งที่ 2) แล้วเดินต่อ — ไม่ข้าม เพราะไม่มีหลักฐานว่ามันเสร็จ · ไม่เริ่มทั้งรอบใหม่ เพราะตารางบอกว่าขั้นไหนเสร็จแล้ว · ถ้าคุณตอบว่า "ข้าม เพราะมันเริ่มไปแล้ว" ลองนึกว่ามันตายตอนเพิ่งเปิดไฟล์ได้บรรทัดเดียว · ถ้าตอบว่า "เริ่มใหม่ทั้งรอบ" นั่นคือเช้าวันที่ 1 มิถุนายนของบทที่ 1

ผลรันนี้คือแนวคิดที่ต้องติดตัวไปทั้งเล่ม: รอบหนึ่งคือคอลัมน์ในตาราง ไม่ใช่โปรเซสที่กำลังวิ่ง · โปรเซสตายได้ คอลัมน์ไม่ตายตาม และใครก็ตามที่อ่านตารางได้ ก็เดินรอบนั้นต่อได้

แต่สังเกตราคาที่จ่าย: transform ถูกทำซ้ำ · ถ้ามันเขียนข้อมูลออกไปครึ่งหนึ่งก่อนตาย การทำซ้ำต้องไม่ทำให้ข้อมูลเบิ้ล — เรื่องเดียวกับเช้าวันที่ 1 มิถุนายน แค่เล็กลงเหลือขั้นเดียว (บทที่ 12)

⚠️ recover() ของเรามีสมมติฐานซ่อนอยู่: มีตัวรันได้ทีละตัว ช่องที่ค้าง running จึงแปลว่าคนรันตายแน่นอน · ตัวคุมจริงมีหลายโปรเซสทำงานพร้อมกัน ต้องแยกให้ออกว่าช่องที่ค้างคือคนรันตายแล้ว หรือแค่ยังทำไม่เสร็จ ซึ่งยากกว่านี้มาก

4.5 ทุกคำสั่งคือการแก้ตาราง

จากนี้ไป ทุกครั้งที่เจอคำสั่งใหม่ของ Airflow ให้ถามคำถามเดียว: มันแก้ตารางตรงไหน · ตารางนี้คือคำตอบล่วงหน้าของคำสั่งหลักที่จะเจอ

คำสั่งตารางก่อนตารางหลังเจอเต็ม ๆ ในบทที่
เริ่มรอบใหม่ (trigger)ยังไม่มีคอลัมน์ของรอบนี้คอลัมน์ใหม่ ทุกช่องเป็น none6
ลองใหม่ (retry)ช่องพัง แต่ยังมีสิทธิ์up_for_retry → รอเวลา → running อีกครั้ง10
ล้างสถานะแล้วรันต่อ (clear)ช่องที่พังกับปลายน้ำของมันกลับเป็น none · ต้นน้ำไม่ถูกแตะ · ตัวนับครั้งที่ลองนับต่อ ไม่เริ่มที่ 111
ย้อนรัน (backfill)ไม่มีคอลัมน์ของเดือนที่ขาดเพิ่มคอลัมน์ย้อนหลัง13
ตัวรันตายกลางงานช่องค้าง runningตัวคุมคืนช่องนั้นให้รันใหม่4 (บทนี้)

4.6 แบบจำลองนี้หยุดตรงไหน

แบบจำลองที่ไม่บอกว่าตัวเองหยุดตรงไหน จะกลายเป็นความเข้าใจผิดเสียเอง · ของจริงต่างจากตารางของบทนี้ห้าข้อ

  1. ช่องไม่ได้เปลี่ยนพร้อมกันเป็นจังหวะ — ของจริงเปลี่ยนทีละช่องเมื่อมีคนทำเสร็จ และมีสถานะผ่านทางระหว่าง none กับ running (scheduled · queued)
  2. วงรอบไม่ได้อยู่ในโปรเซสเดียว — Airflow แบ่งงานของวงรอบให้หลายโปรเซส · เครื่องทดลองที่ลงในบทที่ 5 จะพิมพ์บันทึกออกมาจากสี่ส่วน (บทที่ 15)
  3. ตารางไม่ได้มีแผ่นเดียว — ของจริงอยู่ในฐานข้อมูล แยกเป็นตารางของรอบ (dag_run) กับตารางของช่อง (task_instance)
  4. ตัวงานไม่ได้เขียนตารางเอง — ใน Airflow 3 โค้ดของงานแตะฐานข้อมูลของ Airflow ตรง ๆ ไม่ได้แล้ว การเปลี่ยนสถานะทุกครั้งวิ่งผ่าน API ของมัน (หน้า Upgrading to Airflow 3 ในเอกสารทางการ)
  5. คนรันตายจริงกับคนรันช้า ต้องแยกให้ออก — ข้อควรระวังท้ายหัวข้อ 4.4

ทั้งห้าข้อเป็นเรื่องว่าใครเขียนตาราง และเขียนเมื่อไหร่ · ตัวตารางกับคำถามของวงรอบยังเป็นแบบเดียวกับบทนี้

4.7 ปริศนาที่จะไขในบทหลัง

  1. วงรอบมีวงเดียว แต่เครื่องทดลองของบทที่ 5 จะพิมพ์บันทึกออกมาจากสี่ส่วน — แต่ละส่วนทำอะไร → บทที่ 15
  2. สถานะของรอบดูจากขั้นปลายสุดเท่านั้น — ถ้าขั้นปลายสุดถูกตั้งให้รันแม้ต้นน้ำพัง รอบจะเป็นสีอะไร → บทที่ 7 และ 11
  3. transform ถูกทำซ้ำหลังตัวรันตาย — ถ้ามันเขียนไปครึ่งหนึ่งแล้วจะเป็นยังไง → บทที่ 12

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

  1. คอลัมน์รอบเดือนกันยายนตอนนี้เป็นแบบนี้: receive_a success 1 · receive_b success 1 · receive_c up_for_retry 1 · ขั้นที่เหลือ none 0 · ทุกขั้นได้สิทธิ์ลองใหม่ 1 ครั้ง · ถ้า receive_c พังอีกครั้ง สามจังหวะถัดไปตารางจะเปลี่ยนยังไง และรอบนี้จบเป็นสถานะอะไร
  2. เพื่อนบอกว่า "ตัวรันของฉันเก็บสถานะใน dict ตลอดอายุโปรเซส ไม่เคยเขียนลงดิสก์ เร็วกว่าเยอะ" — คำถามข้อไหนในสามข้อของบทที่ 1 ที่ตัวรันนี้ตอบไม่ได้ และจะรู้ตัวเมื่อไหร่
  3. เพิ่มขั้น archive ที่รอแค่ receive ทั้งสาม (และไม่มีใครรอมัน) เข้าไปในท่อของหัวข้อ 4.3 · รอบเดือนสิงหาคมจะจบเป็นสถานะอะไร เพราะอะไร

เฉลยคำถามท้ายบท#

บทที่ 1

  1. ลูกค้า 80 คนแรกได้อีเมลซ้ำเป็นฉบับที่สอง · ความรู้ที่หายไปคือ "ส่งถึงใครไปแล้วบ้าง" ซึ่งอยู่ในตัวแปรของลูป — ไม่มีใครจดไว้นอกโปรเซส จึงไม่มีทางเริ่มที่คนที่ 81 ได้ นอกจากไปไล่ดูกล่องอีเมลขาออกเอง
  2. รหัสจบการทำงานบอกแค่ว่าโปรเซสจบโดยไม่มี error ที่ไม่มีใครจับ ไม่ได้บอกว่าข้อมูลที่ได้ถูกต้อง · เดือนกรกฎาคมทุกขั้นทำงานตามที่เขียนไว้ทุกบรรทัด แต่ขั้นอ่านไฟล์ได้ศูนย์แถว และไม่มีขั้นไหนถามว่าศูนย์แถวสมเหตุสมผลไหม
  3. เลขแถวทางซ้ายกระโดด (4 ไป 38) · ปุ่มที่หัวคอลัมน์เป็นรูปกรวย · แถบสถานะขึ้น Filter Mode — เจอข้อใดข้อหนึ่ง = มีแถวถูกซ่อนด้วยตัวกรอง ล้างตัวกรองหรืออ่านไฟล์ตรง ๆ ก่อนสรุป

บทที่ 2

  1. รันใหม่แล้วไม่มีอะไรถูกจดไว้เลย ตัวรันเริ่มขั้นที่ 1 ใหม่ทั้งหมด — อาการเดียวกับสคริปต์ก้อนเดียว · ที่ขาดคือการจดทันทีที่แต่ละขั้นเสร็จ
  2. ตัวรันอ่าน state.json เจอว่าทุกขั้นสำเร็จแล้ว จึงข้ามหมด เดือนมิถุนายนไม่ถูกทำเลย และไม่มี error ให้เห็น · ต้นเหตุคือไฟล์สถานะไม่รู้จักคำว่า "รอบ" — บทที่ 4 แก้ด้วยการให้แต่ละรอบเป็นคอลัมน์ของตัวเอง
  3. ทางหนึ่งที่ทำได้: ในบล็อก except ตรงที่ตรวจ cell["tries"] > RETRIES ให้เขียนไฟล์ก่อน return เช่น (WORK / "ALERT.txt").write_text(f"{name} ลองครบ {cell['tries']} ครั้งแล้วยังพัง: {e}", encoding="utf-8") · ของจริงเปลี่ยนการเขียนไฟล์เป็นส่งข้อความหาคน แต่ตำแหน่งในโค้ดเป็นจุดเดียวกัน — จุดที่ตัวรันรู้แน่แล้วว่ายอมแพ้

บทที่ 3

  1. ขั้นที่ไม่รอใครเลยคือดึงรายการสั่งซื้อ ดึงรายชื่อสินค้า และดึงรายชื่อสาขา จึงอยู่จังหวะแรกทั้งสามขั้น · จับคู่กับสินค้ารอสองขั้นแรก (จังหวะที่ 2) · จับคู่กับสาขารอผลจับคู่กับสินค้าและรายชื่อสาขา (จังหวะที่ 3) · สรุปรายสาขา (จังหวะที่ 4) · ส่งรายงาน (จังหวะที่ 5) — ห้าจังหวะ
  2. ไม่มีเลย notify เป็นขั้นปลายสุด ไม่มีใครรอมัน · ความพังไหลลงปลายน้ำเท่านั้น และขั้นนี้ไม่มีปลายน้ำ
  3. เพราะขั้นที่อยู่ในวงไม่มีวันพร้อม รันส่วนที่เหลือไปก่อนก็ได้ท่อที่จบไม่ได้ทุกรอบ และรู้ตัวช้าไปหนึ่งคืนเสมอ · เจอตอนประกาศกราฟคือการบอกความผิดพลาด ณ จุดที่มันเกิด — ตอนแก้ NEEDS — ก่อนจะกลายเป็นรอบที่ค้างอยู่กลางทางทุกคืน

บทที่ 4

  1. จังหวะที่ 1: receive_c พังครั้งที่ 2 หมดสิทธิ์ → failed 2 · จังหวะที่ 2: merge_check → upstream_failed 0 · จังหวะที่ 3: transform → upstream_failed 0 · (ต่ออีกสองจังหวะ load กับ notify ก็เป็น upstream_failed 0) · รอบจบเป็น failed เพราะขั้นปลายสุด notify เป็น upstream_failed
  2. ข้อ "ทำถึงไหนแล้ว" — สถานะอยู่ในโปรเซส ตายแล้วหายตาม · จะรู้ตัวในคืนแรกที่ตัวรันตายกลางงาน แล้วเปิดใหม่ต้องเริ่มทั้งรอบ ไม่ใช่ตอนที่ทุกอย่างปกติดี
  3. failed · ท่อนี้มีขั้นปลายสุดสองขั้น archive สำเร็จ (receive ทั้งสามผ่าน) แต่ notify เป็น upstream_failed · กติกาคือมีขั้นปลายสุดตัวไหน failed หรือ upstream_failed รอบก็เป็น failed — ขั้นปลายสุดที่ผ่านไม่ได้ช่วยกลบตัวที่พัง

ถ้าเนื้อหานี้มีประโยชน์ —เลี้ยงกาแฟสักแก้วหรือโอนผ่านพร้อมเพย์