← คู่มือ Airflow — ทำไมงานหลายขั้นต้องมีคนคุมลำดับ
LEVEL 0 · ปูพื้นจากศูนย์
ลองคุมลำดับเองด้วย Python ล้วน
บทที่แล้วจบด้วยประโยคว่าความรู้ว่างานทำถึงไหนแล้วต้องอยู่นอกตัวงาน บทนี้ลองทำตามประโยคนั้นด้วยมือ — เขียนตัวรันเล็ก ๆ ด้วย Python ล้วน แล้วดูว่าต้องเพิ่มชิ้นอะไรบ้าง กว่าจะรอดคืนแบบวันที่ 1 มิถุนายนได้
ไม่ต้องติดตั้งอะไรเลยในบทนี้ และชิ้นส่วนที่โผล่มาระหว่างทางคือชิ้นเดียวกับที่ตัวคุมลำดับงานทุกตัวมี
2.1 แยกงานเป็นขั้น แล้วจดทุกครั้งที่ขั้นหนึ่งเสร็จ
เปลี่ยนจากสคริปต์ก้อนเดียวสองจุด
- แยกงานเป็นขั้นที่มีชื่อ — หนึ่งขั้นหนึ่งฟังก์ชัน แล้วเก็บไว้ในรายการที่บอกลำดับ · ชื่อขั้นเป็นภาษาอังกฤษเพราะจะใช้ชื่อชุดนี้ไปจนจบเล่ม:
receive_areceive_breceive_c(รับไฟล์สามสาขา) ·merge_check(รวมและตรวจคอลัมน์) ·transform(แปลง) ·load(โหลด) ·notify(ส่งรายงาน) - จดลงไฟล์ทันทีที่แต่ละขั้นเสร็จ — ไฟล์
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 สิ่งที่ตัวรันนี้ยังทำไม่ได้
- มันเดินตามรายการทีละขั้น — รับไฟล์สามสาขาไม่ได้รอกันเลย แต่ต้องต่อคิวกัน · และถ้า
receive_bพังถาวรreceive_cจะไม่ถูกลองเลยสักครั้ง ทั้งที่ไม่เกี่ยวกัน → บทที่ 3 - มันมีรอบเดียว —
state.jsonไม่รู้ว่าเป็นของเดือนไหน · ถ้าเอาตัวรันนี้ไปปิดยอดเดือนมิถุนายนโดยไม่ลบไฟล์ มันจะเห็นว่าทุกขั้นสำเร็จแล้ว แล้วข้ามหมด เดือนมิถุนายนไม่ถูกทำเลยทั้งที่ทุกอย่างดูเรียบร้อย → บทที่ 4 - ถ้าตัวรันเองตายระหว่างขั้น —
state.jsonจดแค่ขั้นที่เสร็จแล้ว ไม่ได้จดว่ากำลังทำอะไรอยู่ → บทที่ 4
ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)
- เพื่อนส่งตัวรันมาให้ดู มันจด
state.jsonครั้งเดียวตอนจบรอบ · ถ้าพังที่ขั้นที่ 6 แล้วรันใหม่ จะเกิดอะไรขึ้น และขาดชิ้นไหน - ถ้าใช้ตัวรันในบทนี้ปิดยอดเดือนมิถุนายนต่อจากเดือนพฤษภาคม โดยไม่ลบ
state.jsonจะเกิดอะไรขึ้น - ลองเขียนเอง: เพิ่มให้ตัวรันในหัวข้อ 2.2 เขียนไฟล์
ALERT.txtบอกชื่อขั้นกับข้อความ error ทันทีที่ขั้นไหนลองครบแล้วยังพัง