Top 5 Airflow DAG mistakes

Top 5 Airflow DAG Mistakes and How You Can Prevent Them!

5 ข้อผิดพลาดที่เจอบ่อยที่สุดใน Airflow DAG

Top 5 Airflow DAG Mistakes and Best Practices Cover Banner

Introduction

ทำไม Airflow Scheduler ของทีมถึงช้าลงเรื่อยๆ ทั้งที่จำนวน DAG ก็ไม่ได้เยอะขนาดนั้น หรือทำไมพอ Task พังกลางดึกแล้วกด Retry ไปกลับได้ข้อมูลซ้ำซ้อนแทนที่จะได้ผลลัพธ์เดิม

Airflow Scheduler bottleneck and retry error illustration

มีโอกาสได้ไล่รีวิว Airflow DAG มาหลายตัวจากหลายๆที่ก็เริ่มเห็น Pattern ที่ซ้ำกันจนสังเกตได้ว่าปัญหาส่วนใหญ่ที่ทำให้ Cluster ช้าลง Debug ยากขึ้น หรือ Resource ไม่พอ ไม่ได้มาจากเรื่องซับซ้อนอะไรเลยครับ แต่มาจาก 5 เรื่องนี้ที่เจอซ้ำแล้วซ้ำอีก วันนี้เลยขอรวบรวมมาเล่าให้ฟัง พร้อมตัวอย่างโค้ดแบบ Bad กับ Good เทียบกันให้เห็นชัดๆ


Top 5 common mistakes

1. Top-level Code

อันนี้เจอบ่อยที่สุดเลยครับ คนมักลืมไปว่าไฟล์ DAG ไม่ได้ถูกรันแค่ตอน Deploy ครั้งเดียวแบบสคริปต์ปกติ แต่ Scheduler จะ Parse ไฟล์นี้ซ้ำๆ ตามรอบ dag_dir_list_interval (Default ทุก 30 วินาที) เพื่อเช็คว่ามี DAG ใหม่หรือมีการแก้ไขไหม ซึ่งหมายความว่าโค้ดตัวไหนที่อยู่นอก Task หรือ Operator (เรียกว่า Top-level Code) จะถูกรันซ้ำทุก 30 วินาทีแบบไม่มีวันหยุดครับ

# Bad: เรียก API ตอน Parse DAG ทุกครั้ง
import requests
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
config = requests.get("https://api.internal.company.com/dag-config").json()
with DAG(
dag_id="bad_top_level_call",
start_date=datetime(2024, 1, 1),
schedule="@daily",
) as dag:
task = PythonOperator(
task_id="process",
python_callable=lambda: print(config),
)
# Good: ย้าย Logic ที่หนักไปไว้ใน Task ให้รันตอน Execute จริงเท่านั้น
import requests
from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime
def process(**context):
config = requests.get("https://api.internal.company.com/dag-config").json()
print(config)
with DAG(
dag_id="good_top_level_call",
start_date=datetime(2024, 1, 1),
schedule="@daily",
catchup=False,
) as dag:
task = PythonOperator(
task_id="process",
python_callable=process,
)

ถ้าใน Top-level Code นั้นมีการยิง API เรียก DB หรืออ่านไฟล์ นั่นแปลว่าระบบปลายทางโดนยิงถี่ยิบตลอดเวลาโดยที่ DAG ยังไม่ได้ทำงานจริงด้วยซ้ำ แถม Scheduler เองก็จะช้าลงเพราะต้อง Parse โค้ดหนักๆ ซ้ำไปเรื่อยๆ ลองมาดูตัวอย่างเวลาของ DAG parsing จะเห็นเลยว่าต่างกันขนาดไหน (อันนี้ขนาดว่ามีแค่อันเดียวนะ ถ้าเป็นร้อย เป็นพันละ😱)

Airflow DAG parsing latency comparison between bad and good top-level code

เสริมอีกนิดเรื่อง Airflow Architecture จะได้เข้าใจมากขึ้นว่าทำไมการ parse ซ้ำๆของโค้ดที่ไม่ดีจะส่งผลกระทบกับ resource ยังไง โดยจะเห็นว่า Scheduler จะทำงานเป็นทั้งตัวอ่านไฟล์และส่งงานไปที่ worker ทำให้ถ้า DAG parsing ยิ่งนานก็ยิ่งกระทบไปเรื่อยๆ

Apache Airflow core architecture diagram showing DAG directories, Scheduler, and Workers

