← คู่มือ Airflow — ทำไมงานหลายขั้นต้องมีคนคุมลำดับ
LEVEL 1 · พื้นฐานที่ใช้ทุกวัน
DAG แรก ทีละบรรทัด
ลองเดาก่อน: ไฟล์ DAG ข้างล่างบอกว่า receive_a คืนค่า 120 — ถ้าสั่ง Python รันไฟล์นี้ตรง ๆ แบบสคริปต์ธรรมดา จะเห็นเลข 120 ไหม · จดคำตอบไว้ เฉลยอยู่หัวข้อ 6.3
6.1 จาก NEEDS สู่ไฟล์ DAG
ท่อเจ็ดขั้นของบทที่ 3 เขียนด้วยไวยากรณ์ของ Airflow · ยังไม่อ่าน Excel ใช้ยอดเดือนพฤษภาคมของบทที่ 1 แทน งานจริงมาในบทที่ 8 · บันทึกไว้ในโฟลเดอร์ dags ใต้ AIRFLOW_HOME ซึ่งเป็นที่ที่ Airflow มองหาไฟล์ DAG
# ~/airflow-lab/dags/stock_first.py
from datetime import datetime
from airflow.sdk import dag, task
COUNTS = {"a": 120, "b": 340, "c": 560} # ยอดเดือนพฤษภาคมของบทที่ 1 · บทที่ 8 เปลี่ยนเป็นอ่านไฟล์จริง
@dag(schedule=None, start_date=datetime(2026, 5, 1), catchup=False)
def stock_first():
@task
def receive(branch):
print(f"receive_{branch}: {COUNTS[branch]} ชิ้น")
return COUNTS[branch]
@task
def merge_check(counts):
print(f"merge_check: ได้รับ {counts} ชิ้น ครบ {len(counts)} สาขา")
return counts
@task
def transform(counts):
print(f"transform: รวม {sum(counts)} ชิ้น")
return sum(counts)
@task
def load(total):
print(f"load: รับ {total} ชิ้น (ยังไม่มีฐานข้อมูล — บทที่ 8)")
return total
@task
def notify(total):
print(f"notify: ส่งรายงาน สต็อกรวม {total:,} ชิ้น")
# ฟังก์ชันเดียวใช้ซ้ำสามแถว — override ตั้งชื่อแถวให้ไม่ชนกัน
counts = [receive.override(task_id=f"receive_{b}")(b) for b in COUNTS]
notify(load(transform(merge_check(counts))))
stock_first()ไล่กลับไปหาตารางของบทที่ 4
@dag= ท่อหนึ่งท่อ ชื่อฟังก์ชันคือชื่อ DAG ·schedule=None= ยังไม่มีนาฬิกามาเพิ่มคอลัมน์ให้ ต้องสั่งเอง (start_dateกับcatchupเป็นเรื่องของนาฬิกา บทที่ 13)@taskแต่ละตัวที่ถูกเรียก = หนึ่งแถว ·override(task_id=...)ตั้งชื่อแถวเองreceiveตัวเดียวจึงเป็นสามแถว- ส่งค่าที่ฟังก์ชันหนึ่งคืนเข้าอีกฟังก์ชัน = เส้นของกราฟ ·
merge_check(counts)ทำให้merge_checkรอ receive ทั้งสาม เหมือนบรรทัดของมันในNEEDS· ไฟล์นี้ไม่มีคำว่า "รอ" สักคำ เส้นทั้งหมดมาจากการส่งค่า
6.2 รันหนึ่งรอบ แล้วดูว่าค่าข้ามขั้นไปได้ยังไง
airflow dags test เพิ่มคอลัมน์ใหม่แล้วเดินวงรอบจนจบในหน้าต่างที่สั่ง ไม่ต้องรอ scheduler · grep เก็บเฉพาะบรรทัดที่งานพิมพ์เอง (ทุกบรรทัดมีคำว่า "ชิ้น") จากบันทึกเป็นร้อยบรรทัด
airflow dags test stock_first 2>&1 | grep ชิ้นผลรันจริง
receive_b: 340 ชิ้น
receive_a: 120 ชิ้น
receive_c: 560 ชิ้น
merge_check: ได้รับ [120, 340, 560] ชิ้น ครบ 3 สาขา
transform: รวม 1020 ชิ้น
load: รับ 1020 ชิ้น (ยังไม่มีฐานข้อมูล — บทที่ 8)
notify: ส่งรายงาน สต็อกรวม 1,020 ชิ้นสามบรรทัดแรกสลับลำดับได้ทุกครั้ง — แล็บสองรอบติดกันได้ b · a · c แล้วได้ b · c · a เพราะอยู่จังหวะเดียวกันในบทที่ 3 · แต่ merge_check ได้ [120, 340, 560] เรียงตามสาขาทุกครั้ง ทั้งที่แต่ละขั้นรันแยกเป็นงานของตัวเอง — ค่าข้ามขั้นไปได้ยังไง
💡 XCom — ที่เก็บค่าเล็ก ๆ ที่ขั้นหนึ่งส่งให้อีกขั้น · ค่าที่ @task คืนถูกเขียนลง XCom เอง (ที่เก็บตั้งต้นคือฐานข้อมูลของ Airflow) แล้วขั้นที่รับอ่านกลับขึ้นมาตอนเริ่ม · เอกสารทางการให้เหตุผลว่างานแต่ละงานแยกขาดจากกัน และอาจรันคนละเครื่องด้วยซ้ำ
นี่คือคำตอบของคำเตือนในหัวข้อ 2.1 ที่ว่าผลของแต่ละขั้นต้องอยู่นอกโปรเซส · merge_check ได้ลำดับถูกเพราะมันหยิบ XCom ของแต่ละแถวตามตำแหน่งที่ประกาศไว้ ไม่ใช่ตามลำดับที่ใครเสร็จก่อน — หัวข้อถัดไปจะเห็นใบสั่งหยิบนั้นกับตา
6.3 พิสูจน์ด้วยตาเอง — ไฟล์ DAG คือการประกาศ ไม่ใช่การรัน
ไฟล์ทดลองนี้ import stock_first.py แบบที่ Airflow ทำตอนอ่านไฟล์ แล้วถามตัว DAG ว่ามีแถวอะไร แต่ละแถวรอใคร · วางไว้นอกโฟลเดอร์ dags แล้วรันด้วย python ~/airflow-lab/probe_declare.py
# ~/airflow-lab/probe_declare.py — ไฟล์ทดลอง วางนอก dags/ จึงไม่ถูกหยิบไปเป็น DAG
import sys
from pathlib import Path
from airflow.sdk import DAG
sys.path.insert(0, str(Path.home() / "airflow-lab" / "dags"))
import stock_first # ← รันโค้ดชั้นบนสุดของไฟล์ DAG ทั้งไฟล์ แบบเดียวกับที่ Airflow ทำทุกครั้งที่อ่านไฟล์
dag = stock_first.stock_first() # เรียกฟังก์ชัน @dag อีกครั้ง เพื่อเก็บตัว DAG ที่ได้ไว้ดู
print("ฟังก์ชัน @dag คืนตัว DAG:", isinstance(dag, DAG))
for task_id in dag.task_ids:
print(f"{task_id:<12} รอ {sorted(dag.task_dict[task_id].upstream_task_ids)}")
print("ช่องแรกของรายการที่ merge_check จะได้รับ:", dag.task_dict["merge_check"].op_args[0][0])ผลรันจริง
ฟังก์ชัน @dag คืนตัว DAG: True
receive_a รอ []
receive_b รอ []
receive_c รอ []
merge_check รอ ['receive_a', 'receive_b', 'receive_c']
transform รอ ['merge_check']
load รอ ['transform']
notify รอ ['load']
ช่องแรกของรายการที่ merge_check จะได้รับ: {{ task_instance.xcom_pull(task_ids='receive_a', dag_id='stock_first', key='return_value') }}เจ็ดบรรทัดกลางคือ NEEDS ของบทที่ 3 ทุกตัวอักษร — Airflow สร้างมันจากการส่งค่าในไฟล์ · และไม่มีบรรทัดไหนมีคำว่า "ชิ้น" ไม่มีงานไหนถูกรันเลย ทั้งที่โค้ดในไฟล์ถูกรันครบทุกบรรทัด
เฉลยคำถามเปิดบท: ไม่เห็น 120 · ตอนประกาศกราฟ receive_a() ไม่ได้รันเนื้อในฟังก์ชัน มันคืนใบสั่งว่า "ตอนรัน ให้ไปหยิบ XCom ของแถว receive_a" — บรรทัดสุดท้ายของผลรัน · ถ้าตอบว่าเห็น 120 นั่นคือภาพของสคริปต์ในบทที่ 1 ที่ทุกบรรทัดทำงานจริงตอนรัน
ไฟล์ DAG คือการประกาศกราฟ ไม่ใช่การรันงาน — เนื้อในฟังก์ชัน
@taskรันเมื่อช่องของมันถูกหยิบเท่านั้น
@task (พื้นจาง) รันเฉพาะตอนช่องของมันถูกหยิบ · การอ่าน Excel ที่ย้ายขึ้นไปอยู่ในบรรทัดที่มีแถบเน้น จะถูกอ่านทุกครึ่งนาทีทั้งวัน⚠️ ความเข้าใจผิด: "โค้ดในไฟล์ DAG รันเฉพาะตอน DAG ทำงาน" — หน้า Best Practices ทางการมีหัวข้อ "Top level Python Code" เตือนเรื่องนี้โดยเฉพาะ · ผิดเพราะ Airflow ต้องรันไฟล์ทั้งไฟล์ทุกครั้งที่อยากรู้ว่ากราฟหน้าตาเป็นยังไง · กลไกที่ถูกคือโค้ดนอก @task รันทุกครั้งที่ไฟล์ถูกอ่าน เนื้อใน @task รันเฉพาะตอนช่องถูกหยิบ · พิสูจน์ได้ในหัวข้อถัดไป
6.4 ไฟล์ถูกอ่านซ้ำทุก 30 วินาที แม้ไม่มีใครแก้
แล็บของผู้เขียนวางไฟล์ DAG ที่มี print("PA_TOPLEVEL pid=...") ไว้ชั้นบนสุด แล้วปล่อยทิ้ง 25.5 นาทีโดยไม่แตะไฟล์ · บันทึกของ dag-processor (~/airflow-lab/logs/dag_processor/<วันที่>/dags-folder/<ชื่อไฟล์>.log) มีบรรทัดนั้นแบบนี้
{"timestamp":"2026-10-05T09:24:59.538490Z","level":"info","event":"PA_TOPLEVEL pid=620","logger":"dag_processor.stdout"}
{"timestamp":"2026-10-05T09:25:32.011243Z","level":"info","event":"PA_TOPLEVEL pid=680","logger":"dag_processor.stdout"}
{"timestamp":"2026-10-05T09:26:02.901418Z","level":"info","event":"PA_TOPLEVEL pid=748","logger":"dag_processor.stdout"}
… (อีก 47 บรรทัด)
{"timestamp":"2026-10-05T09:50:30.236556Z","level":"info","event":"PA_TOPLEVEL pid=1875","logger":"dag_processor.stdout"}รวม 51 ครั้ง ห่างกันราว 30 วินาที (ต่ำสุด 30.2 · กลางสุด 30.5 · สูงสุด 32.5) และ pid ไม่ซ้ำเลย — ทุกครั้งเป็นโปรเซสใหม่ · 30 คือค่าตั้งต้นของ [dag_processor] min_file_process_interval
ย้ายการอ่าน Excel ขึ้นไปชั้นบนสุด = อ่านทุกครึ่งนาทีทั้งวันทั้งคืน · Best Practices จึงห้ามงานหนัก การต่อฐานข้อมูล และการเรียกเครือข่ายที่ชั้นบนสุด — ใครเป็นคนอ่าน และทำไมต้องเปิดโปรเซสใหม่ทุกครั้ง คือปริศนาของบทที่ 15
6.5 ไฟล์ใหม่ไม่โผล่ทันที
การหาไฟล์ใหม่ในโฟลเดอร์เป็นอีกนาฬิกาหนึ่ง: [dag_processor] refresh_interval ค่าตั้งต้น 300 วินาที · แล็บวัดสองครั้ง ไฟล์ใหม่ถูกเจอหลังวาง 282 และ 217 วินาที แล้วแต่ว่าวางตอนไหนของรอบสแกน · ไม่อยากรอ สั่ง airflow dags reserialize (แล็บ: ขึ้นภายใน 7 วินาที)
DAG ใหม่เริ่มในสถานะหยุด (dags_are_paused_at_creation ค่าตั้งต้น True) — สั่งเริ่มรอบได้ แต่คอลัมน์ใหม่ค้างเป็น queued ทุกช่องว่าง จนกว่าจะ airflow dags unpause <ชื่อ> · ส่วน airflow dags test ไม่สนสถานะหยุด ใช้ได้ทันทีที่บันทึกไฟล์
ลองตอบโดยไม่ย้อนกลับไปอ่าน (เฉลยท้ายเล่ม)
- ไฟล์ DAG ของเพื่อนมีบรรทัด
rows = read_all_branches()อยู่ชั้นบนสุด แล้วส่งrowsให้@taskตัวหนึ่ง · ฟังก์ชันนั้นอ่าน Excel สามไฟล์ใช้เวลา 20 วินาที — จะเกิดอะไรขึ้นกับเครื่อง แม้ทั้งวันไม่มีรอบเลย และควรย้ายบรรทัดนั้นไปไว้ไหน - ถ้าแก้บรรทัดสุดท้ายของ
stock_firstเป็นnotify(load(transform(merge_check(counts[:2]))))·receive_cยังเป็นแถวในตารางไหม รอใคร มีใครรอมัน และถ้าreceive_cพัง รอบจะจบเป็นอะไร - ทำไม
merge_checkได้[120, 340, 560]เรียงตามสาขาทุกครั้ง ทั้งที่ receive ทั้งสามเสร็จคนละลำดับในแต่ละรอบ