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

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

LEVEL 4 · ระดับเทพ

เลือกเครื่อง⁠มือ · อ่านโค้ดรุ่น 2 · บทสรุป

บทสุดท้ายตอบคำถามของบทที่ 1 ว่างานแบบไหนไม่ต้องมีตัวคุม · แล้วอ่านโค้ดรุ่น 2 ที่ยังเต็มอินเทอร์เน็ต ให้รู้ว่าพังตรงไหน และตรงไหนไม่พังแต่ผิด

16.1 ต้องมีตัวคุมไหม — ถามความจริงสามข้อ

เลือกเครื่องมือด้วยความจริงคงที่สามข้อของเล่ม ไม่ใช่ด้วยตารางฟีเจอร์

  1. ความรู้ว่าทำถึงไหนต้องอยู่นอกตัวงาน (บทที่ 1) — งานของคุณพังกลางทางแล้วต้องรู้ไหมว่าทำถึงไหน · ถ้ามีขั้นเดียว หรือรันใหม่ทั้งก้อนแล้วไม่เสียอะไรเลย ความรู้นี้ไม่มีค่า
  2. ลำดับคือกราฟ (บทที่ 3) — มีขั้นที่รอขั้นอื่น และขั้นที่ไม่ต้องรอกันไหม · ถ้าเป็นเส้นตรงสามบรรทัด กราฟไม่ช่วยอะไร
  3. ทุกขั้นต้องรันซ้ำได้ (บทที่ 12) — ข้อนี้ไม่ใช่เกณฑ์เลือก เพราะไม่มีเครื่องมือไหนทำแทนคุณ · ใช้ Task Scheduler หรือ Airflow ก็ต้องทำเอง

ตอบ "ไม่" ทั้งข้อ 1 และ 2 = Task Scheduler หรือ cron พอ — คำตอบของหัวข้อ 1.5 ข้อ 2 · ตัวคุมเพิ่มของที่ต้องดูแล (ฐานข้อมูล · ห้าส่วนของบทที่ 15 · แรมราว 1 GB) โดยไม่ได้อะไรกลับมา

ตอบ "ใช่" = ต้องมีตารางสถานะนอกตัวงานกับวงรอบ · Prefect กับ Dagster อยู่กลุ่มเดียวกับ Airflow แต่เล่มนี้ไม่ได้ลอง จึงไม่เปรียบฟีเจอร์ให้ · ถามเอกสารของมันด้วยคำถามของเล่มแทน: สถานะของแต่ละขั้นเก็บที่ไหน · ประกาศ "อะไรรออะไร" ยังไง · clear แตะตารางยังไง · รอบรู้ช่วงข้อมูลจากไหน · ข้อมูลมาแล้วปลุกงานได้ไหม (คำตอบของ Airflow: บทที่ 4 · 3 · 11 · 13 · 14)

16.2 อ่านโค้ดรุ่น 2

Airflow 2 หมดการดูแล 22 เมษายน 2026 (หน้า version lifecycle ทางการ) แต่บทเรียนออนไลน์ส่วนใหญ่ยังเป็นรุ่น 2 · สัญญาณ: from airflow import DAG · airflow.decorators · airflow.operators.python · schedule_interval= · context["execution_date"] · days_ago(…)

ไฟล์นี้หน้าตาเหมือนบทเรียนรุ่น 2 ทั่วไป · วางใน legacy/ ที่ Airflow ไม่อ่าน

# ~/airflow-lab/legacy/stock_daily_v2.py — โค้ดแบบบทเรียนรุ่น 2 ที่ยังเจอได้ทั่วไป (ยังไม่แปลง)
from datetime import datetime

from airflow import DAG
from airflow.operators.python import PythonOperator


def report(**context):
    day = context["execution_date"].strftime("%Y-%m-%d")
    print(f"รายงานสต็อกของวันที่ {day}")


with DAG(
    dag_id="stock_daily_v2",
    schedule_interval="0 6 * * *",
    start_date=datetime(2026, 9, 1),
    catchup=False,
) as dag:
    PythonOperator(task_id="report", python_callable=report)

หน้า Upgrading ทางการแนะนำให้ตรวจด้วย ruff ที่มีกฎ AIR3 (ruff 0.13.1 ขึ้นไป · แล็บใช้ 0.16.10 ผ่าน uvx ซึ่งมากับ uv ของบทที่ 5)

cd ~/airflow-lab/legacy && uvx ruff@0.16.10 check stock_daily_v2.py --select AIR3 --preview --output-format concise

ผลรันจริง

