← คู่มือ Airflow — ทำไมงานหลายขั้นต้องมีคนคุมลำดับ
LEVEL 2 · ระดับกลาง
retry ช่วยได้กับความพังแบบไหน
บทที่ 7 มีคู่ที่น่าเทียบอยู่แล้ว · receive_c ของรอบที่สองพังครั้งแรกแล้วผ่านครั้งที่ 2 · transform ของรอบที่สามพังสองครั้งด้วยข้อความเดียวกันทุกตัวอักษร — ครั้งที่ 2 เสียเวลาเปล่า 10 วินาที · สิทธิ์ลองใหม่เท่ากัน ต่างกันมิติเดียว: รอแล้วเหตุของความพังหายเองไหม
retry ซื้อเวลา — ช่วยได้เฉพาะความพังที่เวลาแก้ให้
เดือนกันยายนมีทั้งสองแบบ · ไฟล์สาขา ค มาช้า (อยู่ใน late/) — รอแล้วจะมา · ไฟล์สาขา ก ไม่มีคอลัมน์ "จำนวน" — รอเท่าไหร่ก็ไม่มี
10.1 ให้โค้ดบอกว่าความพังแบบไหนถาวร
ตัวคุมเห็นแค่ว่าขั้นพัง ไม่รู้ว่าทำไม · คนที่รู้คือโค้ดของขั้นนั้น ฉบับนี้จึงแยกสองชนิดไว้ในโค้ด
- ชั่วคราว —
receiveเช็กก่อนว่าไฟล์มาหรือยัง ถ้ายังก็ยกFileNotFoundErrorที่บอกสาขาและเดือน แล้วปล่อยให้ retries ทำงาน - ถาวร —
gateจับBadFileของบทที่ 9 แล้วยกAirflowFailExceptionแทน · หน้า Tasks ทางการเขียนว่ามันทำให้ขั้นfailedทันทีโดยไม่สนสิทธิ์ที่เหลือ
แล้วเปิดสิทธิ์ลองใหม่ให้ทุกขั้นสองครั้งผ่าน default_args · แล็บรอ 20 วินาทีให้อ่านผลทัน ของจริงควรรอเป็นนาที
💡 ค่าตั้งต้น: Airflow 3.3.2 ตั้ง retries = 0 และรอ 300 วินาที (หน้า config reference · ยันในแล็บ) — ไม่มีขั้นไหนถูกลองซ้ำจนกว่าคุณจะบอก · ถ้าลองใหม่ดูมีแต่ได้ ทำไมค่าตั้งต้นถึงเป็น 0 — บทที่ 12 จะตอบ
# ~/airflow-lab/dags/stock_monthly.py — บทที่ 10: ลองใหม่ได้ 2 ครั้ง ยกเว้นไฟล์ผิดรูปที่หยุดทันที
from datetime import datetime, timedelta
from pathlib import Path
import stock_gate
import stock_steps as steps
from airflow.sdk import dag, task
from airflow.sdk.exceptions import AirflowFailException, AirflowSkipException
LAB = Path.home() / "airflow-lab"
@dag(
schedule=None,
start_date=datetime(2026, 5, 1),
catchup=False,
params={"month": "2026-07"},
default_args={"retries": 2, "retry_delay": timedelta(seconds=20)}, # แล็บรอ 20 วินาที ของจริงหลายนาที
)
def stock_monthly():
@task
def arrived(**context):
month = context["params"]["month"]
files = sorted((LAB / "inbox").glob(f"stock_{month}_*.xlsx"))
if not files: # ยังไม่มีไฟล์ของเดือนนี้สักสาขา — ไม่มีอะไรให้ทำ ไม่ใช่ความผิดของใคร
raise AirflowSkipException(f"ยังไม่มีไฟล์ของเดือน {month}")
return len(files)
@task
def receive(branch, **context):
month = context["params"]["month"]
path = LAB / "inbox" / f"stock_{month}_{branch}.xlsx"
if not path.exists(): # ไฟล์มาช้า — ชั่วคราว ปล่อยให้ retries ลองใหม่
raise FileNotFoundError(f"ไฟล์สาขา {stock_gate.NAMES[branch]} เดือน {month} ยังมาไม่ถึง")
rows, missing = steps.read_branch(path)
staged = steps.stage(branch, rows, missing, LAB / "staging" / month / f"{branch}.json")
return {**staged, "problems": stock_gate.inspect(path, rows, missing, month)}
@task
def gate(staged):
for s in staged:
print(f"ด่าน ← receive_{s['branch']}: {s['rows']} แถว · ปัญหา {len(s['problems'])} ข้อ")
try:
return stock_gate.check(staged)
except stock_gate.BadFile as e: # ไฟล์ผิดรูป ลองกี่ครั้งก็ไม่หาย — หยุดเลย ไม่ใช้ retries ที่เหลือ
raise AirflowFailException(str(e)) from e
@task
def transform(staged, **context):
return steps.transform(staged, LAB / "staging" / context["params"]["month"] / "all.json")
@task
def load(summary, **context):
return steps.load(summary, LAB / "warehouse.db", context["params"]["month"])
@task
def notify(staged, summary, in_db):
print(steps.report(staged, summary, in_db))
received = [receive.override(task_id=f"receive_{b}")(b) for b in "abc"]
arrived() >> received
staged = gate(received)
summary = transform(staged)
notify(staged, summary, load(summary))
stock_monthly()10.2 ไฟล์ไม่มาภายในโควตา
รันเดือนกันยายนโดยปล่อยไฟล์สาขา ค ไว้ใน late/
airflow dags test stock_monthly -c '{"month": "2026-09"}' 2>&1 | grep "^FileNotFoundError"ผลรันจริง
FileNotFoundError: ไฟล์สาขา ค เดือน 2026-09 ยังมาไม่ถึง
FileNotFoundError: ไฟล์สาขา ค เดือน 2026-09 ยังมาไม่ถึง
FileNotFoundError: ไฟล์สาขา ค เดือน 2026-09 ยังมาไม่ถึงpython ~/airflow-lab/states.py stock_monthlyผลรันจริง
รอบล่าสุดของ stock_monthly {"month": "2026-09"} → failed
arrived success 1
receive_a success 1
receive_b success 1
receive_c failed 3
gate upstream_failed 0
transform upstream_failed 0
load upstream_failed 0
notify upstream_failed 0สามบรรทัด = สามครั้ง: ครั้งแรกบวกสิทธิ์อีกสอง ห่างกันราว 20 วินาที (ทั้งรอบราว 51 วินาทีในแล็บ) · สิทธิ์หมดแล้วไฟล์ยังไม่มา receive_c จึงเป็น failed 3 และ gate ไม่ได้รันเลย
10.3 ไฟล์มาระหว่างรอ — และด่านที่ไม่ลองซ้ำ
รอบเดียวกันอีกครั้ง แต่ย้ายไฟล์สาขา ค เข้า inbox ระหว่างที่ receive_c รอรอบสอง · ผลขึ้นกับจังหวะที่ย้ายไฟล์ ตัวตรวจอัตโนมัติรันซ้ำให้ไม่ได้ ข้างล่างคือบันทึกจากแล็บ (ย้ายทันทีที่บรรทัด FileNotFoundError แรกขึ้น)
บันทึกจากแล็บ
── หน้าต่างที่ 1 ──
$ airflow dags test stock_monthly -c '{"month": "2026-09"}' 2>&1 | grep -E "^FileNotFoundError|AirflowFailException:"
FileNotFoundError: ไฟล์สาขา ค เดือน 2026-09 ยังมาไม่ถึง
airflow.sdk.exceptions.AirflowFailException: ด่านหยุดท่อ — สาขา ก: ไม่มีคอลัมน์ จำนวน
── หน้าต่างที่ 2 ระหว่างที่ receive_c รอรอบสอง ──
$ mv ~/airflow-lab/inbox/late/stock_2026-09_c.xlsx ~/airflow-lab/inbox/
── หลังรอบจบ ──
$ python ~/airflow-lab/states.py stock_monthly
รอบล่าสุดของ stock_monthly {"month": "2026-09"} → failed
arrived success 1
receive_a success 1
receive_b success 1
receive_c success 2
gate failed 1
transform upstream_failed 0
load upstream_failed 0
notify upstream_failed 0AirflowFailException — หยุดครั้งแรก ไม่ใช้สิทธิ์ที่เหลือ · ไฟล์ผิดรูปที่ยก error ธรรมดา (บทที่ 7) — ลองครบแล้วพังด้วยข้อความเดิม เสียเวลารอเปล่า · แถวบนสองแถวคือคอลัมน์เดียวกันในบันทึกของหัวข้อ 10.3คอลัมน์เดียวได้ทั้งคู่เทียบ · receive_c success 2 — รออีก 20 วินาทีแล้วไฟล์มา · gate failed 1 — สิทธิ์อีกสองครั้งยังอยู่ แต่ AirflowFailException บอกว่าไม่ต้องใช้ · ถ้า gate ยก BadFile ตรง ๆ แบบบทที่ 9 มันจะพังครบสามครั้งด้วยข้อความเดิม แล้วเสียเวลารอเปล่าอีกราว 40 วินาที (ถ้ารอ 5 นาทีแบบค่าตั้งต้น = 10 นาที) ก่อนมีใครรู้
⚠️ ตั้ง retries สูงไว้ทุกขั้นไม่ได้ทำให้ปลอดภัยขึ้น · ความพังถาวรแค่ถูกรายงานช้าลง · และทุกครั้งที่ลองใหม่คือการรันขั้นนั้นซ้ำ — ขั้นที่รันซ้ำแล้วได้ผลไม่เท่าเดิม ยิ่งลองยิ่งเสียหาย (บทที่ 12)
10.4 ตั้งแต่ 3.3: ประกาศกฎรายชนิด error
เนื้อหลักของบทใช้ AirflowFailException เพราะไม่พึ่งของใหม่ · Airflow 3.3.0 เพิ่ม retry policy (release notes: Pluggable Retry Policies) ให้ประกาศกฎแยกตามชนิด error บนขั้นได้ เช่น retry_policy=ExceptionRetryPolicy(rules=[RetryRule(exception=BadFile, action=RetryAction.FAIL)]) จาก airflow.sdk · ผลในแล็บรุ่น 3.3.2 ระหว่างเตรียมเล่ม: กฎ FAIL หยุดที่ครั้งที่ 1 แม้ตั้ง retries ไว้ 3 · กฎ RETRY ใส่เวลารอของตัวเองได้ (ตั้ง 3 วินาที ลองใหม่หลัง ~3.4 วินาที แทน 30 วินาทีของขั้น) · กฎลองเกิน retries ของขั้นไม่ได้ · AirflowFailException หยุดทันทีโดยไม่ถามกฎ · ข้อดีคือกฎแยกจากเนื้องาน ข้อเสียคือผูกรุ่น
ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)
- ขั้นหนึ่งเรียก API ภายนอก เจอ error สี่แบบ: หมดเวลารอ · รหัสเข้าใช้หมดอายุ · ข้อมูลที่ได้ขาดฟิลด์บังคับ · ปลายทางตอบว่าส่งถี่เกิน — แบบไหนควรลองใหม่ แบบไหนควรหยุดทันที
- ในหัวข้อ 10.2 หลัง
receive_cพังครั้งที่ 3 ตารางของบทที่ 4 เปลี่ยนต่อยังไงทีละจังหวะ และทำไมรอบเป็นสีแดงทั้งที่ท่อนี้ไม่มี watcher - ตัวรันของบทที่ 2 ลองใหม่ทุก error เท่ากัน · ถ้าจะให้หยุดทันทีเมื่อไฟล์ผิดรูป ต้องแก้ตรงไหน และชิ้นนั้นเทียบได้กับอะไรใน Airflow