← คู่มือ Airflow — ทำไมงานหลายขั้นต้องมีคนคุมลำดับ
LEVEL 4 · ระดับเทพ
เลือกเครื่องมือ · อ่านโค้ดรุ่น 2 · บทสรุป
บทสุดท้ายตอบคำถามของบทที่ 1 ว่างานแบบไหนไม่ต้องมีตัวคุม · แล้วอ่านโค้ดรุ่น 2 ที่ยังเต็มอินเทอร์เน็ต ให้รู้ว่าพังตรงไหน และตรงไหนไม่พังแต่ผิด
16.1 ต้องมีตัวคุมไหม — ถามความจริงสามข้อ
เลือกเครื่องมือด้วยความจริงคงที่สามข้อของเล่ม ไม่ใช่ด้วยตารางฟีเจอร์
- ความรู้ว่าทำถึงไหนต้องอยู่นอกตัวงาน (บทที่ 1) — งานของคุณพังกลางทางแล้วต้องรู้ไหมว่าทำถึงไหน · ถ้ามีขั้นเดียว หรือรันใหม่ทั้งก้อนแล้วไม่เสียอะไรเลย ความรู้นี้ไม่มีค่า
- ลำดับคือกราฟ (บทที่ 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)
ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)
- งานของคุณ: ทุกเช้าดาวน์โหลดไฟล์หนึ่งไฟล์จากเว็บไซต์ราชการ แล้วคัดลอกไปไว้ในโฟลเดอร์ร่วม · ถ้าพังก็แค่สั่งใหม่ · ต้องมีตัวคุมไหม และถ้าวันหนึ่งงานนี้เพิ่มขั้น "แปลงไฟล์ → โหลดเข้าฐานข้อมูล → ส่งเมลสรุป" คำตอบเปลี่ยนไหม
- โค้ดรุ่น 2 ใช้
schedule_interval="@daily"และอ่านcontext["ds"]เพื่อดึงยอดขายของวันนั้น · ย้ายมารุ่น 3 แล้วแก้แค่schedule_interval→schedule· รันผ่านทุกวันไม่มี error · มีอะไรผิด และแก้ยังไง — ใช้บทที่ 13 ตอบ - ทีมคุณกำลังเลือกระหว่าง Airflow กับเครื่องมืออื่นที่เล่มนี้ไม่ได้สอน · เขียนคำถามสามข้อที่จะถามเอกสารของเครื่องมือนั้น โดยแต่ละข้อมาจากคนละบทของเล่ม
เฉลยคำถามท้ายบท#
บทที่ 1
- ลูกค้า 80 คนแรกได้อีเมลซ้ำเป็นฉบับที่สอง · ความรู้ที่หายไปคือ "ส่งถึงใครไปแล้วบ้าง" ซึ่งอยู่ในตัวแปรของลูป — ไม่มีใครจดไว้นอกโปรเซส จึงไม่มีทางเริ่มที่คนที่ 81 ได้ นอกจากไปไล่ดูกล่องอีเมลขาออกเอง
- รหัสจบการทำงานบอกแค่ว่าโปรเซสจบโดยไม่มี error ที่ไม่มีใครจับ ไม่ได้บอกว่าข้อมูลที่ได้ถูกต้อง · เดือนกรกฎาคมทุกขั้นทำงานตามที่เขียนไว้ทุกบรรทัด แต่ขั้นอ่านไฟล์ได้ศูนย์แถว และไม่มีขั้นไหนถามว่าศูนย์แถวสมเหตุสมผลไหม
- เลขแถวทางซ้ายกระโดด (4 ไป 38) · ปุ่มที่หัวคอลัมน์เป็นรูปกรวย · แถบสถานะขึ้น Filter Mode — เจอข้อใดข้อหนึ่ง = มีแถวถูกซ่อนด้วยตัวกรอง ล้างตัวกรองหรืออ่านไฟล์ตรง ๆ ก่อนสรุป
บทที่ 2
- รันใหม่แล้วไม่มีอะไรถูกจดไว้เลย ตัวรันเริ่มขั้นที่ 1 ใหม่ทั้งหมด — อาการเดียวกับสคริปต์ก้อนเดียว · ที่ขาดคือการจดทันทีที่แต่ละขั้นเสร็จ
- ตัวรันอ่าน
state.jsonเจอว่าทุกขั้นสำเร็จแล้ว จึงข้ามหมด เดือนมิถุนายนไม่ถูกทำเลย และไม่มี error ให้เห็น · ต้นเหตุคือไฟล์สถานะไม่รู้จักคำว่า "รอบ" — บทที่ 4 แก้ด้วยการให้แต่ละรอบเป็นคอลัมน์ของตัวเอง - ทางหนึ่งที่ทำได้: ในบล็อก
exceptตรงที่ตรวจcell["tries"] > RETRIESให้เขียนไฟล์ก่อนreturnเช่น(WORK / "ALERT.txt").write_text(f"{name} ลองครบ {cell['tries']} ครั้งแล้วยังพัง: {e}", encoding="utf-8")· ของจริงเปลี่ยนการเขียนไฟล์เป็นส่งข้อความหาคน แต่ตำแหน่งในโค้ดเป็นจุดเดียวกัน — จุดที่ตัวรันรู้แน่แล้วว่ายอมแพ้
บทที่ 3
- ขั้นที่ไม่รอใครเลยคือดึงรายการสั่งซื้อ ดึงรายชื่อสินค้า และดึงรายชื่อสาขา จึงอยู่จังหวะแรกทั้งสามขั้น · จับคู่กับสินค้ารอสองขั้นแรก (จังหวะที่ 2) · จับคู่กับสาขารอผลจับคู่กับสินค้าและรายชื่อสาขา (จังหวะที่ 3) · สรุปรายสาขา (จังหวะที่ 4) · ส่งรายงาน (จังหวะที่ 5) — ห้าจังหวะ
- ไม่มีเลย
notifyเป็นขั้นปลายสุด ไม่มีใครรอมัน · ความพังไหลลงปลายน้ำเท่านั้น และขั้นนี้ไม่มีปลายน้ำ - เพราะขั้นที่อยู่ในวงไม่มีวันพร้อม รันส่วนที่เหลือไปก่อนก็ได้ท่อที่จบไม่ได้ทุกรอบ และรู้ตัวช้าไปหนึ่งคืนเสมอ · เจอตอนประกาศกราฟคือการบอกความผิดพลาด ณ จุดที่มันเกิด — ตอนแก้
NEEDS— ก่อนจะกลายเป็นรอบที่ค้างอยู่กลางทางทุกคืน
บทที่ 4
- จังหวะที่ 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 - ข้อ "ทำถึงไหนแล้ว" — สถานะอยู่ในโปรเซส ตายแล้วหายตาม · จะรู้ตัวในคืนแรกที่ตัวรันตายกลางงาน แล้วเปิดใหม่ต้องเริ่มทั้งรอบ ไม่ใช่ตอนที่ทุกอย่างปกติดี
failed· ท่อนี้มีขั้นปลายสุดสองขั้นarchiveสำเร็จ (receive ทั้งสามผ่าน) แต่notifyเป็นupstream_failed· กติกาคือมีขั้นปลายสุดตัวไหนfailedหรือupstream_failedรอบก็เป็นfailed— ขั้นปลายสุดที่ผ่านไม่ได้ช่วยกลบตัวที่พัง
บทที่ 5
- ใช้ 3.13 ของเครื่อง: เปลี่ยนสองที่ให้ตรงกัน —
uv venv --python 3.13กับ URL เป็นconstraints-3.13.txt· เรื่อง sudo ไม่ต้องแก้ uv ลงในโฟลเดอร์บ้านอยู่แล้ว · อยากได้ผลตรงเล่มทุกบรรทัด: ใช้--python 3.12ตามเล่ม uv ดาวน์โหลดมาให้โดยไม่แตะของเครื่อง - ในไฟล์ SQLite
~/airflow-lab/airflow.db· ไม่หาย — ตารางอยู่บนดิสก์ ไม่ได้อยู่ในโปรเซสของ standalone · แล็บของเล่มนี้ปิดเปิด standalone มาหลายรอบ รอบของhello_pipelineตั้งแต่วันแรกยังอยู่ครบ - ปุ่มบนหน้าจอคือคำสั่งแก้ตาราง ไม่ใช่แค่ที่ดู — Trigger เพิ่มคอลัมน์ · Clear ล้างช่องให้รันใหม่ · Delete Dag ลบ DAG ออกจากระบบ · ถ้าท่อเขียนข้อมูลจริง ใครที่เข้าถึงพอร์ตได้ก็สั่งรันซ้ำหรือล้างสถานะได้ (บทที่ 4: ทุกคำสั่งคือการแก้ตาราง)
บทที่ 6
- บรรทัดชั้นบนสุดรันทุกครั้งที่ไฟล์ถูกอ่าน — อ่าน Excel 20 วินาทีทุกราว 30 วินาที ทั้งวันทั้งคืน แม้ไม่มีรอบ · ย้ายเข้าไปเป็นเนื้อในของ
@taskขั้นรับไฟล์ แล้วส่งใบสรุปต่อแบบบทที่ 8 - ยังเป็นแถว — มันถูกเรียกในรายการ
countsจึงถูกประกาศ แม้ไม่มีใครรับค่า · ไม่รอใคร ไม่มีใครรอมัน ⇒ กลายเป็นขั้นปลายสุดอีกขั้น ·receive_cพัง = ขั้นปลายสุดตัวหนึ่งfailed⇒ รอบfailedทั้งที่notifyส่งรายงานยอด 460 ชิ้นไปแล้ว (ยันในแล็บ: ขั้นอื่นsuccessทั้งหมด · รอบจบfailed) merge_checkรับรายการที่ประกาศไว้ตอนอ่านไฟล์ — ช่องที่หนึ่งคือใบสั่งหยิบ XCom ของreceive_a(หัวข้อ 6.3) ช่องที่สองของreceive_bช่องที่สามของreceive_c· ลำดับที่สามขั้นเสร็จไม่มีความหมาย เพราะอยู่จังหวะเดียวกัน (บทที่ 3) และmerge_checkไม่เริ่มจนกว่าทั้งสามจะsuccessครบ (บทที่ 4)
บทที่ 7
- แดง (
failed) — ขั้นปลายสุดnotifyเป็นupstream_failedและไม่มีall_doneมากลบ · พังที่merge_checkช่องส้มสามช่องคือผลตาม · เปิดmerge_checkก่อน ดู Try Number แล้วหาบรรทัด ERROR ของครั้งสุดท้าย - trigger rule ต่างกัน ·
loadใช้ค่าตั้งต้นall_success—transformfailedจึงเป็นupstream_failed(กติกาข้อ 2 ของวงรอบในบทที่ 4) ·notifyตั้งall_doneขอแค่ต้นน้ำจบ —loadจบแบบupstream_failedก็นับว่าจบnotifyจึงรัน - สี่ปุ่ม — ครั้งแรกบวกสิทธิ์ลองใหม่สามครั้ง · ช่วงรอเพิ่มจากหนึ่งเป็นสามช่วง ช่วงละ 10 วินาที ⇒ นานขึ้นอย่างน้อย 20 วินาที · และไม่ได้อะไรเลย คอลัมน์หายคือความพังถาวร (คำเตือนท้ายหัวข้อ 2.2 · บทที่ 10)
บทที่ 8
- รับไฟล์: อ่านไฟล์ของวันนั้น เขียนแถวลงที่พักข้อมูลของคุณ คืนใบสรุป (path · จำนวนแถว · คอลัมน์ที่ขาด) · ตรวจ: โครงผิดยก error ไม่งั้นคืนใบสรุปเดิม · โหลด: อ่านแถวจาก path เขียนลงปลายทาง คืนจำนวนแถว · สิ่งที่ไม่ควรผ่าน XCom คือแถวข้อมูลทั้งก้อน
- ท่อรันผ่านเพราะ 33 แถวยังเล็ก · ที่เสียคือแถวทั้งหมดถูกเขียนลงฐานข้อมูลของ Airflow ทุกขั้นทุกรอบ ฐานข้อมูลเดียวกับที่เก็บสถานะของทุกช่อง ไฟล์ยิ่งใหญ่ยิ่งถ่วงทั้งระบบ · "ถูก" ตามบทที่ 2 เพราะข้อมูลอยู่นอกโปรเซสจริง แต่ผิดตามบทนี้ — ข้อมูลควรอยู่ในที่เก็บของคุณ ไม่ใช่ในฐานข้อมูลของตัวคุม
- เป็น
ในคลังเดือนนี้ 198 แถวขณะที่ยอดในบรรทัดเดียวกันยังเป็น 951 ชิ้น เพราะยอดคิดจากรอบนี้รอบเดียว — สองตัวเลขในบรรทัดเดียวขัดกันเอง (ยันในแล็บ: ได้ 198 จริง) ·loadต่อท้าย รันซ้ำจึงเบิ้ล อาการเดียวกับเช้าวันที่ 1 มิถุนายนในบทที่ 1 ที่ยอดกลายเป็น 1,480 · ทางแก้คือบทที่ 12
บทที่ 9
- วันอาทิตย์ไม่มีไฟล์ =
skippedไม่มีใครต้องลงมือ · ไฟล์ไม่มีคอลัมน์ยอดขาย =failedขั้นถัดไปใช้ไม่ได้และต้องมีคนไปขอไฟล์ใหม่ · ร้านใหม่ = ปล่อยผ่าน ร้านเป็นค่าที่ต้นทางเป็นเจ้าของ ไม่ใช่โครงของไฟล์ - เขียว —
gateเป็นskippedแล้วไหลลงปลายน้ำเป็นskippedทั้งหมด ขั้นปลายสุดที่ข้ามนับเป็นผ่าน (หัวข้อ 9.3 รอบเดือนตุลาคมคือกลไกเดียวกัน) · ไม่มีใครรู้ — ไม่มีอะไรโหลด ไม่มีรายงาน และไม่มีรอบแดงให้ใครเปิดดู จนสิ้นเดือนมีคนถามหายอดเดือนสิงหาคม - ผ่าน — ตัวอ่านของบทที่ 8 อ่านทุกแถว ได้ 33 แถว คอลัมน์ครบ วันที่ถูกเดือน ส่วนแถวซ่อนเป็นเรื่องของคำเตือนในรายงาน · ถ้าใช้ตัวอ่านเก่าที่ข้ามแถวซ่อน จะได้ 0 แถว ด่านหยุดท่อด้วย "มี 0 แถว นอกช่วง 10–200" ก่อนถึงขั้นโหลด — เช้าวันที่ 1 สิงหาคมของบทที่ 1 จะเป็นรอบแดงแทนรายงาน 0 ชิ้น
บทที่ 10
- หมดเวลารอ กับส่งถี่เกิน = ลองใหม่ (รอแล้วหาย · แบบหลังควรรอนานขึ้น) · รหัสเข้าใช้หมดอายุ กับข้อมูลขาดฟิลด์บังคับ = หยุดทันที รอเท่าไหร่ก็ไม่หาย ต้องมีคนแก้ — หน้า Tasks ทางการยกตัวอย่างรหัสเข้าใช้ไม่ถูกต้องไว้เป็นกรณีที่ควรพังทันทีตรงตัว
- จังหวะถัดไป
gateเป็นupstream_failed 0· แล้วtransform·load·notifyตามลงทีละจังหวะ (กติกาข้อ 2 ของวงรอบ) · รอบแดงเพราะขั้นปลายสุดnotifyเป็นupstream_failedและไม่มีall_doneมากลบ — ตรงกับผลของstates.pyในหัวข้อ 10.2 (watcher จำเป็นเฉพาะท่อที่มีขั้นปลายสุดแบบall_done) - ในบล็อก
exceptของตัวรัน — ตรงที่ตัดสินว่าup_for_retryหรือfailed· เพิ่มเงื่อนไขว่า error ชนิดที่รู้ว่าถาวร ให้เป็นfailedทันทีโดยไม่ดูจำนวนครั้ง · ชิ้นนั้นคือAirflowFailException(หรือกฎFAILของ retry policy ตั้งแต่ 3.3) — ส่วนการแยกว่า error ไหนถาวร ยังเป็นงานของโค้ดขั้นนั้นเหมือนเดิม
บทที่ 11
- clear
receive_cพร้อม Downstream — ช่องที่อ่านไฟล์ที่เปลี่ยน · หน้าต่างบอก 7 ช่องเหมือนเดิม:receive_cกับปลายน้ำทั้งหกขั้น (gateถึงwatcher) ส่วนreceive_areceive_barrivedไม่ถูกแตะ - ครั้งที่ 2 —
gateจับBadFileแล้วยกAirflowFailExceptionซึ่งหยุดทันทีไม่สนสิทธิ์ที่เหลือ จึงไม่ได้ใช้เพดาน 3 · ช่องจบfailed 2และwatcherทำให้รอบแดงอีกครั้ง - ตัวนับครั้งที่ลองนับต่อ ไม่ถูกล้างตอน clear ·
watcherเคยรันมาแล้ว 1 ครั้ง (พังในรอบแรก) หลัง clear เป็นnone 1· รอบนี้ไม่มีขั้นไหนพัง วงรอบจึงข้ามมันโดยไม่ได้รัน ตัวนับไม่ขยับ =skipped 1· ช่องในบทที่ 9 ไม่เคยรันเลยตั้งแต่ต้น จึงเป็นskipped 0
บทที่ 12
- ไม่ได้ — อีเมลที่ส่งแล้วเรียกคืนไม่ได้ การส่งครั้งที่สองเปลี่ยนโลกข้างนอกเสมอ · กันได้ด้วยการจดลงที่เก็บของคุณว่า "ส่งยอดของเดือน 2026-07 แล้ว" ในธุรกรรมเดียวกับการตัดสินใจส่ง แล้วเช็กก่อนส่งทุกครั้ง · หรือยอมให้ส่งซ้ำได้แต่หัวเรื่องบอกเดือนและบอกว่าเป็นฉบับแก้ ให้ผู้จัดการรู้ว่าฉบับไหนใหม่กว่า
load— ได้ เพราะรันซ้ำแล้วคลังเท่าเดิม และ error ที่loadเจอบ่อยคือฐานข้อมูลสะดุด ซึ่งรอแล้วหาย (บทที่ 10) ·notify— ระวัง ทุกครั้งที่ลองใหม่หลังส่งไปแล้วคือข้อความเพิ่มหนึ่งฉบับ ควรคง 0 หรือกันส่งซ้ำแบบข้อ 1 ก่อน- 1,020 ชิ้น — สาขา ก กับ ข ถูกเขียนทับด้วยของเดือนเดียวกัน ไม่ใช่ต่อท้าย (ตารางของบทที่ 1 ต้องมีคอลัมน์เดือนก่อน ถึงจะลบเฉพาะของเดือนนั้นได้) · ที่ยังขาด: ทำถึงไหนแล้ว · อะไรรออะไร · ใครรู้ว่ามันพัง — ยังอยู่ในโปรเซสเหมือนเดิม การรันซ้ำปลอดภัยแล้ว แต่คืนวันที่ 1 มิถุนายนยังต้องรอคนตื่นมาเห็นรหัสจบการทำงาน 1 และรันใหม่ทั้งก้อนเอง
บทที่ 13
- ช่วง 01:00 วันที่ 9 → 01:00 วันที่ 10 — รอบเกิดตอนปลายช่วง และ
logical_dateคือต้นช่วง · งานควรสรุปของวันที่ 9 โดยคิดจากdata_interval_start - รอบเดียว คือกันยายน งวดล่าสุดที่จบแล้ว (แล็บ unpause วันที่ 7 ตุลาคมได้กันยายนรอบเดียว) · ถ้าดับข้ามหลายต้นเดือน เดือนที่ข้ามไปต้อง backfill เอง · clear ได้ปลอดภัย: clear รันช่องเดิมซ้ำ (บทที่ 11) ทุกขั้นคิดเดือนจากช่วงของรอบ และโหลดเขียนทับทั้งเดือน (บทที่ 12) — รันซ้ำวันไหนก็ได้กันยายน
- ทั้งสามรอบรันวันที่ 7 ตุลาคม จึงได้เดือนกันยายนหมด · พฤษภาคมถึงกรกฎาคมไม่ถูกโหลดเลย · ทั้งสามรอบพยายามประมวลผลกันยายนเดือนเดียวกัน (ถ้าไฟล์ครบ คลังยังเป็นชุดเดียวเพราะเขียนทับทั้งเดือน) · Grid โชว์สามคอลัมน์ของพฤษภาคมถึงกรกฎาคม แต่ไม่มีอะไรบอกว่าข้างในคือกันยายน —
now()ทำให้ backfill เสียความหมายทั้งหมด
บทที่ 14
- ไฟล์ยอดขายรายวัน = นาฬิกา มาตรงเวลาเสมอ ไม่มีอะไรต้องรอ · ไฟล์ราคา = Asset ระบบนั้นเรียก API บอกได้ว่ามาแล้ว จึงไม่ต้องเดาเวลา · ถ้าระบบราคาเรียก API ไม่ได้ ค่อยถอยไปใช้ sensor แบบ
reschedule - ผู้ผลิตประกาศเดือน 2026-09 ของสาขา ค ทั้งที่ไม่มีไฟล์ · ครบสามประกาศ
stock_monthlyถูกปลุก ·arrivedผ่าน (มีไฟล์สองสาขา) ·receive_cยกFileNotFoundError→up_for_retryสองครั้ง →failed 3(บทที่ 10) ·gateถึงnotifyเป็นupstream_failed·morningผ่าน ·watcherพัง รอบแดง · ประกาศที่ไม่มีของจริงย้ายความพังจากผู้ผลิตไปไว้ที่ผู้ใช้ และเสียเวลารอ retries เปล่า — ผู้ผลิตยันก่อนประกาศคือด่านแบบบทที่ 9 ที่ต้นทาง - ถ้ารอบใหม่เกิด ท่อโหลดกันยายนซ้ำด้วยไฟล์ฉบับแก้ของสาขา ข ·
loadลบของเดือนนั้นแล้วเขียนใหม่ทั้งเดือน คลังจึงเป็นตัวเลขฉบับแก้ 99 แถว ไม่เบิ้ล (บทที่ 12) · สิ่งเดียวที่ซ้ำคือข้อความรายงาน (notify·morning) ซึ่งต้องกันส่งซ้ำแบบหัวข้อ 12.5
บทที่ 15
queued= scheduler ส่งให้ executor แล้วแต่ยังไม่มี worker รับ — สงสัย executor/worker ก่อน (ในแล็บคือ worker ของLocalExecutorที่เป็นลูกของ scheduler) และเพดานอย่างparallelism· dag-processor ไม่น่าใช่ เพราะถึงqueuedแล้วแปลว่ากราฟถูกอ่านและ scheduler ตัดสินแล้ว ·task_queued_timeout(ตั้งต้น 600 วินาที ตามหน้า config reference) จัดการช่องที่ค้างนานเกิน- ไม่ — heartbeat บอกว่าคนรันยังมีชีวิต ไม่ได้ดูว่างานคืบหน้า · ต้องตั้ง
execution_timeoutที่ขั้นนั้น · แล็บลองงานค้าง 30 วินาทีกับexecution_timeout5 วินาที ได้AirflowTaskTimeoutทุกครั้ง แล้ว retries ลองซ้ำตามปกติ (failed 2ที่ retries 1) — API ที่ค้างเป็นพัก ๆ จึงได้ประโยชน์จากการลองใหม่ (บทที่ 10) - worker ว่างอยู่ 31 ตัวตลอดรอบอยู่แล้ว ช่องว่างไม่ได้มาจากการรอ worker · มันอยู่ระหว่างต้นน้ำจบกับ
scheduledคือรอวงรอบของ scheduler หมุนมาเห็น · แล็บย่อช่องว่างได้ด้วยscheduler_idle_sleep_timeไม่ใช่ด้วยการเพิ่ม worker
บทที่ 16
- ไม่ต้อง — ขั้นเดียว สั่งใหม่ทั้งก้อนไม่เสียอะไร Task Scheduler พอ · สี่ขั้นแล้วคำตอบเปลี่ยน: พังที่โหลดแล้วไม่อยากดาวน์โหลดใหม่ และส่งเมลต้องรอโหลดเสร็จ — ความรู้ว่าทำถึงไหนเริ่มมีค่า · เลือกอะไรก็ตาม ขั้นโหลดต้องรันซ้ำได้และขั้นส่งเมลต้องกันส่งซ้ำ (บทที่ 12)
@dailyในรุ่น 3 ก็เป็น cron สตริง ได้CronTriggerTimetableช่วงยาว 0 (ยันในแล็บ) ·dsคิดจากlogical_dateซึ่งตอนนี้คือวันที่รัน ไม่ใช่เมื่อวาน ⇒ ทุกรอบดึงยอดของวันนี้ที่ยังไม่จบแทนเมื่อวาน ได้ตัวเลขครึ่งวันโดยไม่มี error · แก้ด้วยschedule=CronDataIntervalTimetable("0 0 * * *", timezone=…)ให้ DAG นั้น แล้วอ่านช่วงจากdata_interval_start(บทที่ 13)- ตัวอย่างหนึ่งชุด: สถานะของแต่ละขั้นของแต่ละรอบเก็บที่ไหน และโปรเซสที่ตายกลางงานแล้วเปิดใหม่ทำต่อได้ไหม (บทที่ 4) · clear หรือรันซ้ำแตะช่องไหน ตัวนับครั้งนับต่อไหม (บทที่ 11) · รอบรู้ได้ยังไงว่าตัวเองประมวลผลช่วงไหน และย้อนรันช่วงเก่ายังไง (บทที่ 13)
อ่านต่อ: System Design#
เล่มนี้จบที่ท่อหนึ่งท่อบนเครื่องเดียว และบทเรียนที่หนักที่สุดคือบทที่ 12: ระบบที่ลองใหม่ได้ บังคับให้ทุกขั้นรันซ้ำแล้วได้ผลเท่าเดิม · พอท่อเริ่มคุยกับบริการอื่นผ่านคิวและ API คำถามเดียวกันกลับมาในอีกรูป — คิวส่วนใหญ่ส่งข้อความซ้ำได้ ผู้รับจึงต้องรับซ้ำแล้วได้ผลเดิม · และข้อความที่พังถาวรต้องมีที่ไปแทนการลองใหม่ไม่รู้จบ (บทที่ 10 ในโลกของคิว)
System Design Handbook — ไปต่อที่ บทที่ 09 · Message Queue & Async Communication ซึ่งพา retry กับ idempotency ออกจากท่อข้อมูลไปสู่ระบบที่หลายบริการคุยกัน
If this was useful —buy me a coffeeor PromptPay · Lightning