stock_daily_v2.py:13:6: airflow3-suggested-update: `airflow.DAG` is removed in Airflow 3.0; It still works in Airflow 3.0 but is expected to be removed in a future version.
stock_daily_v2.py:15:5: airflow3-removal: [*] `schedule_interval` is removed in Airflow 3.0
stock_daily_v2.py:19:5: airflow3-suggested-to-move-to-provider: `airflow.operators.python.PythonOperator` is deprecated and moved into `standard` provider in Airflow 3.0; It still works in Airflow 3.0 but is expected to be removed in a future version.
Found 3 errors.
[*] 1 fixable with the `--fix` option (2 hidden fixes can be enabled with the `--unsafe-fixes` option).

ข้อที่พังจริงมีข้อเดียว (airflow3-removal) — รันใน 3.3.2 ได้ TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval' · อีกสองข้อยังทำงานแต่จะถูกถอดในอนาคต · แก้แค่ข้อที่พังแล้วรัน

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

$ airflow dags test stock_daily_v2 2026-10-07 2>&1 | grep KeyError
KeyError: 'execution_date'

ruff ไม่เห็น เพราะ execution_date เป็นแค่สตริงในคีย์ของ dict · หน้า Upgrading ระบุว่ารุ่น 3 ถอด execution_date · prev_ds · next_ds และคีย์ทำนองนี้ออกจาก context ⇒ ผ่าน ruff ไม่ได้แปลว่ารันได้

กับดักที่เงียบที่สุดอยู่ถัดไป · แก้ execution_date เป็น logical_date ตรง ๆ รันผ่าน แต่ได้วันผิดหนึ่งงวด: รุ่น 2 ให้ cron สตริงมีช่วงข้อมูล logical_date จึงเป็นต้นช่วง (เมื่อวาน) · รุ่น 3 ช่วงยาว 0 logical_date จึงเป็นเวลาที่รัน (วันนี้) — หัวข้อ 13.1 · หน้า Upgrading เตือนว่า ds และค่าที่คิดจาก logical_date "shift between the two timetables" · ฉบับแปลงจึงบอกชนิดตารางเวลาเองและอ่านช่วงข้อมูลตรง ๆ

# ~/airflow-lab/dags/stock_daily.py — แปลงเป็นรุ่น 3: import จาก airflow.sdk · บอกชนิดตารางเวลาเอง · อ่านช่วงข้อมูลแทน execution_date
from datetime import datetime

from airflow.sdk import CronDataIntervalTimetable, dag, task


@dag(
    schedule=CronDataIntervalTimetable("0 6 * * *", timezone="UTC"),  # ช่วงข้อมูลแบบรุ่น 2: รอบเช้านี้ = ข้อมูลของเมื่อวาน
    start_date=datetime(2026, 9, 1),
    catchup=False,
)
def stock_daily():
    @task
    def report(**context):
        start = context.get("data_interval_start")  # รอบที่กดรันเองไม่มีคีย์นี้ (บทที่ 13)
        if start is None:
            raise ValueError("รอบนี้ไม่มีช่วงข้อมูล")
        print(f"รายงานสต็อกของวันที่ {start:%Y-%m-%d}")

    report()


stock_daily()

ฉบับนี้ผ่าน ruff ชุดเดียวกัน (All checks passed!) และ dags test จบ success ในแล็บ

ลำดับแปลงทุกไฟล์: (1) ruff AIR3 · (2) ไล่คีย์ context ที่ถูกถอด · (3) DAG ที่ใช้ cron สตริงและอ่านช่วงข้อมูลหรือ ds — บอก CronDataIntervalTimetable เอง · (4) dags test ทีละไฟล์แล้วเทียบผลกับระบบเดิมหนึ่งงวด

16.3 บทสรุป

เล่มนี้เริ่มจากสคริปต์ที่พังตรงกลางแล้วไม่มีใครรู้ว่าทำถึงไหน · ทุกบทหลังจากนั้นคือตารางแผ่นเดียวของบทที่ 4 ในมุมต่างกัน: Grid คือตารางที่มองเห็น (7) · ด่านคือช่องที่พังก่อนข้อมูลผิดไหลลงปลายน้ำ (9) · retry · clear · backfill คือล้างหรือเพิ่มช่องแล้วให้วงรอบเดินต่อ (10 · 11 · 13) · Asset ปลุกโดยไม่ใช้นาฬิกา (14) · ห้าส่วนของ standalone คือคนเขียนตาราง (15) · และทั้งหมดใช้ได้เพราะทุกขั้นรันซ้ำแล้วได้ผลเท่าเดิม (12) ซึ่งเป็นงานของคุณ ไม่ใช่ของตัวคุม

