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

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

LEVEL 4 · ระดับเทพ

กลไกภายใน: หลายโปรเซสแบ่งกันถือตารางเดียว

ปริศนาที่ค้างมาทั้งเล่มมีคำตอบอยู่ข้างในเครื่องที่ standalone เปิด: บันทึกสี่ชื่อทั้งที่แบบจำลองมีวงรอบเดียว (หัวข้อ 4.7 · 5.4) · worker ว่าง 32 ตัว (5.5) · ไฟล์ DAG ถูกอ่านซ้ำด้วยโปรเซสใหม่ (6.4) · ขั้นถัดไปเริ่มช้าราว 1 วินาที (7.2) · XCom เก็บที่ไหน · states.py อ่านตารางอะไร (9.2) · คนรันตายกับคนรันช้าแยกกันยังไง (4.6)

แบบจำลองของบทที่ 4 ยังถูก — ของจริงแค่แบ่งงานของวงรอบเดียวให้หลายโปรเซส และทุกโปรเซสคุยกันผ่านตารางในฐานข้อมูล

15.1 ใครเขียนตาราง

ทุกส่วนในบทนี้ตอบสามคำถามเดียวกัน: (ก) ทำอะไร (ข) อ่านหรือเขียนตารางตรงไหน (ค) หลักฐานในแล็บ

แล็บพิมพ์ต้นไม้โปรเซสระหว่างที่งานหนึ่งกำลังรัน (ตัด worker ว่างออก)

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

airflow standalone
 ├─ airflow scheduler
 │   ├─ airflow serve-logs
 │   ├─ airflow worker -- LocalExecutor: <idle>          ← 31 ตัวแบบนี้
 │   └─ airflow worker -- LocalExecutor: 01a115cb-…      ← ตัวที่รับงาน
 │       └─ airflow worker -- 01a115cb-…                 ← โปรเซสที่รันโค้ดของงาน
 ├─ airflow dag-processor
 ├─ airflow api_server -- host:0.0.0.0 port:8080
 └─ airflow triggerer
     ├─ airflow serve-logs
     └─ airflow triggerer
