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

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

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

ลองคุมลำดับเองด้วย 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 ทันทีที่ขั้นไหนลองครบแล้วยังพัง

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