สี่ข้อที่หน้าแรกสัญญา: ตัดสินว่าต้องมีตัวคุมไหม (16.1) · DAG ที่มีด่าน retry เฉพาะความพังชั่วคราว และรันซ้ำได้ (8–12) · อ่าน Grid ตอนเช้าแล้วสั่งรันเฉพาะส่วนที่ต้องรัน (7 · 11) · อ่านโค้ดรุ่น 2 แล้วแปลง (16.2)

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

  1. งานของคุณ: ทุกเช้าดาวน์โหลดไฟล์หนึ่งไฟล์จากเว็บไซต์ราชการ แล้วคัดลอกไปไว้ในโฟลเดอร์ร่วม · ถ้าพังก็แค่สั่งใหม่ · ต้องมีตัวคุมไหม และถ้าวันหนึ่งงานนี้เพิ่มขั้น "แปลงไฟล์ → โหลดเข้าฐานข้อมูล → ส่งเมลสรุป" คำตอบเปลี่ยนไหม
  2. โค้ดรุ่น 2 ใช้ schedule_interval="@daily" และอ่าน context["ds"] เพื่อดึงยอดขายของวันนั้น · ย้ายมารุ่น 3 แล้วแก้แค่ schedule_interval → schedule · รันผ่านทุกวันไม่มี error · มีอะไรผิด และแก้ยังไง — ใช้บทที่ 13 ตอบ
  3. ทีมคุณกำลังเลือกระหว่าง Airflow กับเครื่องมืออื่นที่เล่มนี้ไม่ได้สอน · เขียนคำถามสามข้อที่จะถามเอกสารของเครื่องมือนั้น โดยแต่ละข้อมาจากคนละบทของเล่ม

เฉลยคำถามท้ายบท#

บทที่ 1

  1. ลูกค้า 80 คนแรกได้อีเมลซ้ำเป็นฉบับที่สอง · ความรู้ที่หายไปคือ "ส่งถึงใครไปแล้วบ้าง" ซึ่งอยู่ในตัวแปรของลูป — ไม่มีใครจดไว้นอกโปรเซส จึงไม่มีทางเริ่มที่คนที่ 81 ได้ นอกจากไปไล่ดูกล่องอีเมลขาออกเอง
  2. รหัสจบการทำงานบอกแค่ว่าโปรเซสจบโดยไม่มี error ที่ไม่มีใครจับ ไม่ได้บอกว่าข้อมูลที่ได้ถูกต้อง · เดือนกรกฎาคมทุกขั้นทำงานตามที่เขียนไว้ทุกบรรทัด แต่ขั้นอ่านไฟล์ได้ศูนย์แถว และไม่มีขั้นไหนถามว่าศูนย์แถวสมเหตุสมผลไหม
  3. เลขแถวทางซ้ายกระโดด (4 ไป 38) · ปุ่มที่หัวคอลัมน์เป็นรูปกรวย · แถบสถานะขึ้น Filter Mode — เจอข้อใดข้อหนึ่ง = มีแถวถูกซ่อนด้วยตัวกรอง ล้างตัวกรองหรืออ่านไฟล์ตรง ๆ ก่อนสรุป

บทที่ 2

  1. รันใหม่แล้วไม่มีอะไรถูกจดไว้เลย ตัวรันเริ่มขั้นที่ 1 ใหม่ทั้งหมด — อาการเดียวกับสคริปต์ก้อนเดียว · ที่ขาดคือการจดทันทีที่แต่ละขั้นเสร็จ
  2. ตัวรันอ่าน state.json เจอว่าทุกขั้นสำเร็จแล้ว จึงข้ามหมด เดือนมิถุนายนไม่ถูกทำเลย และไม่มี error ให้เห็น · ต้นเหตุคือไฟล์สถานะไม่รู้จักคำว่า "รอบ" — บทที่ 4 แก้ด้วยการให้แต่ละรอบเป็นคอลัมน์ของตัวเอง
  3. ทางหนึ่งที่ทำได้: ในบล็อก except ตรงที่ตรวจ cell["tries"] > RETRIES ให้เขียนไฟล์ก่อน return เช่น (WORK / "ALERT.txt").write_text(f"{name} ลองครบ {cell['tries']} ครั้งแล้วยังพัง: {e}", encoding="utf-8") · ของจริงเปลี่ยนการเขียนไฟล์เป็นส่งข้อความหาคน แต่ตำแหน่งในโค้ดเป็นจุดเดียวกัน — จุดที่ตัวรันรู้แน่แล้วว่ายอมแพ้

