PipeCraft

Roadmap Progress

ข้อมูลเบื้องต้นเกี่ยวกับวิศวกรรมข้อมูล (Introduction to Data Engineering) Python สำหรับ Data Engineering พื้นฐาน Linux & CLI สำหรับวิศวกรข้อมูล Git สำหรับระบบงานข้อมูลและทีมพัฒนา พื้นฐานเครือข่าย เว็บเทคโนโลยี และระบบกระจายศูนย์ (Web Fundamentals, Networking & Distributed Systems) SQL ขั้นสูง (Window Functions & Optimization) ฐานข้อมูลไม่ใช่เชิงสัมพันธ์ (NoSQL Databases & Modern Stores) การออกแบบโครงสร้างข้อมูล (Star/Snowflake & SCD) โครงสร้างพื้นฐานเครือข่ายระบบคลาวด์ (Cloud-based Networking) สถาปัตยกรรม Data Lakehouse (Apache Iceberg & MinIO) Local-first Data Engineering ด้วย DuckDB โครงสร้างพื้นฐานในรูปแบบโค้ด (IaC ด้วย Terraform) การแปลงข้อมูลระดับโปรด้วย dbt-core ข้อตกลงร่วมด้านข้อมูล (Data Contracts) การบริการและการส่งต่อข้อมูลวิเคราะห์ (Data Serving & Reverse ETL) การประมวลผลข้อมูลขนาดใหญ่แบบกระจายด้วย Apache Spark ระบบจัดการเวิร์กโฟลว์ข้อมูล (Airflow, Dagster & Prefect) การประมวลผลข้อมูลแบบเรียลไทม์ด้วย Kafka & Redpanda การจัดการคอนเทนเนอร์และคลัสเตอร์ (Containers & Kubernetes) ระบบ CI/CD และการเฝ้าระวังคุณภาพข้อมูล (CI/CD, Monitoring & Testing) ระบบความปลอดภัยและการควบคุมข้อมูล (Security, Governance & Privacy) ระบบปฏิบัติการและประมวลผลโมเดล (Machine Learning & MLOps) รากฐานของ AWS VPC การจัดการเส้นทางขั้นสูงและ NAT Gateway การเชื่อมต่อหลาย VPC และ Hybrid Cloud VPC Endpoints (PrivateLink) ระบบรักษาความปลอดภัยขอบเขตเครือข่าย สถาปัตยกรรมเครือข่ายสำหรับ EKS IPv6 และการจัดการ IP Address สถาปัตยกรรมเครือข่ายระดับโลก แก่นแท้ของสถาปัตยกรรม Kubernetes สถาปัตยกรรมเครื่องและคอมโพเนนต์ เลเยอร์ที่เชื่อมต่อได้ (Pluggable Layers) การรันบนโปรดักชันระดับองค์กร Apache Flink Stateful Stream Processing Real-Time CDC & Event Sourcing at Scale Low-Latency Stream-Table Joins & Windowing Data Lineage & Metadata Graph Engineering Statistical Anomaly Detection & Data Drift Data Incident Management & Automated DLQ Remediation Cloud Data FinOps & Cost Optimization Mechanics Data Mesh & Multi-Tenant Platform Architecture Vector Databases & AI-Ready Data Infrastructure API Fundamentals, Architectures & Core Components API Versioning, Docs, Real-Time & Microservices Production Reliability, Security & Resiliency World-Class Master Architecture & Traffic Management
บทที่ 5: ระบบประมวลผลขนาดใหญ่และการจัดลำดับงาน (Orchestration, Scale & Streaming)

ระบบจัดการเวิร์กโฟลว์ข้อมูล (Airflow, Dagster & Prefect)

5.2 Workflow Orchestration ด้วย Apache Airflow

ระบบจัดการสายงาน (Workflow Orchestration) ทำหน้าที่สแตนด์บายประสานเวลารันและลำดับขั้นจราจรของท่อส่งข้อมูลหลากหลายแบบให้เป็นไปตามเงื่อนไขอย่างราบรื่น

Technical Architecture Diagram
Click to Zoom
3D Isometric Architecture Breakdown:
  • 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 สำเร็จ

extract_data Idle
clean_data Idle
load_warehouse Idle
เลือกคำสั่งจับคู่ทิศทางเวิร์กโฟลว์ (Task Flow Operators):
extract_data clean_data load_warehouse

ข้อดี / จุดเด่น & ข้อเสีย / ข้อควรระวัง

ข้อดี / จุดเด่น

  • สามารถจัดการลำดับความเชื่อมโยงของงานที่สลับซับซ้อนได้อย่างน่าเชื่อถือ
  • มีระบบ User Interface คอยมอนิเตอร์และสั่งทำงานย้อนหลัง (Backfilling) ได้โดยง่าย

ข้อเสีย / ข้อควรระวัง

  • ต้องการทรัพยากรจัดเตรียมระบบสูง (ฐานข้อมูล Metadata, Web Server, Scheduler)
  • มีความยุ่งยากในการพัฒนาและรันชุดทดสอบความถูกต้องของสคริปต์ที่เครื่องตนเอง

ปฏิบัติการจริง (Lab Practice)

Bilingual Guide

แล็บ: สร้างและรันสคริปต์ควบคุม DAG สำหรับท่อส่งข้อมูล

วิธีรันแล็บปฏิบัติการบนเครื่องจริง (Local Terminal Execution Blueprint)

เนื่องจากปฏิบัติการ Data Engineering / DevOps ระดับสูงต้องรันบน Environment จริง ขั้นตอนด้านล่างนี้คือคำสั่งสำหรับนำไปรันบน Terminal / Docker ในเครื่องของคุณ:

1. สร้างโฟลเดอร์ปฏิบัติการและเตรียมไฟล์ Environment
mkdir -p pipecraft-lab && cd pipecraft-lab
python3 -m venv venv && source venv/bin/activate
pip install --upgrade pip pandas polars pytest requests
2. จำลองการสร้างและรันระบบประมวลผล (Run Pipeline Command)
# 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)
" 
3. ตรวจสอบการผ่านเกณฑ์และการทำงาน (Data Quality Assertions)
# Verify clean execution exit status
echo " PipeCraft Local Lab Execution Completed Successfully!"

💡 คำแนะนำ & ทริกเด็ด

ห้ามใช้โหนด Worker ของ Airflow ในการคำนวณหรือดึงข้อมูลปริมาณมากตรงๆ ให้ใช้คำสั่ง Operators ไปสั่งรันงานหนักภายนอก (เช่น KubernetesPodOperator หรือ BigQueryInsertJobOperator) แทน เพื่อไม่ให้ Airflow Scheduler แฮงก์