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
- Idempotent — Pipeline bisa di-rerun tanpa duplikasi data (gunakan UPSERT / delete-then-insert)
- Logging — Log setiap tahap untuk debugging
- Monitoring — Track jumlah rows, duration, error rate
- Config-driven — Jangan hardcode URL, credentials, atau parameter
- Testing — Unit test untuk transformasi logic