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

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

LEVEL 3 · ระดับสูง

รอบนี้ประมวล⁠ผลข้อมูลช่วงไหน

ตั้งแต่บทที่ 8 เดือนถูกส่งเข้าท่อด้วยมือผ่าน -c · บทนี้ให้นาฬิกาบอกแทน — และคำตอบของ Airflow 3 ไม่ใช่สิ่งที่คนย้ายมาจากรุ่น 2 คาด

ลองเดาก่อน: ตั้งท่อให้รันตอนเที่ยงคืนวันที่ 1 ทุกเดือน · รอบที่รัน 1 สิงหาคมควรประมวลผลเดือนไหน และ Airflow 3 บอกงานว่าเดือนไหน · จดไว้ เฉลยอยู่หัวข้อ 13.1

13.1 รอบที่รันวันนี้ ไม่ได้ประมวลผลวันนี้

รอบหนึ่งรอบประมวลผลช่วงที่เพิ่งจบ ไม่ใช่ช่วงที่มันรันอยู่

รอบ 1 สิงหาคมต้องประมวลผลกรกฎาคม — ไฟล์ของเดือนครบเมื่อเดือนจบแล้ว · Airflow เรียกช่วงนี้ว่าช่วงข้อมูล (data_interval_start → data_interval_end) · และ Airflow 3 ไม่บอกอะไรเลย ถ้าส่ง cron เป็นสตริง

แล็บให้ DAG สองตัวใช้ cron */2 * * * * เดียวกัน ต่างกันมิติเดียว: ตัวแรกส่งสตริง ตัวที่สองบอกชนิดตารางเวลาเองว่าเป็นแบบมีช่วง (CronDataIntervalTimetable) แล้วให้งานพิมพ์สิ่งที่เห็นจากรอบที่ scheduler สร้าง

บันทึกจากแล็บ

cron สตริง                 logical_date= 09:26  start= 09:26  end= 09:26
CronDataIntervalTimetable  logical_date= 09:24  start= 09:24  end= 09:26

cron สตริงถูกเก็บเป็น CronTriggerTimetable ช่วงยาว 0 — หน้า Timetables ทางการเขียนว่าตารางเวลาแบบนี้ไม่มีแนวคิดเรื่องช่วงข้อมูล · แบบมีช่วงให้งวดก่อนหน้าแบบรุ่น 2 และ logical_date คือต้นช่วง ไม่ใช่เวลาที่รัน

⚠️ ความเข้าใจผิด: "ใช้ datetime.now() ในงานเพื่อรู้ว่ารอบนี้คือเดือนไหน" — หน้า Best Practices ทางการเตือนไว้ว่าห้ามใช้ now() ตัดสินอะไรสำคัญในขั้น · ผิดเพราะ now() คือเวลาที่รัน: retry ข้ามเดือน · clear หลังได้ไฟล์แก้ · backfill เดือนเก่า — รันวันอื่นแล้วได้เดือนผิดโดยไม่มี error · กลไกที่ถูกคือให้ตัวคุมบอกช่วงมา แล้วงานคิดเดือนจากช่วงนั้น · พิสูจน์ที่ backfill ในหัวข้อ 13.3 ซึ่งรันเดือนตุลาคมแต่ได้พฤษภาคมถึงกรกฎาคมถูกทุกเดือน

13.2 ให้ DAG บอกชนิดตารางเวลาเอง

ฉบับนี้เปลี่ยนสองที่ · schedule เป็น CronDataIntervalTimetable รันเที่ยงคืนวันที่ 1 (UTC) · month_of ตัดสินเดือนให้ทุกขั้น: -c ก่อน · ไม่บอกก็ใช้ต้นช่วงข้อมูล · ไม่มีทั้งคู่ = หยุดพร้อมเหตุ

# ~/airflow-lab/dags/stock_monthly.py — บทที่ 13: รอบตามตารางรู้เดือนจากช่วงข้อมูลของตัวเอง
from datetime import datetime, timedelta
from pathlib import Path

import stock_gate
import stock_load
import stock_steps as steps
from airflow.sdk import CronDataIntervalTimetable, dag, task
from airflow.sdk.exceptions import AirflowFailException, AirflowSkipException

LAB = Path.home() / "airflow-lab"


def month_of(context):
    """เดือนที่รอบนี้ประมวลผล: บอกมาทาง -c ก่อน · ไม่บอกก็ใช้ต้นช่วงข้อมูลของรอบ · ไม่มีทั้งคู่ = หยุด"""
    if context["params"]["month"]:
        return context["params"]["month"]
    start = context.get("data_interval_start")  # รอบที่กดรันเองไม่มีคีย์นี้เลย — ["..."] จะยก KeyError
    if start is None:
        raise AirflowFailException("รอบนี้ไม่มีช่วงข้อมูล — บอกเดือนด้วย -c หรือย้อนรันด้วย backfill")
    return f"{start:%Y-%m}"


