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

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

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

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

ก่อนอ่านต่อ ลองเดาก่อน: คืนหนึ่งไฟล์สาขา ข เสียจนอ่านไม่ได้ ตัวรันของบทที่ 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. ทำไมการเจอวงตั้งแต่ตอนประกาศกราฟ ถึงดีกว่าปล่อยให้ท่อรันขั้นที่ไม่อยู่ในวงไปก่อนแล้วค่อยรู้ทีหลัง

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