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

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

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

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

บทนี้สำคัญที่สุดในเล่ม · ทุกคำสั่งของ 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 — ขั้นปลายสุดที่ผ่านไม่ได้ช่วยกลบตัวที่พัง

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

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