@dag(
    schedule=CronDataIntervalTimetable("0 0 1 * *", timezone="UTC"),  # ทุกต้นเดือน รอบละหนึ่งเดือนที่จบแล้ว
    start_date=datetime(2026, 5, 1),
    catchup=False,
    params={"month": ""},
    default_args={"retries": 2, "retry_delay": timedelta(seconds=20)},  # แล็บรอ 20 วินาที ของจริงหลายนาที
)
def stock_monthly():
    @task
    def arrived(**context):
        month = month_of(context)
        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 = month_of(context)
        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" / month_of(context) / "all.json")

    @task
    def load(summary, **context):
        return stock_load.load_month(summary, LAB / "warehouse.db", month_of(context))  # ← เดิม steps.load

    @task
    def notify(staged, summary, in_db):
        print(steps.report(staged, summary, in_db))

    @task(trigger_rule="all_done")  # ส่งทุกเช้า ไม่ว่าท่อจะผ่านหรือพัง — แบบ notify ของบทที่ 7
    def morning():
        print("morning: ส่งข้อความเช้าแล้ว")

    @task(trigger_rule="one_failed")  # watcher: มีขั้นไหนพัง ขั้นนี้ได้รันแล้วพังเสมอ
    def watcher():
        raise AirflowFailException("มีขั้นที่พังในรอบนี้")

    found = arrived()
    received = [receive.override(task_id=f"receive_{b}")(b) for b in "abc"]
    found >> received
    staged = gate(received)
    summary = transform(staged)
    loaded = load(summary)
    done = notify(staged, summary, loaded)
    note = morning()
    done >> note
    [found, *received, staged, summary, loaded, done, note] >> watcher()  # ต่อจากทุกขั้น


stock_monthly()

ทำไม context.get · แล็บลอง context["data_interval_start"] ก่อน แล้วรอบที่กดรันเองพังด้วย KeyError — 3.3.2 ไม่ใส่คีย์นี้เลยในรอบแบบนั้น · และ KeyError เป็น error ธรรมดา retries จึงลองซ้ำครบสามครั้งแบบไม่มีทางผ่าน (บทที่ 10)

ลองรันโดยบอกวันที่ให้รอบ

airflow dags test stock_monthly 2026-07-01 2>&1 | grep "ในคลัง"

ผลรันจริง

สต็อกรวม 1,114 ชิ้น (ก 140 · ข 364 · ค 610) จาก 99 แถว · ในคลังเดือนนี้ 99 แถว

1,114 ชิ้นคือยอดของมิถุนายน (ตาราง 8.2: 140 + 364 + 610) — ประโยคแก่นของหัวข้อ 13.1 อีกครั้ง: 1 กรกฎาคมคือจังหวะที่รอบเกิด ช่วงที่เพิ่งจบคือมิถุนายน · สั่ง 15 มิถุนายนได้พฤษภาคม · ไม่ใส่วันที่ (7 ตุลาคม) ได้กันยายน

13.3 ปล่อยให้นาฬิกาเดิน — และย้อนรันเดือนที่ขาด

ฉากนี้ใช้ scheduler จริง จึงเป็นบันทึกจากแล็บ · DAG ที่ยังไม่เคยรัน กด unpause วันที่ 7 ตุลาคม

บันทึกจากแล็บ — ช่วงข้อมูลของแต่ละรอบ กับจำนวนแถวในคลังต่อเดือน

$ airflow dags unpause stock_monthly
scheduled  2026-09-01 → 2026-10-01  failed      ← รอบเดียว: งวดล่าสุดที่จบแล้ว

$ airflow backfill create --dag-id stock_monthly --from-date 2026-05-01 --to-date 2026-07-01
backfill   2026-05-01 → 2026-06-01  success
backfill   2026-06-01 → 2026-07-01  success
backfill   2026-07-01 → 2026-08-01  success
คลัง: 2026-05 99 แถว · 2026-06 99 แถว · 2026-07 99 แถว

$ airflow dags trigger stock_monthly                    ← กดรันเองโดยไม่บอกเดือน
manual     None       → None        failed      arrived failed 1
AirflowFailException: รอบนี้ไม่มีช่วงข้อมูล — บอกเดือนด้วย -c หรือย้อนรันด้วย backfill