2. Sensor ที่ไปจับ Worker Slot ค้างไว้ ทั้งที่มี Deferrable Operator ให้ใช้แล้ว

Deferrable Operator คือ Operator ที่พักตัวเองแล้วปล่อย Worker Slot คืนระหว่างรอ แทนที่จะจับ Slot ไว้เฉยๆ จนกว่าเงื่อนไขจะเป็นจริง

Sensor แบบ Standard (Mode poke หรือ reschedule แบบเก่า) จะจับ Worker Slot ไว้ตลอดเวลาที่รอ ต่อให้ไม่ได้ทำอะไรเลยก็ตาม ลองนึกภาพว่ามี Worker 100 Slot แล้วมี 100 DAG ที่รอ Sensor ตัวเดิมพร้อมกัน Cluster ทั้งก้อนจะโดนล็อกจนรันอย่างอื่นไม่ได้เลยทั้งที่จริงๆแล้วแทบไม่มีอะไรทำงานอยู่เลยครับ

Comparison between traditional sensor blocking worker slots and deferrable operator with triggerer

ทางแก้คือใช้ Deferrable Sensor ซึ่ง Airflow มีให้ใช้แทนตัวเดิมแบบ Drop-in ได้เลยในหลาย Operator (ผ่าน Parameter deferrable=True หรือใช้ Class รุ่น Async เช่น TimeSensorAsync) โดยงานรอจะถูกโยนไปให้ Triggerer จัดการแทน Worker

# Bad: Sensor จับ Worker Slot ไว้ 5 ชั่วโมงเต็มๆ ระหว่างรอ
from airflow.sensors.external_task import ExternalTaskSensor
wait_for_upstream = ExternalTaskSensor(
task_id="wait_for_upstream",
external_dag_id="upstream_dag",
poke_interval=60,
timeout=60 * 60 * 5,
mode="poke",
)
# Good: เปิด deferrable ปล่อย Worker Slot คืนระหว่างรอ ใช้ Triggerer แทน
from airflow.sensors.external_task import ExternalTaskSensor
wait_for_upstream = ExternalTaskSensor(
task_id="wait_for_upstream",
external_dag_id="upstream_dag",
poke_interval=60,
timeout=60 * 60 * 5,
deferrable=True,
)

3. Task ที่ไม่ Atomic ไม่ Idempotent และไม่ Modular

อันนี้เป็นเรื่องที่ทำให้ Debug กับ Retry กลายเป็นฝันร้ายเลยครับ หลายทีมชอบยัด Extract, Transform, Load รวมเป็น Function เดียวใน PythonOperator ตัวเดียว พอ Step ท้ายๆพังจะต้องรันใหม่ทั้งหมดตั้งแต่ต้น แถมถ้า Load เป็นการ INSERT ธรรมดาไม่ใช่ Upsert พอ Retry ก็จะได้ข้อมูลซ้ำเข้าไปอีกชุดทันที

Monolithic ETL vs Atomic modular tasks in Airflow pipeline

Task ที่ดีควรทำงานเดียวจบในตัว (Atomic) และรันซ้ำกี่ครั้งก็ได้ผลลัพธ์เหมือนเดิมเป๊ะ (Idempotent) ไม่ใช่ยิ่งรันยิ่งพังข้อมูลมากขึ้น

# Bad: รวม Extract + Transform + Load ไว้ Function เดียว และ Insert ตรงๆ
def extract_transform_load():
data = extract_from_api()
df = transform(data)
load_to_warehouse(df) # Insert ตรงๆ, Retry ทีก็ได้ข้อมูลซ้ำที
etl_task = PythonOperator(
task_id="etl_all_in_one",
python_callable=extract_transform_load,
)
# Good: แยกเป็น Task ย่อย Atomic แต่ละตัว และ Load แบบ Upsert ตาม Partition
from airflow.decorators import task
@task
def extract(logical_date=None):
partition = logical_date.format("YYYY-MM-DD")
return extract_from_api(partition=partition)
@task
def transform(raw_data):
return transform_fn(raw_data)
@task
def load(transformed_data, logical_date=None):
partition = logical_date.format("YYYY-MM-DD")
upsert_to_warehouse(transformed_data, partition=partition)
load(transform(extract()))

4. Over-provisioning และ Concurrency ที่ไม่มีขอบเขต

ข้อนี้จริงๆแยกได้เป็นสองปัญหาที่มักเกิดคู่กันครับ เลยขอแยกอธิบายทีละอย่างให้ชัดๆ