บทที่ 3

  1. ขั้นที่ไม่รอใครเลยคือดึงรายการสั่งซื้อ ดึงรายชื่อสินค้า และดึงรายชื่อสาขา จึงอยู่จังหวะแรกทั้งสามขั้น · จับคู่กับสินค้ารอสองขั้นแรก (จังหวะที่ 2) · จับคู่กับสาขารอผลจับคู่กับสินค้าและรายชื่อสาขา (จังหวะที่ 3) · สรุปรายสาขา (จังหวะที่ 4) · ส่งรายงาน (จังหวะที่ 5) — ห้าจังหวะ
  2. ไม่มีเลย notify เป็นขั้นปลายสุด ไม่มีใครรอมัน · ความพังไหลลงปลายน้ำเท่านั้น และขั้นนี้ไม่มีปลายน้ำ
  3. เพราะขั้นที่อยู่ในวงไม่มีวันพร้อม รันส่วนที่เหลือไปก่อนก็ได้ท่อที่จบไม่ได้ทุกรอบ และรู้ตัวช้าไปหนึ่งคืนเสมอ · เจอตอนประกาศกราฟคือการบอกความผิดพลาด ณ จุดที่มันเกิด — ตอนแก้ NEEDS — ก่อนจะกลายเป็นรอบที่ค้างอยู่กลางทางทุกคืน

บทที่ 4

  1. จังหวะที่ 1: receive_c พังครั้งที่ 2 หมดสิทธิ์ → failed 2 · จังหวะที่ 2: merge_check → upstream_failed 0 · จังหวะที่ 3: transform → upstream_failed 0 · (ต่ออีกสองจังหวะ load กับ notify ก็เป็น upstream_failed 0) · รอบจบเป็น failed เพราะขั้นปลายสุด notify เป็น upstream_failed
  2. ข้อ "ทำถึงไหนแล้ว" — สถานะอยู่ในโปรเซส ตายแล้วหายตาม · จะรู้ตัวในคืนแรกที่ตัวรันตายกลางงาน แล้วเปิดใหม่ต้องเริ่มทั้งรอบ ไม่ใช่ตอนที่ทุกอย่างปกติดี
  3. failed · ท่อนี้มีขั้นปลายสุดสองขั้น archive สำเร็จ (receive ทั้งสามผ่าน) แต่ notify เป็น upstream_failed · กติกาคือมีขั้นปลายสุดตัวไหน failed หรือ upstream_failed รอบก็เป็น failed — ขั้นปลายสุดที่ผ่านไม่ได้ช่วยกลบตัวที่พัง

บทที่ 5

  1. ใช้ 3.13 ของเครื่อง: เปลี่ยนสองที่ให้ตรงกัน — uv venv --python 3.13 กับ URL เป็น constraints-3.13.txt · เรื่อง sudo ไม่ต้องแก้ uv ลงในโฟลเดอร์บ้านอยู่แล้ว · อยากได้ผลตรงเล่มทุกบรรทัด: ใช้ --python 3.12 ตามเล่ม uv ดาวน์โหลดมาให้โดยไม่แตะของเครื่อง
  2. ในไฟล์ SQLite ~/airflow-lab/airflow.db · ไม่หาย — ตารางอยู่บนดิสก์ ไม่ได้อยู่ในโปรเซสของ standalone · แล็บของเล่มนี้ปิดเปิด standalone มาหลายรอบ รอบของ hello_pipeline ตั้งแต่วันแรกยังอยู่ครบ
  3. ปุ่มบนหน้าจอคือคำสั่งแก้ตาราง ไม่ใช่แค่ที่ดู — Trigger เพิ่มคอลัมน์ · Clear ล้างช่องให้รันใหม่ · Delete Dag ลบ DAG ออกจากระบบ · ถ้าท่อเขียนข้อมูลจริง ใครที่เข้าถึงพอร์ตได้ก็สั่งรันซ้ำหรือล้างสถานะได้ (บทที่ 4: ทุกคำสั่งคือการแก้ตาราง)

