ระบบจัดการเวิร์กโฟลว์ข้อมูล (Airflow, Dagster & Prefect)
5.2 Workflow Orchestration ด้วย Apache Airflow
ระบบจัดการสายงาน (Workflow Orchestration) ทำหน้าที่สแตนด์บายประสานเวลารันและลำดับขั้นจราจรของท่อส่งข้อมูลหลากหลายแบบให้เป็นไปตามเงื่อนไขอย่างราบรื่น
- Step 1: Data Ingestion & Event Processing
- Step 2: Distributed Computation & State Management
- Step 3: Orchestration & Resource Allocation
1. Airflow Architecture: Scheduler, Webserver, และ Executors (Celery vs K8s)
Learning Progression
- [BASIC] ปูพื้นฐานภาษาเข้าใจง่าย - เข้าใจคอนเซปต์ภาพรวมและการแก้ปัญหาเบื้องต้น
- [INTERMEDIATE] โค้ด/คอนฟิกไวยากรณ์จริง - การเขียนโค้ดเพื่อใช้งานจริงในระบบ
- [PROFESSIONAL] Under-the-hood & Performance/FinOps - กลไกเบื้องลึกและการรีดประสิทธิภาพ
Real-World Enterprise Scenario
เคสระบบการเงิน/Big Tech: การรองรับ Transaction จำนวนมหาศาลต่อวินาทีพร้อมประกัน Data Integrity สูงสุด โดยใช้สถาปัตยกรรมที่ยืดหยุ่นและการมอนิเตอร์ระดับสูง
ทฤษฎีและกลไกการทำงาน (How it works): Apache Airflow มีสถาปัตยกรรมแบบแยกส่วน: - **Webserver**: ให้บริการหน้า UI สำหรับผู้ใช้งานเพื่อดูสถานะ DAGs และจัดการระบบ - **Scheduler**: สแกนไฟล์ DAG อย่างต่อเนื่อง ตัดสินใจว่า Task ไหนถึงเวลารัน และบันทึกสถานะลง **Metadata Database (เช่น PostgreSQL)** - **Executors**: กลไกการรัน Task ที่รับคำสั่งจาก Scheduler - **CeleryExecutor**: ใช้ Message Broker (Redis/RabbitMQ) กระจายงานให้ Worker nodes ที่รันรออยู่ตลอดเวลา เหมาะกับงานที่คงที่ - **KubernetesExecutor**: สร้าง Pod ใหม่ขึ้นมาทำงานแต่ละ Task แยกกันโดยเฉพาะ (Dynamic provisioning) คืนทรัพยากรทันทีเมื่อจบงาน เหมาะกับงานที่มีไลบรารีต่างกันและต้องการแยก Resource ชัดเจน
# airflow.cfg - Executor configuration
[core]
executor = KubernetesExecutor
# executor = CeleryExecutor
Use Case ในชีวิตจริง (Real-world Scenario): บริษัทอีคอมเมิร์ซที่มี Data Pipeline แตกต่างกัน 500 เส้น เลือกใช้ KubernetesExecutor เพื่อแยกการจัดการ Library (เช่น เส้น A ใช้ Pandas 1.0, เส้น B ใช้ Pandas 2.0) โดยที่ Pod จะถูกสร้างและทำลายทิ้งทันทีเมื่อ Task เสร็จสิ้น ทำให้ไม่เปลืองทรัพยากร
ข้อควรระวังและวิธีแก้ (Pitfalls & Mitigations): หากสเกล Worker ไม่พอใน CeleryExecutor งานจะไปกองที่คิว (Queue buildup) จนเกิดคอขวด แก้ไขโดยการตั้งค่าการสเกล Worker อัตโนมัติด้วย KEDA หรือย้ายไปใช้ KubernetesExecutor ที่รองรับการสเกลตามโควต้าของคลัสเตอร์
2. Idempotency และกลไกการทำ Backfilling อย่างปลอดภัย
ทฤษฎีและกลไกการทำงาน (How it works): ท่อส่งที่ดีต้องมีคุณสมบัติ **Idempotent** หมายความว่ารันโค้ดประมวลผลข้อมูลกี่สิบรอบด้วยพารามิเตอร์วันเดิม จะต้องได้ผลลัพธ์ปลายทางคงเดิมเสมอ ห้ามเกิดปัญหายอดข้อมูลซ้อนหรือผลลัพธ์งอกเพิ่ม นอกจากนี้ Airflow รองรับระบบ **Backfill** โดยเราสามารถระบุสั่งรันย้อนเวลาตามประวัติวันเวลารัน (Logical Date / Execution Date) ย้อนหลังแบบอัตโนมัติได้เมื่อต้องการเก็บข้อมูลย้อนหลัง
# Idempotent Airflow DAG with logical execution date filters
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime
with DAG(
"idempotent_etl_pipeline",
start_date=datetime(2026, 7, 1),
schedule_interval="@daily",
catchup=False # Prevents automatic run of unexecuted past DAG runs on deployment
) as dag:
# Process only partition corresponding to the execution date using Jinja templates ({{ ds }})
process_daily_data = BashOperator(
task_id="run_pipeline_partition",
bash_command="python3 /opt/etl.py --date {{ ds }}"
)
Use Case ในชีวิตจริง (Real-world Scenario): ระบบวิเคราะห์ยอดขายพบคิวรีล่มในประวัติ 5 วันที่แล้ว วิศวกรข้อมูลสามารถทำ Backfill เฉพาะข้อมูลพาร์ทิชัน 5 วันนั้นใหม่ได้อย่างปลอดภัย โดยสคริปต์จะใช้ตัวแปร `{{ ds }}` ในการสั่ง Delete และ Insert ข้อมูลเก่าทิ้งโดยไม่ทำลายโครงสร้างตารางวันอื่น
ข้อควรระวังและวิธีแก้ (Pitfalls & Mitigations): การกำหนดวันเวลารันโดยดึงค่า `datetime.now()` ในเนื้อโค้ด Python โดยตรง จะทำให้คุณสมบัติ Idempotency และระบบ Backfill ใช้งานไม่ได้ เพราะมันจะได้เวลาปัจจุบันเสมอ แก้ไขโดยส่งผ่านตัวแปรแฝงระบบของ Airflow `{{ ds }}` เข้าไปในพารามิเตอร์เสมอ
3. Dynamic Task Mapping (ฟีเจอร์เด่นเพื่อการสเกล Task อัตโนมัติ)
ทฤษฎีและกลไกการทำงาน (How it works): Dynamic Task Mapping เป็นฟีเจอร์ (ใน Airflow 2.3+) ที่อนุญาตให้สร้าง Task ย่อยแบบไดนามิกตอนรันไทม์ (Runtime) ตามข้อมูลหรืออาร์เรย์ที่ถูกส่งมาจาก Task ก่อนหน้า โดยใช้ฟังก์ชัน `.expand()` และ `.partial()` ซึ่งแตกต่างจากการใช้ `for loop` ธรรมดาตอนนิยาม DAG ที่จะคงที่ (Static) ไปตลอด
# Dynamic Task Mapping Example
from airflow.decorators import dag, task
from datetime import datetime
@dag(start_date=datetime(2026, 1, 1), schedule_interval=None)
def dynamic_mapping_dag():
@task
def get_files():
return ["file1.csv", "file2.csv", "file3.csv"] # Returns list at runtime
@task
def process_file(file_name: str):
print(f"Processing {file_name}")
# Dynamically maps a task for each file in the list returned by get_files
process_file.expand(file_name=get_files())
dynamic_mapping_dag()
Use Case ในชีวิตจริง (Real-world Scenario): สคริปต์สแกนตรวจสอบว่าวันนี้มีไฟล์ลูกค้าอัปโหลดเข้ามาใน S3 กี่ไฟล์ หากมี 10 ไฟล์ Airflow จะสร้าง 10 Tasks คู่ขนานกันไปรันให้เองทันที หากมี 0 ไฟล์ก็จะข้ามไป ทำให้มีความยืดหยุ่นสูงมากโดยไม่ต้องแก้โค้ด DAG
ข้อควรระวังและวิธีแก้ (Pitfalls & Mitigations): หาก Task ก่อนหน้าส่งลิสต์ที่มีสมาชิกจำนวนมหาศาล (เช่น 10,000 ไฟล์) จะทำให้ Scheduler และ Metadata Database สร้าง Task ล้นจนช้า (Overload) แก้ไขโดยตั้งค่า `max_map_length` หรือจับกลุ่ม (Chunk) ลิสต์ให้มีขนาดพอเหมาะ (เช่น กลุ่มละ 500 ไฟล์)
4. การสร้าง Custom Operators และการเลือกใช้ KubernetesPodOperator
ทฤษฎีและกลไกการทำงาน (How it works): - **Custom Operators**: หากตรรกะในธุรกิจมีความเฉพาะตัว การสืบทอด (Inherit) คลาส `BaseOperator` มาเขียนพฤติกรรมเองจะช่วยลดโค้ดซ้ำซ้อน - **KubernetesPodOperator (KPOs)**: คือ Operator มาตรฐานสำหรับการสั่งสร้างกล่องคอนเทนเนอร์ Pod แยกต่างหากบนคลัสเตอร์ Kubernetes คัดแยกความเสี่ยงของรหัสหลุดล่มออกไปโดยสิ้นเชิง
# Creating a Custom Operator
from airflow.models import BaseOperator
class MyAPIETLOperator(BaseOperator):
def __init__(self, api_endpoint, *args, **kwargs):
super().__init__(*args, **kwargs)
self.api_endpoint = api_endpoint
def execute(self, context):
# Specific business logic to extract and load
print(f"Extracting data from {self.api_endpoint} for date {context['ds']}")
Use Case ในชีวิตจริง (Real-world Scenario): ทีมวิศวกรสร้าง `CompanySecurityScanOperator` เป็น Custom Operator เพื่อบังคับให้ทุกท่อส่งข้อมูลต้องเรียกใช้เช็คความปลอดภัยก่อนดึงข้อมูล ทำให้การใช้งานง่ายและเป็นมาตรฐานเดียวกันทั้งบริษัท
ข้อควรระวังและวิธีแก้ (Pitfalls & Mitigations): การเริ่มต้นเปิดคอนเทนเนอร์ Pod ใหม่ด้วย KPOs หรือการรัน Custom Operator ที่มีขนาดใหญ่มากในทุก Task ย่อยสร้างความล่าช้า แก้ไขโดยหลีกเลี่ยง KPOs กับงานประเภทประมวลผลขนาดสั้นต่ำกว่าวินาที
5. กลไกกู้ภัยสายงานล้มเหลว (Retries/Backoff) และ Cross-DAG Dependencies
ทฤษฎีและกลไกการทำงาน (How it works): เพื่อเสถียรภาพสูงสุด: - **Retries & Exponential Backoff**: ตั้งค่าให้รันแก้ตัวและหน่วงเวลาก่อนสู้ใหม่เป็นทวีคูณเมื่อภารกิจพังเสียหาย - **Cross-DAG Dependencies**: ท่อส่งข้อมูลที่ต้องรันต่อกันข้ามไฟล์ สามารถเชื่อมโยงผ่านเครื่องมือ `TriggerDagRunOperator` หรือเฝ้ารอคอยความสมบูรณ์ผ่าน `ExternalTaskSensor`
# Triggering another pipeline file using TriggerDagRunOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
trigger_reporting_pipeline = TriggerDagRunOperator(
task_id="trigger_downstream_report",
trigger_dag_id="financial_reporting_dashboard_pipeline",
execution_date="{{ execution_date }}",
reset_dag_run=True,
wait_for_completion=False
)
Use Case ในชีวิตจริง (Real-world Scenario): ท่อส่งโหลด API ข้อมูลหุ้นล้มเหลวเนื่องจากเซิร์ฟเวอร์ปลายทางขัดข้องชั่วขณะ การตั้ง Retry Exponential Backoff ช่วยให้ระบบพักรอ 5 นาทีและกลับไปดึงใหม่จนผ่านได้ จากนั้นใช้ `TriggerDagRunOperator` ยิงปลุกท่อรันโมเดล AI ให้ทำงานต่อ
ข้อควรระวังและวิธีแก้ (Pitfalls & Mitigations): ปัญหาวงจรรันระบบรอสายค้างคา (Deadlocks) เมื่อ Sensor เฝ้ารอคอยสายงานอื่นนานเกินไปแย่งโควต้า Workers รันค้างเต็มระบบ แก้ไขโดยตั้งค่าโหมดประหยัดทรัพยากร `mode='reschedule'` เพื่อให้ Sensor ปล่อยสล็อต Workers ไปทำจ๊อบอื่นระหว่างรอ
Weekend Sandbox Challenge: Build an Idempotent DAG with Dynamic Mapping
โจทย์ปฏิบัติการ: จงเขียนไฟล์ Airflow DAG ที่สาธิตการใช้ Dynamic Task Mapping ดึงข้อมูลไฟล์เพื่อนำไปประมวลผลแบบคู่ขนาน และใช้ `catchup=False` เพื่อความปลอดภัย
# local_airflow_dynamic_dag.py
from airflow.decorators import dag, task
from datetime import datetime, timedelta
default_args = {
'owner': 'data_ops',
'retries': 2,
'retry_delay': timedelta(seconds=30),
}
@dag(start_date=datetime(2026, 7, 1), schedule_interval='@daily', catchup=False, default_args=default_args)
def dynamic_etl_pipeline():
@task
def extract_file_list(execution_date=None):
print(f"Finding files for date: {execution_date}")
return ["file_A.parquet", "file_B.parquet"] # Mock dynamic discovery
@task
def process_file(file_name: str):
print(f"Transforming {file_name}")
# Map the task dynamically over the list
files = extract_file_list(execution_date="{{ ds }}")
process_file.expand(file_name=files)
dynamic_etl_pipeline()
Senior Technical Interview Q&A
Q1: Dynamic Task Mapping ใน Airflow 2.3+ ต่างจากการใช้ For Loop ปกติใน DAG อย่างไร?
A1: การใช้ For Loop ปกติตอนเขียน DAG จะถูกประเมินและยึดเป็นโครงสร้างตายตัว (Static) โดย Scheduler ตอนสแกนไฟล์ ซึ่งไม่สามารถปรับจำนวน Task ตามข้อมูลในแต่ละรอบการรันได้ แต่ Dynamic Task Mapping อนุญาตให้สร้าง Task อัตโนมัติตามจำนวนข้อมูลที่ถูกส่งออกมาจาก Task ก่อนหน้าในขณะรันไทม์ (Runtime) ได้อย่างยืดหยุ่น
Q2: จงวิเคราะห์ความแตกต่างระหว่างการรันโปรเซสด้วย CeleryExecutor และ KubernetesExecutor ในระบบ Airflow ขนาดใหญ่?
A2: CeleryExecutor ต้องการเปิด Worker สแตนด์บายตลอดเวลา มีข้อจำกัดในการสเกลเครื่องและเสี่ยงต่อปัญหาไลบรารีโค้ดขัดแย้งข้ามจ๊อบ ในขณะที่ KubernetesExecutor จะสั่งเปิด Pod ใหม่ขนานบนคลัสเตอร์ขึ้นมารันงานรายจ๊อบย่อยเฉพาะกิจและทำลายโหนดทิ้งทันทีเมื่อจบงาน คัดแยกสภาพแวดล้อม (Isolation) และจัดการทรัพยากรได้ดีกว่ามากในคลัสเตอร์ Kubernetes
Interactive DAG Workflow Orchestrator
Interactive DAGเป้าหมาย: กำหนดลำดับการทำงาน (Dependencies) ของเวิร์กโฟลว์ข้อมูลให้ถูกต้อง โดยงาน Clean จะทำงานหลังจาก Extract สำเร็จ และงาน Load จะทำงานหลังจาก Clean สำเร็จ
ข้อดี / จุดเด่น & ข้อเสีย / ข้อควรระวัง
ข้อดี / จุดเด่น
- สามารถจัดการลำดับความเชื่อมโยงของงานที่สลับซับซ้อนได้อย่างน่าเชื่อถือ
- มีระบบ User Interface คอยมอนิเตอร์และสั่งทำงานย้อนหลัง (Backfilling) ได้โดยง่าย
ข้อเสีย / ข้อควรระวัง
- ต้องการทรัพยากรจัดเตรียมระบบสูง (ฐานข้อมูล Metadata, Web Server, Scheduler)
- มีความยุ่งยากในการพัฒนาและรันชุดทดสอบความถูกต้องของสคริปต์ที่เครื่องตนเอง
ปฏิบัติการจริง (Lab Practice)
แล็บ: สร้างและรันสคริปต์ควบคุม DAG สำหรับท่อส่งข้อมูล
วิธีรันแล็บปฏิบัติการบนเครื่องจริง (Local Terminal Execution Blueprint)
เนื่องจากปฏิบัติการ Data Engineering / DevOps ระดับสูงต้องรันบน Environment จริง ขั้นตอนด้านล่างนี้คือคำสั่งสำหรับนำไปรันบน Terminal / Docker ในเครื่องของคุณ:
mkdir -p pipecraft-lab && cd pipecraft-lab python3 -m venv venv && source venv/bin/activate pip install --upgrade pip pandas polars pytest requests
# Run execution pipeline test
python3 -c "
import polars as pl
print(' [PipeCraft Lab] Running local pipeline engine...')
df = pl.DataFrame({'id': [1, 2, 3], 'status': ['SUCCESS', 'SUCCESS', 'AUDITED']})
print(df)
"
# Verify clean execution exit status echo " PipeCraft Local Lab Execution Completed Successfully!"
💡 คำแนะนำ & ทริกเด็ด
ห้ามใช้โหนด Worker ของ Airflow ในการคำนวณหรือดึงข้อมูลปริมาณมากตรงๆ ให้ใช้คำสั่ง Operators ไปสั่งรันงานหนักภายนอก (เช่น KubernetesPodOperator หรือ BigQueryInsertJobOperator) แทน เพื่อไม่ให้ Airflow Scheduler แฮงก์