Over-provisioning คือการตั้ง Resource Request (CPU/Memory) ให้ Task สูงเกินความจำเป็นจริงมาก เช่นใช้ KubernetesPodOperator แล้ว Request CPU 4 Core Memory 16Gi ทั้งที่ Task จริงๆใช้ Resource แค่เสี้ยวเดียว พอมี Task แบบนี้เยอะขึ้นเรื่อยๆ Cluster ก็จะเต็มเร็วกว่าที่ควรจะเป็น เพราะ Kubernetes Scheduler จอง Resource ตาม Request ไม่ใช่ตาม Usage จริง ทำให้ Task อื่นต้องรอคิวทั้งที่ Resource จริงยังเหลือเฟืออยู่เลยครับ

# Bad: Request CPU/Memory เกินความจำเป็นจริงของ Task
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
from kubernetes.client import models as k8s
heavy_request_task = KubernetesPodOperator(
task_id="transform_small_file",
image="my-transform-image:latest",
container_resources=k8s.V1ResourceRequirements(
requests={"cpu": "4", "memory": "16Gi"}, # Task จริงใช้ไม่ถึง 10%
limits={"cpu": "4", "memory": "16Gi"},
),
)
# Good: Profile Usage จริงก่อน แล้วตั้ง Resource ให้ใกล้เคียงของจริง
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
from kubernetes.client import models as k8s
right_sized_task = KubernetesPodOperator(
task_id="transform_small_file",
image="my-transform-image:latest",
container_resources=k8s.V1ResourceRequirements(
requests={"cpu": "250m", "memory": "512Mi"},
limits={"cpu": "500m", "memory": "1Gi"},
),
)
Resource allocation comparison showing over-provisioning versus right-sized Kubernetes pods

ส่วน Concurrency ที่ไม่มีขอบเขต คือทีมเปิด Catchup ทิ้งไว้ หรือไม่ตั้ง Limit อะไรเลย พอ Deploy DAG ใหม่หรือแก้ start_date ปุ๊บ Airflow ก็จะพยายามรัน Backfill ย้อนหลังทีเดียวหลายร้อย Run พร้อมกัน แต่ละ Run ก็ยิง Connection เข้า Database หรือ API ปลายทางพร้อมๆกันจนระบบเป้าหมายรับไม่ไหว

ปัญหาคือสองอย่างนี้พอมาเจอกันจะยิ่งหนักขึ้นไปอีก เพราะ Task ที่ Over-provisioned อยู่แล้ว พอมารันพร้อมกันแบบไม่มี Limit ก็จะกิน Resource เกินจริงไปหลายเท่าตัว กลายเป็น Resource Contention ทั้ง Cluster ไม่ใช่แค่ DAG ตัวเดียวที่ได้รับผลกระทบด้วยครับ

# Bad: ไม่จำกัด Concurrency เลย แถมเปิด Catchup ทิ้งไว้
with DAG(
dag_id="bad_unbounded_concurrency",
schedule="@hourly",
start_date=datetime(2024, 1, 1),
catchup=True,
) as dag:
...
# Good: จำกัดทั้งจำนวน Run พร้อมกัน, Task พร้อมกัน และ Connection ผ่าน Pool
with DAG(
dag_id="good_bounded_concurrency",
schedule="@hourly",
start_date=datetime(2024, 1, 1),
catchup=False,
max_active_runs=1,
max_active_tasks=10,
) as dag:
query_task = PythonOperator(
task_id="query_warehouse",
python_callable=query_fn,
pool="warehouse_connections",
)

5. เรื่องเล็กๆที่ดูไม่สำคัญ แต่พอสะสมหลายร้อย DAG กลายเป็นปัญหาใหญ่

พวกนี้แยกเป็นเรื่องเดียวไม่คุ้มเขียนหัวข้อ แต่พอเจอรวมกันในหลาย DAG พร้อมกันก็สร้างความปวดหัวได้ไม่แพ้ 4 ข้อข้างบนเลยครับ

  • Owner ไม่ชัดเจน DAG ส่วนใหญ่ปล่อย owner เป็น "airflow" ตาม Default พอ DAG พังกลางดึกไม่มีใครรู้ว่าต้อง Ping ใคร
  • Schedule ไม่ Timezone-aware ใช้ datetime.utcnow() หรือ Naive Datetime ตรงๆ พอทีมอยู่คนละ Timezone จะงงกันหมดว่า DAG รันตอนกี่โมงกันแน่ในเวลาท้องถิ่น (โดยเฉพาะทำงานกับบริษัทที่เป็น International)
  • Retry Config ตั้งผิด บางทีตั้ง retries=0 ทำให้ Task ที่พังจาก Network Blip เล็กๆน้อยๆ ล้มทันทีไม่มีโอกาสลองใหม่ หรือตรงข้ามคือตั้ง retry_delay สั้นเกินไปจน Retry รัวใส่ระบบปลายทางที่กำลังมีปัญหาอยู่แล้วซ้ำเข้าไปอีก