บทที่ 6

  1. บรรทัดชั้นบนสุดรันทุกครั้งที่ไฟล์ถูกอ่าน — อ่าน Excel 20 วินาทีทุกราว 30 วินาที ทั้งวันทั้งคืน แม้ไม่มีรอบ · ย้ายเข้าไปเป็นเนื้อในของ @task ขั้นรับไฟล์ แล้วส่งใบสรุปต่อแบบบทที่ 8
  2. ยังเป็นแถว — มันถูกเรียกในรายการ counts จึงถูกประกาศ แม้ไม่มีใครรับค่า · ไม่รอใคร ไม่มีใครรอมัน ⇒ กลายเป็นขั้นปลายสุดอีกขั้น · receive_c พัง = ขั้นปลายสุดตัวหนึ่ง failed ⇒ รอบ failed ทั้งที่ notify ส่งรายงานยอด 460 ชิ้นไปแล้ว (ยันในแล็บ: ขั้นอื่น success ทั้งหมด · รอบจบ failed)
  3. merge_check รับรายการที่ประกาศไว้ตอนอ่านไฟล์ — ช่องที่หนึ่งคือใบสั่งหยิบ XCom ของ receive_a (หัวข้อ 6.3) ช่องที่สองของ receive_b ช่องที่สามของ receive_c · ลำดับที่สามขั้นเสร็จไม่มีความหมาย เพราะอยู่จังหวะเดียวกัน (บทที่ 3) และ merge_check ไม่เริ่มจนกว่าทั้งสามจะ success ครบ (บทที่ 4)

บทที่ 7

  1. แดง (failed) — ขั้นปลายสุด notify เป็น upstream_failed และไม่มี all_done มากลบ · พังที่ merge_check ช่องส้มสามช่องคือผลตาม · เปิด merge_check ก่อน ดู Try Number แล้วหาบรรทัด ERROR ของครั้งสุดท้าย
  2. trigger rule ต่างกัน · load ใช้ค่าตั้งต้น all_success — transform failed จึงเป็น upstream_failed (กติกาข้อ 2 ของวงรอบในบทที่ 4) · notify ตั้ง all_done ขอแค่ต้นน้ำจบ — load จบแบบ upstream_failed ก็นับว่าจบ notify จึงรัน
  3. สี่ปุ่ม — ครั้งแรกบวกสิทธิ์ลองใหม่สามครั้ง · ช่วงรอเพิ่มจากหนึ่งเป็นสามช่วง ช่วงละ 10 วินาที ⇒ นานขึ้นอย่างน้อย 20 วินาที · และไม่ได้อะไรเลย คอลัมน์หายคือความพังถาวร (คำเตือนท้ายหัวข้อ 2.2 · บทที่ 10)

บทที่ 8

  1. รับไฟล์: อ่านไฟล์ของวันนั้น เขียนแถวลงที่พักข้อมูลของคุณ คืนใบสรุป (path · จำนวนแถว · คอลัมน์ที่ขาด) · ตรวจ: โครงผิดยก error ไม่งั้นคืนใบสรุปเดิม · โหลด: อ่านแถวจาก path เขียนลงปลายทาง คืนจำนวนแถว · สิ่งที่ไม่ควรผ่าน XCom คือแถวข้อมูลทั้งก้อน
  2. ท่อรันผ่านเพราะ 33 แถวยังเล็ก · ที่เสียคือแถวทั้งหมดถูกเขียนลงฐานข้อมูลของ Airflow ทุกขั้นทุกรอบ ฐานข้อมูลเดียวกับที่เก็บสถานะของทุกช่อง ไฟล์ยิ่งใหญ่ยิ่งถ่วงทั้งระบบ · "ถูก" ตามบทที่ 2 เพราะข้อมูลอยู่นอกโปรเซสจริง แต่ผิดตามบทนี้ — ข้อมูลควรอยู่ในที่เก็บของคุณ ไม่ใช่ในฐานข้อมูลของตัวคุม
  3. เป็น ในคลังเดือนนี้ 198 แถว ขณะที่ยอดในบรรทัดเดียวกันยังเป็น 951 ชิ้น เพราะยอดคิดจากรอบนี้รอบเดียว — สองตัวเลขในบรรทัดเดียวขัดกันเอง (ยันในแล็บ: ได้ 198 จริง) · load ต่อท้าย รันซ้ำจึงเบิ้ล อาการเดียวกับเช้าวันที่ 1 มิถุนายนในบทที่ 1 ที่ยอดกลายเป็น 1,480 · ทางแก้คือบทที่ 12