โปรเซสใต้ standaloneตารางในฐานข้อมูลเดียวกันdags/*.pydag-processorอ่านไฟล์ทุก 30 วินาที ด้วยโปรเซสใหม่schedulerวงรอบของบทที่ 4workerโปรเซสลูกรันโค้ดของงานapi-serverรับแทนงาน แล้วเขียนลงตารางtriggererถือการรอของช่อง deferredserialized_dagdag_runtask_instancexcomscheduled · queuedสถานะที่เหลือ · heartbeatส่งช่องให้รันสถานะ · XCom · heartbeat ทุกราว 5 วินาทีช่องที่รอไม่กิน worker · ครบเวลาแล้วกลับไปจบใน workerฐานข้อมูลของแล็บมี 59 ตาราง · สี่ตารางนี้คือที่เล่มใช้มาตลอด (states.py อ่าน dag_run กับ task_instance)
FIG 15.1 — ห้าส่วนที่ standalone เปิด กับตารางที่แต่ละส่วนเขียน (ลูกศรทึบ = เขียน) · ทุกส่วนคุยกันผ่านตารางในฐานข้อมูลเดียว · scheduler เขียนแค่ถึง queued · worker ไม่เขียนตารางเอง ทุกอย่างของงานไปทาง api-server — heartbeat ที่เงียบไปของเส้นนี้คือสิ่งที่ scheduler ใช้จับคนรันที่ตาย (หัวข้อ 15.4)
  • dag-processor (ก) อ่านไฟล์ DAG แล้วเก็บกราฟ (ข) เขียน serialized_dag (ค) อ่านทุก 30 วินาทีด้วยโปรเซสใหม่ — หน้า Dag File Processing ทางการให้เหตุผลว่าส่วนนี้รันโค้ดของผู้ใช้ จึงแยกโปรเซส และอ่านแต่ละไฟล์ในโปรเซสที่มีเวลาจำกัด (dag_file_processor_timeout) โค้ดชั้นบนสุดที่ค้างจึงไม่ลาก scheduler ไปด้วย
  • scheduler (ก) วงรอบของบทที่ 4 ตรงตัว: หาช่องที่ต้นน้ำผ่านแล้ว เปลี่ยนเป็น scheduled → queued ส่งให้ executor · ตัดสินรอบ · จับคนรันที่ตาย (ข) เขียน dag_run กับ task_instance (ค) worker ของ LocalExecutor เป็นลูกของมัน 32 ตัวเท่า parallelism
  • worker (ก) แตกโปรเซสลูกไปรันโค้ดของงาน (ข) ไม่เขียนตารางตรง ส่งผ่าน api-server (ค) pid ที่งาน print ตรงกับโปรเซสลูกในต้นไม้
  • api-server (ก) UI และ Task Execution API (ข) รับสถานะ XCom และ heartbeat จากงานแล้วเขียนลงตารางแทน — ข้อ 4 ของหัวข้อ 4.6 (ค) บันทึกมีคำขอ PUT /execution/task-instances/<id>/heartbeat ทุกราว 5 วินาทีตลอดที่งานรัน
  • triggerer (ก) ถือ "การรอ" แทน worker สำหรับงานแบบ deferrable (ข) ช่องเป็น deferred ระหว่างรอ ไม่กินช่องของ worker (ค) sensor แบบ deferrable ในแล็บเป็น deferred 20 วินาที บันทึกเขียน Pausing task as DEFERRED · การรออยู่ในบันทึกของ trigger · ครบเวลาแล้วงานกลับไปจบใน worker

15.2 เปิดตารางดู

ฐานข้อมูลของแล็บมี 59 ตาราง · ตารางของบทที่ 4 คือ dag_run (คอลัมน์) กับ task_instance (ช่อง: state · try_number · max_tries บวกสิ่งที่แบบจำลองไม่มีอย่าง scheduled_dttm · queued_dttm · pid · last_heartbeat_at) — สองตารางที่ states.py อ่าน

XCom อยู่ตาราง xcom ในฐานข้อมูลเดียวกัน — คำตอบของคำเตือนในหัวข้อ 8.5 เห็นกับตา:

บันทึกจากแล็บ (ตาราง xcom ของรอบ backfill เดือนกรกฎาคม · ตัด path)

arrived    return_value  3
receive_a  return_value  {"branch": "a", "path": ".../staging/2026-07/a.json", "rows": 33, "hidden": 33, "missing": [], "problems": []}
transform  return_value  {"path": ".../staging/2026-07/all.json", "rows": 99, "pieces": {"a": 121, "b": 382, "c": 448}}
load       return_value  99

ค่าที่คืนจากทุกขั้นถูกเขียนลงตารางเดียวกับสถานะของทุกช่องในระบบ · คืน DataFrame ทั้งก้อน = ยัดมันลงตรงนี้ทุกขั้นทุกรอบ

15.3 วินาทีที่หายไประหว่างขั้น

หัวข้อ 7.2 ถามว่าขั้นถัดไปเริ่มช้ากว่าต้นน้ำจบราว 1 วินาทีเพราะอะไร · task_instance บันทึกเวลาทุกจังหวะ แล็บดึงรอบหนึ่งของ stock_grid มาดู

บันทึกจากแล็บ (เวลาเป็นวินาทีของนาทีเดียวกัน)

              ต้นน้ำจบ   scheduled   queued    เริ่ม
merge_check   02.172    02.935      02.951    03.063
transform     03.465    04.040      04.060    04.213
load          04.409    05.188      05.213    05.258
notify        05.538    06.383      06.394    06.409

ช่องว่างเกือบทั้งหมดอยู่ระหว่าง "ต้นน้ำจบ" กับ scheduled — ช่วงที่รอให้วงรอบหมุนมาเห็น · พอเห็นแล้ว ส่งต่อถึงเริ่มใช้ไม่ถึง 0.2 วินาที · ค่า [scheduler] scheduler_idle_sleep_time ตั้งต้น 1 วินาที — หน้า config reference บอกว่าวงรอบที่ไม่มีอะไรให้ทำจะหลับเท่านี้ก่อนหมุนรอบถัดไป · แล็บเปิด standalone ใหม่ด้วยค่า 0.1 แล้วรันซ้ำ: ช่องว่างเหลือราว 0.2 วินาที ทั้งรอบจาก 5.5–7.2 วินาทีเหลือ 1.8–2.0 วินาที · วงรอบของบทที่ 4 ก็เป็นจังหวะแบบนี้ — ของจริงแค่หลับระหว่างจังหวะ

15.4 คนรันตาย กับ คนรันช้า

ข้อ 5 ของหัวข้อ 4.6: ช่อง running ที่ไม่มีใครทำแล้วต้องถูกคืนให้รันใหม่ · ช่อง running ที่แค่นานต้องไม่ถูกแย่ง

คำตอบคือ heartbeat: คนรันที่ยังมีชีวิตส่งสัญญาณเป็นระยะ · เงียบนานเกิน [scheduler] task_instance_heartbeat_timeout (ตั้งต้น 300 วินาที) = ตาย · แล็บตั้ง 30 วินาที แล้วฆ่างานที่นอน 90 วินาทีสองแบบ (retries 1)

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

แบบที่ 1 · ฆ่าเฉพาะโปรเซสที่รันโค้ด (worker ยังอยู่)
10:00:58  kill -9            slow running 1
10:01:03                     slow up_for_retry 1     ← worker เห็นลูกตาย รายงานทันที
10:01:08                     slow running 2

แบบที่ 2 · ฆ่า worker พร้อมลูก (เหมือนเครื่องดับ — ไม่เหลือใครรายงาน)
10:04:18  kill -9            slow running 1
10:04:53                     slow running 1          ← ตารางยังโกหกอยู่
10:04:55  scheduler: Failing 1 TIs without heartbeat · Detected a task instance without a heartbeat
10:04:58                     slow up_for_retry 1
10:05:03                     slow running 2

แบบที่ 2 คือหัวข้อ 4.4 บนของจริง · ตารางบอก running 37 วินาทีทั้งที่ไม่มีใครทำ · scheduler รู้จากความเงียบ ไม่ใช่จากใครบอก · งานที่ช้าแต่ยังส่ง heartbeat ไม่โดนแตะต่อให้รันเป็นชั่วโมง · ราคาคือการรอ (ตั้งต้น 5 นาที) และความเงียบหลอกได้ — เครื่องที่เงียบไปชั่วครู่ก็ดูเหมือนตาย งานของมันจะถูกรันซ้ำ ซึ่งบทที่ 12 ทำให้ไม่พังข้อมูล

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

  1. เช้าหนึ่ง Grid มีช่องค้าง queued นานผิดปกติหลายช่อง แต่ไม่มีช่องไหน running · ส่วนไหนในบทนี้ที่คุณจะสงสัยก่อน และอะไรบ้างที่ไม่น่าเป็นต้นเหตุ
  2. งานหนึ่งเรียก API ภายนอกที่ค้างเงียบ 10 นาทีโดยโปรเซสยังอยู่ · Airflow จะตัดสินว่างานนี้ตายไหม เพราะอะไร และถ้าอยากให้หยุดรอเอง ต้องตั้งอะไรที่ตัวงาน — ใช้บทที่ 10 ประกอบ
  3. ทำไมการเพิ่ม worker เป็น 64 ตัวไม่ทำให้ช่องว่างระหว่างขั้นของ stock_grid สั้นลง

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