Data Pipeline dengan Python — Data Engineering

Membangun Data Pipeline dengan Python Data pipeline adalah serangkaian langkah otomatis yang mengambil data dari sumber, mentransformasi, dan memuatnya ke…

Membangun Data Pipeline dengan Python

Data pipeline adalah serangkaian langkah otomatis yang mengambil data dari sumber, mentransformasi, dan memuatnya ke tujuan. Python adalah bahasa utama untuk membangun pipeline karena ekosistem library-nya yang kaya.

Struktur Pipeline Sederhana

import pandas as pd
import requests
from sqlalchemy import create_engine
from datetime import datetime
import logging

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class SalesPipeline:
    def __init__(self, db_url: str):
        self.engine = create_engine(db_url)

    def extract_from_api(self, date: str) -> list[dict]:
        """Extract: ambil data dari API"""
        logger.info(f"Extracting data for {date}")
        response = requests.get(
            f"https://api.example.com/sales?date={date}",
            headers={"Authorization": "Bearer xxx"}
        )
        response.raise_for_status()
        return response.json()["data"]

    def extract_from_csv(self, path: str) -> pd.DataFrame:
        """Extract: baca dari file CSV"""
        return pd.read_csv(path)

    def transform(self, raw_data: list[dict]) -> pd.DataFrame:
        """Transform: bersihkan dan enrich data"""
        df = pd.DataFrame(raw_data)

        # Drop duplicates
        df = df.drop_duplicates(subset=["transaction_id"])

        # Type casting
        df["amount"] = pd.to_numeric(df["amount"], errors="coerce")
        df["created_at"] = pd.to_datetime(df["created_at"])

        # Filter invalid
        df = df[df["amount"] > 0]

        # Enrich
        df["month"] = df["created_at"].dt.to_period("M")
        df["processed_at"] = datetime.now()

        logger.info(f"Transformed {len(df)} valid records")
        return df

    def load(self, df: pd.DataFrame, table: str):
        """Load: simpan ke database"""
        df.to_sql(table, self.engine, if_exists="append", index=False)
        logger.info(f"Loaded {len(df)} rows to {table}")

    def run(self, date: str):
        """Jalankan full pipeline"""
        try:
            raw = self.extract_from_api(date)
            cleaned = self.transform(raw)
            self.load(cleaned, "fact_sales")
            logger.info("Pipeline completed successfully")
        except Exception as e:
            logger.error(f"Pipeline failed: {e}")
            raise

Error Handling & Retry

from tenacity import retry, stop_after_attempt, wait_exponential

class RobustPipeline:
    @retry(stop=stop_after_attempt(3), wait=wait_exponential(min=1, max=60))
    def extract(self, url: str) -> dict:
        """Retry otomatis dengan exponential backoff"""
        response = requests.get(url, timeout=30)
        response.raise_for_status()
        return response.json()

    def run_with_checkpoint(self, dates: list[str]):
        """Track progress — resume dari checkpoint jika gagal"""
        processed = self.load_checkpoint()
        for date in dates:
            if date in processed:
                continue
            self.process_date(date)
            self.save_checkpoint(date)

Best Practices

Yang akan kamu pelajari