บทที่ 9

  1. วันอาทิตย์ไม่มีไฟล์ = skipped ไม่มีใครต้องลงมือ · ไฟล์ไม่มีคอลัมน์ยอดขาย = failed ขั้นถัดไปใช้ไม่ได้และต้องมีคนไปขอไฟล์ใหม่ · ร้านใหม่ = ปล่อยผ่าน ร้านเป็นค่าที่ต้นทางเป็นเจ้าของ ไม่ใช่โครงของไฟล์
  2. เขียว — gate เป็น skipped แล้วไหลลงปลายน้ำเป็น skipped ทั้งหมด ขั้นปลายสุดที่ข้ามนับเป็นผ่าน (หัวข้อ 9.3 รอบเดือนตุลาคมคือกลไกเดียวกัน) · ไม่มีใครรู้ — ไม่มีอะไรโหลด ไม่มีรายงาน และไม่มีรอบแดงให้ใครเปิดดู จนสิ้นเดือนมีคนถามหายอดเดือนสิงหาคม
  3. ผ่าน — ตัวอ่านของบทที่ 8 อ่านทุกแถว ได้ 33 แถว คอลัมน์ครบ วันที่ถูกเดือน ส่วนแถวซ่อนเป็นเรื่องของคำเตือนในรายงาน · ถ้าใช้ตัวอ่านเก่าที่ข้ามแถวซ่อน จะได้ 0 แถว ด่านหยุดท่อด้วย "มี 0 แถว นอกช่วง 10–200" ก่อนถึงขั้นโหลด — เช้าวันที่ 1 สิงหาคมของบทที่ 1 จะเป็นรอบแดงแทนรายงาน 0 ชิ้น

บทที่ 10

  1. หมดเวลารอ กับส่งถี่เกิน = ลองใหม่ (รอแล้วหาย · แบบหลังควรรอนานขึ้น) · รหัสเข้าใช้หมดอายุ กับข้อมูลขาดฟิลด์บังคับ = หยุดทันที รอเท่าไหร่ก็ไม่หาย ต้องมีคนแก้ — หน้า Tasks ทางการยกตัวอย่างรหัสเข้าใช้ไม่ถูกต้องไว้เป็นกรณีที่ควรพังทันทีตรงตัว
  2. จังหวะถัดไป gate เป็น upstream_failed 0 · แล้ว transform · load · notify ตามลงทีละจังหวะ (กติกาข้อ 2 ของวงรอบ) · รอบแดงเพราะขั้นปลายสุด notify เป็น upstream_failed และไม่มี all_done มากลบ — ตรงกับผลของ states.py ในหัวข้อ 10.2 (watcher จำเป็นเฉพาะท่อที่มีขั้นปลายสุดแบบ all_done)
  3. ในบล็อก except ของตัวรัน — ตรงที่ตัดสินว่า up_for_retry หรือ failed · เพิ่มเงื่อนไขว่า error ชนิดที่รู้ว่าถาวร ให้เป็น failed ทันทีโดยไม่ดูจำนวนครั้ง · ชิ้นนั้นคือ AirflowFailException (หรือกฎ FAIL ของ retry policy ตั้งแต่ 3.3) — ส่วนการแยกว่า error ไหนถาวร ยังเป็นงานของโค้ดขั้นนั้นเหมือนเดิม

บทที่ 11

  1. clear receive_c พร้อม Downstream — ช่องที่อ่านไฟล์ที่เปลี่ยน · หน้าต่างบอก 7 ช่องเหมือนเดิม: receive_c กับปลายน้ำทั้งหกขั้น (gate ถึง watcher) ส่วน receive_a receive_b arrived ไม่ถูกแตะ
  2. ครั้งที่ 2 — gate จับ BadFile แล้วยก AirflowFailException ซึ่งหยุดทันทีไม่สนสิทธิ์ที่เหลือ จึงไม่ได้ใช้เพดาน 3 · ช่องจบ failed 2 และ watcher ทำให้รอบแดงอีกครั้ง
  3. ตัวนับครั้งที่ลองนับต่อ ไม่ถูกล้างตอน clear · watcher เคยรันมาแล้ว 1 ครั้ง (พังในรอบแรก) หลัง clear เป็น none 1 · รอบนี้ไม่มีขั้นไหนพัง วงรอบจึงข้ามมันโดยไม่ได้รัน ตัวนับไม่ขยับ = skipped 1 · ช่องในบทที่ 9 ไม่เคยรันเลยตั้งแต่ต้น จึงเป็น skipped 0