$ airflow backfill create --dag-id stock_monthly --from-date 2026-05-01 --to-date 2026-07-01
ไม่มีรอบใหม่ — ทั้งสามเดือนถูกข้ามด้วยเหตุผล already exists
2026-052026-062026-072026-082026-092026-10scheduledbackfillmanualไม่มีรอบ — catchup ปิด ไม่มีอะไรเตือนfailedsuccesssuccesssuccessไม่มีรอบกดรันเอง: ไม่มีช่วงข้อมูล (None → None) · arrived failed 17 ต.ค. · วันที่สั่งทุกคำสั่งแท่ง = ช่วงข้อมูลของรอบ (data_interval_start → end) · backfill สามเดือนเข้าคลังเดือนละ 99 แถว
FIG 13.1 — ฉาก stock_monthly ของหัวข้อ 13.3 บนแกนเดือน · ทุกคำสั่งสั่งวันที่ 7 ตุลาคม แต่แต่ละรอบรู้ช่วงข้อมูลของตัวเอง — แท่งยาวเท่าช่วงนั้น · unpause ได้งวดล่าสุดที่จบแล้วงวดเดียว · backfill เพิ่มคอลัมน์ย้อนหลังทีละเดือน · รอบที่กดเองไม่มีแท่งเพราะไม่มีช่วง · สิงหาคมว่างทุกแถว และไม่มีอะไรเตือน
  • unpause ได้รอบเดียว — start_date คือจุดที่นาฬิกาเริ่มนับงวด (หน้า Scheduler ทางการ: รอบแรกสร้างจาก start_date) แต่ catchup ตั้งต้นเป็น False ในรุ่น 3 จึงไม่ย้อนไปถึงมัน · พฤษภาคมถึงสิงหาคมไม่ถูกรันและไม่มีอะไรเตือน (standalone ที่ดับไปสองชั่วโมงก็เหมือนกัน: ช่วงที่ดับไม่ถูกรันย้อน) · กันยายนพังเพราะไฟล์สาขา ค ยังไม่มา — บทถัดไป
  • backfill = เพิ่มคอลัมน์ย้อนหลัง (ตาราง 4.5) แต่ละคอลัมน์ได้ช่วงของตัวเอง · --from-date กับ --to-date นับ logical_date = ต้นช่วง
  • กดรันเองไม่ได้ช่วงข้อมูล — month_of หยุดตั้งแต่ขั้นแรกแทนที่จะเดา · ใส่วันที่ให้ก็ยังไม่ได้ช่วงของวันนั้น (หน้า DAG Runs ทางการเตือนไว้ · แล็บยันแล้ว) ⇒ ย้อนรันเดือนเก่าด้วย backfill
  • backfill ซ้ำข้ามของเดิม — --reprocess-behavior none (ตั้งต้น) ข้ามวันที่มีรอบแล้ว · จะรันซ้ำใส่ completed (หน้า Backfill ทางการ) ปลอดภัยเพราะบทที่ 12 · ระวัง: รอบที่กดเองพร้อมใส่วันที่ "จอง" วันนั้นไว้ แล้ว backfill ข้ามทั้งที่รอบนั้นได้ช่วงผิด (ยันในแล็บ)

ทำไมบทที่ 8 ใช้ params: วันที่ของรอบไม่ใช่เดือนที่ต้องประมวลผล และ cron สตริงในรุ่น 3 ไม่มีช่วงให้ด้วยซ้ำ

13.4 กับดักตอนย้ายจากรุ่น 2

โค้ดรุ่น 2 ที่อ่านช่วงข้อมูลจาก cron สตริง ย้ายมารุ่น 3 แล้วไม่พัง แต่ได้ช่วงยาว 0 · หน้า Upgrading ทางการให้สองทาง: บอก CronDataIntervalTimetable ให้ DAG เอง · หรือตั้ง [scheduler] create_cron_data_intervals = True ที่มีผลกับทุก DAG และถ้าพลิกหลังมีรอบของรุ่น 3 แล้ว รอบถัดไปข้ามไปหนึ่งงวด (แล็บเจองวดที่ไม่มีรอบรับจริง) · เล่มนี้แนะนำทางแรก · บทที่ 16 กลับมาที่โค้ดรุ่น 2

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

  1. ท่อรายวันของคุณสรุปยอดขายของ "เมื่อวาน" ตั้งให้รันตี 1 ทุกวัน · ถ้าใช้ CronDataIntervalTimetable("0 1 * * *") รอบที่รันตี 1 วันที่ 10 ได้ช่วงข้อมูลอะไร และงานควรสรุปของวันไหน
  2. เครื่องดับตั้งแต่ 28 กันยายนถึง 3 ตุลาคม · stock_monthly ฉบับนี้ (catchup ปิด) จะรันเดือนไหนบ้างเมื่อเปิดกลับมา และถ้ามันพังกลางทางวันนั้น คุณ clear แล้วรันซ้ำได้ปลอดภัยไหม — ใช้บทที่ 11 กับ 12 ตอบ
  3. เพื่อนแก้ month_of ให้ใช้ datetime.now() ลบหนึ่งเดือนแทนช่วงข้อมูล "จะได้ไม่ต้องพึ่งตารางเวลา" · backfill พฤษภาคมถึงกรกฎาคมที่สั่งวันที่ 7 ตุลาคมจะโหลดเดือนอะไร และคลังจะหน้าตาเป็นยังไง

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