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

← คู่มือ 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 รันเมื่อช่องของมันถูกหยิบเท่านั้น

stock_first.pyfrom airflow.sdk import dag, taskCOUNTS = {"a": 120, "b": 340, "c": 560}@dag(schedule=None, start_date=..., catchup=False)def stock_first(): @task def receive(branch): print(...) return COUNTS[branch] … merge_check · transform · load · notify counts = [receive.override(task_id=f"receive_{b}")(b) for b in COUNTS] notify(load(transform(merge_check(counts))))stock_first()รันทุกครั้งที่ไฟล์ถูกอ่านราวทุก 30 วินาที แม้ไม่มีรอบโปรเซสใหม่ทุกครั้งแล็บ: 51 ครั้งใน 25.5 นาทีได้: กราฟแบบ NEEDS+ ใบสั่งหยิบ XComรันเมื่อช่องของมันถูกหยิบครั้งละหนึ่งช่องของหนึ่งรอบไม่ถูกแตะตอนอ่านไฟล์ได้: ค่าที่คืน → XCom
FIG 6.1 — ไฟล์ DAG หนึ่งไฟล์มีสองนาฬิกา · บรรทัดที่มีแถบเน้นรันทุกครั้งที่ไฟล์ถูกอ่าน — ราวทุก 30 วินาที แม้ไม่มีรอบเลย และสิ่งที่มันสร้างคือกราฟ ไม่ใช่ผลของงาน · เนื้อในฟังก์ชัน @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 ไม่สนสถานะหยุด ใช้ได้ทันทีที่บันทึกไฟล์

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

  1. ไฟล์ DAG ของเพื่อนมีบรรทัด rows = read_all_branches() อยู่ชั้นบนสุด แล้วส่ง rows ให้ @task ตัวหนึ่ง · ฟังก์ชันนั้นอ่าน Excel สามไฟล์ใช้เวลา 20 วินาที — จะเกิดอะไรขึ้นกับเครื่อง แม้ทั้งวันไม่มีรอบเลย และควรย้ายบรรทัดนั้นไปไว้ไหน
  2. ถ้าแก้บรรทัดสุดท้ายของ stock_first เป็น notify(load(transform(merge_check(counts[:2])))) · receive_c ยังเป็นแถวในตารางไหม รอใคร มีใครรอมัน และถ้า receive_c พัง รอบจะจบเป็นอะไร
  3. ทำไม merge_check ได้ [120, 340, 560] เรียงตามสาขาทุกครั้ง ทั้งที่ receive ทั้งสามเสร็จคนละลำดับในแต่ละรอบ

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