บทที่ 12

  1. ไม่ได้ — อีเมลที่ส่งแล้วเรียกคืนไม่ได้ การส่งครั้งที่สองเปลี่ยนโลกข้างนอกเสมอ · กันได้ด้วยการจดลงที่เก็บของคุณว่า "ส่งยอดของเดือน 2026-07 แล้ว" ในธุรกรรมเดียวกับการตัดสินใจส่ง แล้วเช็กก่อนส่งทุกครั้ง · หรือยอมให้ส่งซ้ำได้แต่หัวเรื่องบอกเดือนและบอกว่าเป็นฉบับแก้ ให้ผู้จัดการรู้ว่าฉบับไหนใหม่กว่า
  2. load — ได้ เพราะรันซ้ำแล้วคลังเท่าเดิม และ error ที่ load เจอบ่อยคือฐานข้อมูลสะดุด ซึ่งรอแล้วหาย (บทที่ 10) · notify — ระวัง ทุกครั้งที่ลองใหม่หลังส่งไปแล้วคือข้อความเพิ่มหนึ่งฉบับ ควรคง 0 หรือกันส่งซ้ำแบบข้อ 1 ก่อน
  3. 1,020 ชิ้น — สาขา ก กับ ข ถูกเขียนทับด้วยของเดือนเดียวกัน ไม่ใช่ต่อท้าย (ตารางของบทที่ 1 ต้องมีคอลัมน์เดือนก่อน ถึงจะลบเฉพาะของเดือนนั้นได้) · ที่ยังขาด: ทำถึงไหนแล้ว · อะไรรออะไร · ใครรู้ว่ามันพัง — ยังอยู่ในโปรเซสเหมือนเดิม การรันซ้ำปลอดภัยแล้ว แต่คืนวันที่ 1 มิถุนายนยังต้องรอคนตื่นมาเห็นรหัสจบการทำงาน 1 และรันใหม่ทั้งก้อนเอง

บทที่ 13

  1. ช่วง 01:00 วันที่ 9 → 01:00 วันที่ 10 — รอบเกิดตอนปลายช่วง และ logical_date คือต้นช่วง · งานควรสรุปของวันที่ 9 โดยคิดจาก data_interval_start
  2. รอบเดียว คือกันยายน งวดล่าสุดที่จบแล้ว (แล็บ unpause วันที่ 7 ตุลาคมได้กันยายนรอบเดียว) · ถ้าดับข้ามหลายต้นเดือน เดือนที่ข้ามไปต้อง backfill เอง · clear ได้ปลอดภัย: clear รันช่องเดิมซ้ำ (บทที่ 11) ทุกขั้นคิดเดือนจากช่วงของรอบ และโหลดเขียนทับทั้งเดือน (บทที่ 12) — รันซ้ำวันไหนก็ได้กันยายน
  3. ทั้งสามรอบรันวันที่ 7 ตุลาคม จึงได้เดือนกันยายนหมด · พฤษภาคมถึงกรกฎาคมไม่ถูกโหลดเลย · ทั้งสามรอบพยายามประมวลผลกันยายนเดือนเดียวกัน (ถ้าไฟล์ครบ คลังยังเป็นชุดเดียวเพราะเขียนทับทั้งเดือน) · Grid โชว์สามคอลัมน์ของพฤษภาคมถึงกรกฎาคม แต่ไม่มีอะไรบอกว่าข้างในคือกันยายน — now() ทำให้ backfill เสียความหมายทั้งหมด

บทที่ 14

  1. ไฟล์ยอดขายรายวัน = นาฬิกา มาตรงเวลาเสมอ ไม่มีอะไรต้องรอ · ไฟล์ราคา = Asset ระบบนั้นเรียก API บอกได้ว่ามาแล้ว จึงไม่ต้องเดาเวลา · ถ้าระบบราคาเรียก API ไม่ได้ ค่อยถอยไปใช้ sensor แบบ reschedule
  2. ผู้ผลิตประกาศเดือน 2026-09 ของสาขา ค ทั้งที่ไม่มีไฟล์ · ครบสามประกาศ stock_monthly ถูกปลุก · arrived ผ่าน (มีไฟล์สองสาขา) · receive_c ยก FileNotFoundError → up_for_retry สองครั้ง → failed 3 (บทที่ 10) · gate ถึง notify เป็น upstream_failed · morning ผ่าน · watcher พัง รอบแดง · ประกาศที่ไม่มีของจริงย้ายความพังจากผู้ผลิตไปไว้ที่ผู้ใช้ และเสียเวลารอ retries เปล่า — ผู้ผลิตยันก่อนประกาศคือด่านแบบบทที่ 9 ที่ต้นทาง
  3. ถ้ารอบใหม่เกิด ท่อโหลดกันยายนซ้ำด้วยไฟล์ฉบับแก้ของสาขา ข · load ลบของเดือนนั้นแล้วเขียนใหม่ทั้งเดือน คลังจึงเป็นตัวเลขฉบับแก้ 99 แถว ไม่เบิ้ล (บทที่ 12) · สิ่งเดียวที่ซ้ำคือข้อความรายงาน (notify · morning) ซึ่งต้องกันส่งซ้ำแบบหัวข้อ 12.5

