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

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

มีโอกาสได้ไล่รีวิว 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 requestsfrom airflow import DAGfrom airflow.operators.python import PythonOperatorfrom datetime import datetimeconfig = 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 requestsfrom airflow import DAGfrom airflow.operators.python import PythonOperatorfrom datetime import datetimedef 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 Architecture จะได้เข้าใจมากขึ้นว่าทำไมการ parse ซ้ำๆของโค้ดที่ไม่ดีจะส่งผลกระทบกับ resource ยังไง โดยจะเห็นว่า Scheduler จะทำงานเป็นทั้งตัวอ่านไฟล์และส่งงานไปที่ worker ทำให้ถ้า DAG parsing ยิ่งนานก็ยิ่งกระทบไปเรื่อยๆ

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 ทั้งก้อนจะโดนล็อกจนรันอย่างอื่นไม่ได้เลยทั้งที่จริงๆแล้วแทบไม่มีอะไรทำงานอยู่เลยครับ

ทางแก้คือใช้ 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 ExternalTaskSensorwait_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 ExternalTaskSensorwait_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 ก็จะได้ข้อมูลซ้ำเข้าไปอีกชุดทันที

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 ตาม Partitionfrom airflow.decorators import tasktaskdef extract(logical_date=None): partition = logical_date.format("YYYY-MM-DD") return extract_from_api(partition=partition)taskdef transform(raw_data): return transform_fn(raw_data)taskdef 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 เกินความจำเป็นจริงของ Taskfrom airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperatorfrom kubernetes.client import models as k8sheavy_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 KubernetesPodOperatorfrom kubernetes.client import models as k8sright_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"}, ),)

ส่วน 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 ผ่าน Poolwith 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 รัวใส่ระบบปลายทางที่กำลังมีปัญหาอยู่แล้วซ้ำเข้าไปอีก
# Baddefault_args = { "owner": "airflow", "retries": 0,}start_date = datetime(2024, 1, 1) # Naive, ไม่มี Timezone...
# Goodimport pendulumdefault_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")

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 Datetimeretriesและ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