# Bad
default_args = {
"owner": "airflow",
"retries": 0,
}
start_date = datetime(2024, 1, 1) # Naive, ไม่มี Timezone
...
# Good
import pendulum
default_args = {
"owner": "team-data-platform",
"retries": 3,
"retry_delay": pendulum.duration(minutes=5),
"retry_exponential_backoff": True,
}
start_date = pendulum.datetime(2024, 1, 1, tz="Asia/Bangkok")
Airflow DAG configuration hygiene showing timezone-aware schedule, owner, and retry backoff

Checklists ตรวจ DAG ของตัวเองก่อน Deploy

ตัวอย่าง Checklist สำหรับให้ทีมใช้ตรวจสอบ DAG

  • ไม่มี Network Call, DB Query หรือ File Read อยู่นอก Task/Operator (ไม่มี Top-level Code หนักๆ)
  • Sensor ทุกตัวเช็คแล้วว่ามี Deferrable Mode ให้ใช้ไหม ถ้ามีให้เปิด deferrable=True
  • แต่ละ Task ทำงานเดียวจบ (Atomic) และรันซ้ำได้โดยไม่ทำข้อมูลซ้ำหรือพัง (Idempotent)
  • Resource Request (CPU/Memory) ของแต่ละ Task ผ่านการ Profile Usage จริงแล้ว ไม่ใช่ตั้งกันเผื่อไว้ก่อนแบบเดาสุ่ม
  • ตั้ง max_active_runs, max_active_tasks และใช้ Pool คุม Connection ที่จำกัดแล้ว
  • owner ใน default_args เป็นชื่อทีมหรือ Alias ที่ Ping ได้จริง ไม่ใช่ "airflow" เฉยๆ
  • start_date เป็น Timezone-aware Object (เช่นผ่าน pendulum) ไม่ใช่ Naive Datetime
  • retries และ retry_delay ตั้งเหมาะสม ไม่ใช่ 0 และไม่ถี่จนซ้ำเติมระบบปลายทาง
  • catchup ปิดไว้ (False) เว้นแต่ตั้งใจจะ Backfill จริงๆ และรู้ว่าต้องจำกัด Concurrency ตอน Backfill ด้วย

Summary

5 ข้อนี้ไม่มีอะไรซับซ้อนเลยสักข้อ แต่เป็นสิ่งที่หลุดง่ายที่สุดเพราะทีมมักมอง DAG File เหมือน Python Script ธรรมดาที่รันครั้งเดียวจบ ทั้งที่จริงๆ Scheduler มัน Parse ซ้ำตลอดเวลา และ Task แต่ละตัวก็มีโอกาสถูก Retry ได้เสมอ

ถ้าจะให้เรียงลำดับว่าควรแก้อะไรก่อน ผมมองว่า ข้อ 1 (Top-level Code) กับข้อ 2 (Sensor แบบเก่า) ควรแก้ก่อนสุด เพราะสองอย่างนี้กระทบ Performance ของ Scheduler และ Cluster ทั้งก้อน ไม่ใช่แค่ DAG ตัวเดียวที่มีปัญหา ส่วนเรื่อง Atomic/Idempotent กับ Concurrency ก็ตามมาติดๆเพราะเกี่ยวกับความถูกต้องของข้อมูลโดยตรง ส่วนข้อ 5 ที่ดูเป็นเรื่อง Hygiene เล็กๆนั้น อย่ามองข้าม เพราะพอมี DAG เป็นร้อยเป็นพันตัวในองค์กร ความไม่เป๊ะเล็กๆพวกนี้คือสิ่งที่ทำให้ทีม On-call ปวดหัวที่สุดตอนกลางดึกครับ


ตัวอย่าง DAG แบบเต็มที่รวมการแก้ทั้ง 5 ข้อไว้ใน Github Repo นี้เลย

https://github.com/kriangsak-puk/top-airflow-common-mistake

REF:https://airflow.apache.org/docs/apache-airflow/stable/authoring-and-scheduling/deferring.html
https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/sensors.html
https://airflow.apache.org/docs/apache-airflow/stable/best-practices.html

Leave a comment