บทที่ 15

  1. queued = scheduler ส่งให้ executor แล้วแต่ยังไม่มี worker รับ — สงสัย executor/worker ก่อน (ในแล็บคือ worker ของ LocalExecutor ที่เป็นลูกของ scheduler) และเพดานอย่าง parallelism · dag-processor ไม่น่าใช่ เพราะถึง queued แล้วแปลว่ากราฟถูกอ่านและ scheduler ตัดสินแล้ว · task_queued_timeout (ตั้งต้น 600 วินาที ตามหน้า config reference) จัดการช่องที่ค้างนานเกิน
  2. ไม่ — heartbeat บอกว่าคนรันยังมีชีวิต ไม่ได้ดูว่างานคืบหน้า · ต้องตั้ง execution_timeout ที่ขั้นนั้น · แล็บลองงานค้าง 30 วินาทีกับ execution_timeout 5 วินาที ได้ AirflowTaskTimeout ทุกครั้ง แล้ว retries ลองซ้ำตามปกติ (failed 2 ที่ retries 1) — API ที่ค้างเป็นพัก ๆ จึงได้ประโยชน์จากการลองใหม่ (บทที่ 10)
  3. worker ว่างอยู่ 31 ตัวตลอดรอบอยู่แล้ว ช่องว่างไม่ได้มาจากการรอ worker · มันอยู่ระหว่างต้นน้ำจบกับ scheduled คือรอวงรอบของ scheduler หมุนมาเห็น · แล็บย่อช่องว่างได้ด้วย scheduler_idle_sleep_time ไม่ใช่ด้วยการเพิ่ม worker

บทที่ 16

  1. ไม่ต้อง — ขั้นเดียว สั่งใหม่ทั้งก้อนไม่เสียอะไร Task Scheduler พอ · สี่ขั้นแล้วคำตอบเปลี่ยน: พังที่โหลดแล้วไม่อยากดาวน์โหลดใหม่ และส่งเมลต้องรอโหลดเสร็จ — ความรู้ว่าทำถึงไหนเริ่มมีค่า · เลือกอะไรก็ตาม ขั้นโหลดต้องรันซ้ำได้และขั้นส่งเมลต้องกันส่งซ้ำ (บทที่ 12)
  2. @daily ในรุ่น 3 ก็เป็น cron สตริง ได้ CronTriggerTimetable ช่วงยาว 0 (ยันในแล็บ) · ds คิดจาก logical_date ซึ่งตอนนี้คือวันที่รัน ไม่ใช่เมื่อวาน ⇒ ทุกรอบดึงยอดของวันนี้ที่ยังไม่จบแทนเมื่อวาน ได้ตัวเลขครึ่งวันโดยไม่มี error · แก้ด้วย schedule=CronDataIntervalTimetable("0 0 * * *", timezone=…) ให้ DAG นั้น แล้วอ่านช่วงจาก data_interval_start (บทที่ 13)
  3. ตัวอย่างหนึ่งชุด: สถานะของแต่ละขั้นของแต่ละรอบเก็บที่ไหน และโปรเซสที่ตายกลางงานแล้วเปิดใหม่ทำต่อได้ไหม (บทที่ 4) · clear หรือรันซ้ำแตะช่องไหน ตัวนับครั้งนับต่อไหม (บทที่ 11) · รอบรู้ได้ยังไงว่าตัวเองประมวลผลช่วงไหน และย้อนรันช่วงเก่ายังไง (บทที่ 13)

อ่านต่อ: System Design#

เล่มนี้จบที่ท่อหนึ่งท่อบนเครื่องเดียว และบทเรียนที่หนักที่สุดคือบทที่ 12: ระบบที่ลองใหม่ได้ บังคับให้ทุกขั้นรันซ้ำแล้วได้ผลเท่าเดิม · พอท่อเริ่มคุยกับบริการอื่นผ่านคิวและ API คำถามเดียวกันกลับมาในอีกรูป — คิวส่วนใหญ่ส่งข้อความซ้ำได้ ผู้รับจึงต้องรับซ้ำแล้วได้ผลเดิม · และข้อความที่พังถาวรต้องมีที่ไปแทนการลองใหม่ไม่รู้จบ (บทที่ 10 ในโลกของคิว)

System Design Handbook — ไปต่อที่ บทที่ 09 · Message Queue & Async Communication ซึ่งพา retry กับ idempotency ออกจากท่อข้อมูลไปสู่ระบบที่หลายบริการคุยกัน

ถ้าเนื้อหานี้มีประโยชน์ —เลี้ยงกาแฟสักแก้วหรือพร้อมเพย์ · Lightning

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