Skip to content
Tayakorn

← คู่มือ 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 0
ไฟล์มาช้าชั่วคราว · หัวข้อ 10.3receive_cup_for_retry 1รอ 20 วินาที · ไฟล์มาsuccess 2ไฟล์ผิดรูป + AirflowFailExceptionถาวร · หัวข้อ 10.3gatefailed 1หยุดทันที ไม่ลองต่อครั้งที่ 2 ไม่ได้ลองครั้งที่ 3 ไม่ได้ลองไฟล์ผิดรูป + error ธรรมดาถาวร · บทที่ 7transformup_for_retry 1รอ 10 วินาที · ข้อความเดิมfailed 2
FIG 10.1 — สามขั้นที่ได้สิทธิ์ลองใหม่ ต่างกันมิติเดียว: รอแล้วเหตุของความพังหายไหม · ไฟล์มาช้า — ลองใหม่แล้วผ่าน · ไฟล์ผิดรูป + AirflowFailException — หยุดครั้งแรก ไม่ใช้สิทธิ์ที่เหลือ · ไฟล์ผิดรูปที่ยก 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 หยุดทันทีโดยไม่ถามกฎ · ข้อดีคือกฎแยกจากเนื้องาน ข้อเสียคือผูกรุ่น

ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)

  1. ขั้นหนึ่งเรียก API ภายนอก เจอ error สี่แบบ: หมดเวลารอ · รหัสเข้าใช้หมดอายุ · ข้อมูลที่ได้ขาดฟิลด์บังคับ · ปลายทางตอบว่าส่งถี่เกิน — แบบไหนควรลองใหม่ แบบไหนควรหยุดทันที
  2. ในหัวข้อ 10.2 หลัง receive_c พังครั้งที่ 3 ตารางของบทที่ 4 เปลี่ยนต่อยังไงทีละจังหวะ และทำไมรอบเป็นสีแดงทั้งที่ท่อนี้ไม่มี watcher
  3. ตัวรันของบทที่ 2 ลองใหม่ทุก error เท่ากัน · ถ้าจะให้หยุดทันทีเมื่อไฟล์ผิดรูป ต้องแก้ตรงไหน และชิ้นนั้นเทียบได้กับอะไรใน Airflow

Read the full book