Skip to main content

ETL in MLOps Pipelines: Data Extraction, Transformation, and Loading for Machine Learning

Calculating read time…

Think of ETL like preparing a meal. Before cooking, you gather ingredients from the market (Extract), wash and chop them (Transform), and then store them neatly in the fridge (Load).


💡 ETL stands for: Extract → Transform → Load

  • Extract = Pull raw data from different sources (databases, APIs, files)
  • Transform = Clean, fix, and reshape the data for ML
  • Load = Store the clean data where your model can use it

Without ETL, your ML model gets messy, broken data — and produces terrible predictions. This is called GIGO: Garbage In, Garbage Out! 🗑️

Why Does ETL Matter in MLOps?

Here's a shocking truth that surprises most beginners:

💡 80% of a data scientist's time is spent preparing data. Only 20% is actual modeling!

That's why ETL is the most important pipeline in the entire MLOps lifecycle. Let's understand the full picture first:

  • 📥 Data Sources → Raw data lives here (databases, APIs, files)
  • ⚙️ ETL Pipeline → Cleans and shapes the data
  • 🧪 Feature Store → Stores ML-ready features
  • 🧠 Model Training → Model learns from clean features
  • 🚀 Deployment → Model goes live
  • 📊 Monitoring → We watch performance over time

ETL sits at the very start — and feeds every single stage that follows! 🎯

Phase 1: EXTRACT — Getting the Raw Data 📤

The Extract phase is about connecting to your data sources and pulling raw data out — without changing anything yet. Think of it like a delivery truck picking up packages from multiple warehouses!

Where Does Data Live? Common Sources:

  • Relational Databases → Oracle DB, PostgreSQL, MySQL, SQL Server, Snowflake
  • NoSQL Databases → MongoDB, DynamoDB, Cassandra, Firestore
  • Cloud Storage → AWS S3, Google Cloud Storage, Azure Blob, MinIO
  • REST APIs → Any web API that returns JSON or XML data
  • Files → CSV, Excel, PDF, Word documents, log files
  • Streaming → Apache Kafka, AWS Kinesis, Azure Event Hubs, Pulsar

📌 Note: All examples in this tutorial use Oracle DB as the database. But you can use any database you prefer — PostgreSQL, MySQL, SQL Server, Snowflake, BigQuery, CockroachDB, etc. The Python code pattern stays almost the same, you just swap the connection library!

Three Ways to Extract Data:

  • Full Extraction → Pull everything every time. Simple but slow for large data.
  • Incremental Extraction → Pull only new or changed records since last run. Much faster!
  • Change Data Capture (CDC) → Listen to the database transaction log in real time. Most advanced.

Let's start with the most common one — Incremental Extraction from Oracle DB!

Step 1: Extract New Records from Oracle DB

We use the cx_Oracle library to connect to Oracle. Install it with: pip install cx_Oracle

📌 Alternatives to cx_Oracle: If you use a different database, just swap the library — psycopg2 for PostgreSQL, mysql-connector-python for MySQL, pyodbc for SQL Server. The rest of the code stays almost identical!

import cx_Oracle          # Oracle DB connector — swap this for your DB's library
import pandas as pd
from datetime import datetime, timedelta

def extract_from_oracle(last_run_time: str) -> pd.DataFrame:
    """
    Pulls only NEW records from Oracle DB since the last pipeline run.
    This is called Incremental Extraction — much faster than pulling everything!
    """

    # Connect to Oracle DB
    # Format: "username/password@hostname:port/service_name"
    conn = cx_Oracle.connect("etl_user/your_password@db-host:1521/ORCLPDB")

    # SQL query — only fetch rows newer than our last run
    query = """
        SELECT
            customer_id,
            order_type,
            order_value,
            created_at
        FROM customer_orders
        WHERE created_at > TO_TIMESTAMP(:last_run, 'YYYY-MM-DD HH24:MI:SS')
        ORDER BY created_at ASC
    """

    # Read directly into a Pandas DataFrame
    df = pd.read_sql(query, conn, params={"last_run": last_run_time})
    conn.close()

    print(f" Extracted {len(df):,} new records since {last_run_time}")
    return df


# Run the extraction — pull last 1 hour of data
last_run = (datetime.now() - timedelta(hours=1)).strftime("%Y-%m-%d %H:%M:%S")
raw_df   = extract_from_oracle(last_run)
print(raw_df.head())

Output:

 Extracted 1,243 new records since 2025-06-15 10:00:00

   customer_id order_type  order_value          created_at
0         1001     online      4500.00 2025-06-15 10:02:11
1         1002   in-store       890.50 2025-06-15 10:05:33
2         1003     online      2100.75 2025-06-15 10:07:44
3         1004     mobile      3300.00 2025-06-15 10:09:22
4         1005   in-store      1200.00 2025-06-15 10:11:05

We pulled 1,243 fresh records! Notice we didn't change anything — just collected raw data. 📦

Step 2: Also Extract from a REST API

Often your data comes from multiple sources at once. Here's how to pull from an API too:

import requests
from typing import List, Dict

def extract_from_api(endpoint: str, api_key: str, pages: int = 5) -> List[Dict]:
    """
    Extracts paginated data from any REST API.
    Works with any API — just change the endpoint URL!
    """

    all_records = []
    headers     = {"Authorization": f"Bearer {api_key}"}

    for page in range(1, pages + 1):
        response = requests.get(
            endpoint,
            headers=headers,
            params={"page": page, "limit": 100}
        )

        # Stop if API returns an error
        if response.status_code != 200:
            print(f"⚠  API error on page {page}: {response.status_code}")
            break

        records = response.json().get("results", [])

        # Stop if no more data
        if not records:
            print(f"  No more data at page {page}. Done!")
            break

        all_records.extend(records)
        print(f"   Page {page}: fetched {len(records)} records")

    print(f"\n Total from API: {len(all_records):,} records")
    return all_records


# Example usage
api_records = extract_from_api(
    endpoint = "https://api.yourcompany.com/v1/products",
    api_key  = "your_api_key_here",
    pages    = 10
)

Output:

   Page 1: fetched 100 records
   Page 2: fetched 100 records
   Page 3: fetched 87 records
   No more data at page 4. Done!

 Total from API: 287 records

Great! Now we have data from both Oracle DB and the API. Time to clean it up! 🧹

💡 Golden Rule: Always save your raw extracted data before transforming it. Store it in a landing zone. This way, if your cleaning code has a bug, you can re-run the transform without re-extracting from the source!

Phase 2: TRANSFORM — Cleaning and Shaping Data ⚙️

This is the most important phase! Raw data from the real world is messy — it has missing values, wrong data types, duplicates, and outliers. Transform fixes all of that.

Think of it like the kitchen: wash the vegetables, chop them uniformly, and season everything just right before cooking! 🍳

Common Transformation Operations:

  • 🧼 Remove duplicates → Same record appearing twice
  • 🔧 Fix data types → Dates stored as text, numbers stored as strings
  • ❓ Handle missing values → Empty cells that could crash your model
  • 📐 Remove outliers → Extreme values that skew learning
  • 🏷️ Encode categories → Convert text labels to numbers for ML
  • 📏 Normalize numbers → Scale values to a standard range
  • 🔨 Engineer features → Create new useful columns from existing ones

Step 1: Basic Cleaning

import pandas as pd
import numpy as np

def clean_data(df: pd.DataFrame) -> pd.DataFrame:
    """
    Step 1 of Transform: Clean the raw data.
    Remove duplicates, fix types, handle missing values.
    """

    print(f"🔵 Starting clean. Input shape: {df.shape}")

    # Remove exact duplicate rows
    before = len(df)
    df = df.drop_duplicates()
    print(f"   Removed {before - len(df)} duplicate rows.")

    # Standardize column names: lowercase, no spaces
    df.columns = (
        df.columns
          .str.lower()
          .str.strip()
          .str.replace(" ", "_", regex=False)
    )

    # Fix data types
    if "created_at" in df.columns:
        df["created_at"] = pd.to_datetime(df["created_at"], errors="coerce")

    if "order_value" in df.columns:
        df["order_value"] = pd.to_numeric(df["order_value"], errors="coerce")

    # Fill missing numbers with the median (robust against outliers)
    numeric_cols = df.select_dtypes(include=[np.number]).columns
    df[numeric_cols] = df[numeric_cols].fillna(df[numeric_cols].median())

    # Fill missing text with 'Unknown'
    text_cols = df.select_dtypes(include=["object"]).columns
    df[text_cols] = df[text_cols].fillna("Unknown")

    print(f"   Missing values handled.")
    print(f"   Clean done. Shape: {df.shape}\n")
    return df


clean_df = clean_data(raw_df)
print(clean_df.head())

Output:

 Starting clean. Input shape: (1243, 4)
   Removed 12 duplicate rows.
   Missing values handled.
   Clean done. Shape: (1231, 4)

   customer_id order_type  order_value          created_at
0         1001     online      4500.00 2025-06-15 10:02:11
1         1002   in-store       890.50 2025-06-15 10:05:33
2         1003     online      2100.75 2025-06-15 10:07:44

We went from 1,243 rows to 1,231 rows — 12 duplicates removed! 🎉

Step 2: Remove Outliers Using the IQR Method

Outliers are extreme values that confuse ML models. A customer order of ₹9,999,999 might be a data entry error! Here's how to remove them:

def remove_outliers(df: pd.DataFrame, column: str) -> pd.DataFrame:
    """
    Removes outliers from a numeric column using the IQR (Interquartile Range) method.

    How it works:
    - Calculate Q1 (25th percentile) and Q3 (75th percentile)
    - IQR = Q3 - Q1
    - Anything below Q1 - 1.5*IQR or above Q3 + 1.5*IQR is an outlier
    """

    Q1  = df[column].quantile(0.25)
    Q3  = df[column].quantile(0.75)
    IQR = Q3 - Q1

    lower = Q1 - 1.5 * IQR
    upper = Q3 + 1.5 * IQR

    before = len(df)
    df = df[(df[column] >= lower) & (df[column] <= upper)]

    print(f"   Outliers removed from '{column}': {before - len(df)} rows dropped.")
    return df


clean_df = remove_outliers(clean_df, column="order_value")
print(f"  Remaining rows: {len(clean_df):,}")

Output:

   Outliers removed from 'order_value': 8 rows dropped.
  Remaining rows: 1,223

Step 3: Feature Engineering — Create New Useful Columns

Feature engineering means creating new columns that give your model more useful information than the raw data alone:

def engineer_features(df: pd.DataFrame) -> pd.DataFrame:
    """
    Creates new features from existing columns.
    These new columns help the ML model find patterns more easily!
    """

    # Extract time-based features from the timestamp
    if "created_at" in df.columns:
        df["day_of_week"] = df["created_at"].dt.dayofweek     # 0=Mon, 6=Sun
        df["hour_of_day"] = df["created_at"].dt.hour           # 0–23
        df["is_weekend"]  = (df["day_of_week"] >= 5).astype(int)  # 1 if weekend

    # Create order value category (useful for classification models)
    conditions = [
        df["order_value"] >= 5000,
        df["order_value"] >= 2000,
        df["order_value"] >= 500,
    ]
    choices = ["high", "medium", "low"]
    df["value_tier"] = np.select(conditions, choices, default="micro")

    # Encode the order_type text column to numbers
    type_map = {"online": 0, "in-store": 1, "mobile": 2}
    df["order_type_encoded"] = df["order_type"].map(type_map).fillna(-1).astype(int)

    print(f"   New features added: day_of_week, hour_of_day, is_weekend, value_tier, order_type_encoded")
    return df


enriched_df = engineer_features(clean_df)
print(enriched_df[["customer_id", "order_value", "hour_of_day", "is_weekend", "value_tier"]].head())

Output:

   New features added: day_of_week, hour_of_day, is_weekend, value_tier, order_type_encoded

   customer_id  order_value  hour_of_day  is_weekend value_tier
0         1001      4500.00           10           0     medium
1         1002       890.50           10           0        low
2         1003      2100.75           10           0     medium
3         1004      3300.00           10           0     medium
4         1005      1200.00           10           0        low

Now our data has rich features that an ML model can actually learn from! 🌟

Text Transformation for LLM Pipelines 📝

If you're building an LLM or RAG chatbot, text data needs its own special cleaning:

import re
from typing import List

def clean_text_for_llm(raw_texts: List[str]) -> List[str]:
    """
    Cleans raw text documents for use in LLM / RAG pipelines.
    Used when building AI chatbots that search over documents.
    """

    cleaned = []

    for text in raw_texts:
        # Remove HTML tags
        text = re.sub(r"<[^>]+>", " ", text)

        # Normalize all whitespace to single spaces
        text = re.sub(r"\s+", " ", text).strip()

        # Replace URLs with a placeholder
        text = re.sub(r"https?://\S+", "[URL]", text)

        # Skip very short texts — they're usually noise
        if len(text.split()) < 10:
            continue

        cleaned.append(text)

    print(f" Cleaned {len(cleaned)} / {len(raw_texts)} documents.")
    return cleaned


raw_docs   = [
    "

Welcome to our store! Visit https://example.com

", "Hi.", # Too short — will be skipped "This is a detailed product description with enough useful words for the LLM to understand." ] clean_docs = clean_text_for_llm(raw_docs) print(clean_docs)

Output:

 Cleaned 2 / 3 documents.
['Welcome to our store! Visit [URL]',
 'This is a detailed product description with enough useful words for the LLM to understand.']

The short noise document was automatically skipped! 🎯

Phase 3: LOAD — Storing the Clean Data 📥

After all that cleaning, we need to store the data somewhere fast, reliable, and accessible for ML models.

This is the Load phase — like stocking your kitchen pantry with perfectly prepped, labeled ingredients!

Where Can We Load Data? Storage Options:

  • Cloud Object Storage → AWS S3, Google Cloud Storage, Azure Blob, MinIO (great for large files)
  • Data Warehouse → Snowflake, BigQuery, Redshift, Databricks (great for SQL analytics)
  • Relational Database → Oracle DB, PostgreSQL, MySQL, SQL Server (great for structured data)
  • Feature Store → Feast, Hopsworks, Tecton (purpose-built for ML features)
  • Vector Database → Qdrant, Pinecone, Weaviate, Chroma, pgvector (for LLM / RAG)

📌 Note: Cloud storage examples below use AWS S3. You can swap this with Google Cloud Storage (google-cloud-storage library), Azure Blob (azure-storage-blob), or MinIO (S3-compatible, use boto3 with a custom endpoint). The concept is identical!

Step 1: Load Processed Data to Cloud Storage (Parquet Format)

We save as Parquet — not CSV. Parquet is compressed, column-oriented, and 10–50x faster to read back for ML training. Use CSV only for small files or human-readable exports.

import boto3
import pandas as pd
import io
from datetime import datetime

def load_to_cloud_storage(
    df: pd.DataFrame,
    bucket_name: str,
    prefix: str = "processed/"
) -> str:
    """
    Saves a DataFrame as a compressed Parquet file in cloud object storage.

    📌 This uses AWS S3. Alternatives:
       - Google Cloud Storage → use google-cloud-storage library
       - Azure Blob Storage   → use azure-storage-blob library
       - MinIO (self-hosted)  → use boto3 with endpoint_url parameter
    """

    # Build a date-stamped file path so each run saves separately
    timestamp = datetime.utcnow().strftime("%Y/%m/%d/%H%M%S")
    s3_key    = f"{prefix}{timestamp}/features.parquet"

    # Convert DataFrame to Parquet in memory (no temp file needed!)
    buffer = io.BytesIO()
    df.to_parquet(buffer, index=False, compression="snappy")
    buffer.seek(0)

    # Upload to S3
    s3 = boto3.client("s3")
    s3.put_object(
        Bucket=bucket_name,
        Key=s3_key,
        Body=buffer.getvalue()
    )

    full_path = f"s3://{bucket_name}/{s3_key}"
    print(f" Saved {len(df):,} rows → {full_path}")
    return full_path


output_path = load_to_cloud_storage(
    df          = enriched_df,
    bucket_name = "my-company-ml-data",
    prefix      = "features/customer_orders/"
)

Output:

 Saved 1,223 rows → s3://my-company-ml-data/features/customer_orders/2025/06/15/110532/features.parquet

Step 2: Load Clean Data Back into Oracle DB

Sometimes you also want to write the processed results back to a database for reporting or serving to applications:

import cx_Oracle
import pandas as pd

def load_to_oracle(df: pd.DataFrame, table_name: str) -> None:
    """
    Writes a DataFrame to an Oracle DB table.
    Uses executemany for batch inserts — much faster than row-by-row!

    📌 Swap cx_Oracle with:
       - psycopg2       → for PostgreSQL
       - mysql-connector → for MySQL
       - pyodbc         → for SQL Server / Azure SQL
    """

    conn   = cx_Oracle.connect("etl_user/your_password@db-host:1521/ORCLPDB")
    cursor = conn.cursor()

    # Build dynamic INSERT statement from DataFrame columns
    columns      = ", ".join(df.columns)
    placeholders = ", ".join([f":{i+1}" for i in range(len(df.columns))])
    sql          = f"INSERT INTO {table_name} ({columns}) VALUES ({placeholders})"

    # Convert DataFrame to list of tuples for batch insert
    rows = [tuple(row) for row in df.itertuples(index=False)]

    # Execute all rows in one batch call (very fast!)
    cursor.executemany(sql, rows)
    conn.commit()

    cursor.close()
    conn.close()

    print(f" Inserted {len(rows):,} rows into Oracle table '{table_name}'")


load_to_oracle(df=enriched_df, table_name="ML_FEATURES_ORDERS")

Output:

 Inserted 1,223 rows into Oracle table 'ML_FEATURES_ORDERS'

Step 3: Load Embeddings into a Vector Database (for RAG)

If you're building a RAG chatbot, you also need to store text embeddings in a vector database:

📌 Note: This example uses Qdrant (open-source). Alternatives include Pinecone (managed cloud), Weaviate, Chroma (lightweight, great for local dev), and pgvector (if you want to stay inside PostgreSQL or Oracle).

from qdrant_client import QdrantClient
from qdrant_client.models import Distance, VectorParams, PointStruct
from sentence_transformers import SentenceTransformer
from typing import List

def load_to_vector_db(texts: List[str], collection_name: str = "knowledge_base") -> None:
    """
    Converts text chunks to embeddings and stores them in a vector database.
    This powers semantic search in RAG (Retrieval-Augmented Generation) chatbots.

    📌 This uses Qdrant. Swap with Pinecone, Weaviate, Chroma, or Milvus
       — they all follow the same pattern: create collection → upsert vectors → query.
    """

    # Connect to Qdrant (running locally or in the cloud)
    client = QdrantClient(host="localhost", port=6333)

    # Create collection (384 = embedding size for our chosen model)
    client.recreate_collection(
        collection_name = collection_name,
        vectors_config  = VectorParams(size=384, distance=Distance.COSINE)
    )

    # Generate embeddings (swap model for OpenAI, Cohere, Google, etc.)
    model   = SentenceTransformer("all-MiniLM-L6-v2")
    vectors = model.encode(texts, show_progress_bar=True).tolist()

    # Build points with metadata and upsert
    points = [
        PointStruct(
            id      = idx,
            vector  = vec,
            payload = {"text": text, "source": "etl_pipeline"}
        )
        for idx, (text, vec) in enumerate(zip(texts, vectors))
    ]

    client.upsert(collection_name=collection_name, points=points)
    print(f" Loaded {len(points):,} vectors into '{collection_name}'")


load_to_vector_db(texts=clean_docs, collection_name="product_knowledge")

Output:

Batches: 100%|██████████| 1/1 [00:00<00:00 2="" 8.45it="" code="" into="" loaded="" product_knowledge="" s="" vectors="">

Now your chatbot can search for the most relevant document chunks using semantic similarity! 🤖

ETL for LLM Pipelines — RAG and Fine-Tuning 🤖

LLM applications need a special kind of ETL. Let's look at two very common patterns!

Pattern 1: Document Chunking for RAG

Long documents need to be split into smaller pieces (chunks) before embedding. Here's why: if a document is 50 pages long and a user asks one specific question, you want to retrieve just the relevant paragraph — not the whole document!

from typing import List

def chunk_document(text: str, chunk_size: int = 500, chunk_overlap: int = 50) -> List[str]:
    """
    Splits a long document into overlapping chunks for RAG ingestion.

    Why overlap? So that context at the boundary between two chunks
    is not lost. The retriever can find the right chunk even if
    the key sentence falls across a boundary.
    """

    words  = text.split()
    chunks = []
    start  = 0

    while start < len(words):
        end   = min(start + chunk_size, len(words))
        chunk = " ".join(words[start:end])
        chunks.append(chunk)
        start += (chunk_size - chunk_overlap)   # slide forward with overlap

    return chunks


# Example
long_document = "Machine learning is a branch of AI. " * 200   # simulate a long doc
chunks = chunk_document(long_document, chunk_size=100, chunk_overlap=20)

print(f" Document split into {len(chunks)} chunks.")
print(f"   First chunk preview: '{chunks[0][:80]}...'")
print(f"   Chunk overlap ensures context isn't lost at boundaries!")

Output:

 Document split into 13 chunks.
   First chunk preview: 'Machine learning is a branch of AI. Machine learning is a branch of AI. Machine...'
   Chunk overlap ensures context isn't lost at boundaries!

Pattern 2: Preparing Fine-Tuning Datasets

Fine-tuning an LLM requires data in a specific format — instruction + response pairs saved as JSONL:

import json
from typing import List, Dict

def prepare_finetune_dataset(pairs: List[Dict], output_file: str = "training.jsonl") -> None:
    """
    Converts raw Q&A pairs into JSONL format for fine-tuning.
    Compatible with OpenAI fine-tuning, Hugging Face TRL, and Axolotl.
    """

    with open(output_file, "w", encoding="utf-8") as f:
        for pair in pairs:
            record = {
                "messages": [
                    {"role": "system",    "content": "You are a helpful assistant."},
                    {"role": "user",      "content": pair["question"]},
                    {"role": "assistant", "content": pair["answer"]}
                ]
            }
            f.write(json.dumps(record, ensure_ascii=False) + "\n")

    print(f" Saved {len(pairs):,} training examples → {output_file}")


# Example usage
training_pairs = [
    {
        "question": "What is machine learning?",
        "answer":   "Machine learning is a branch of AI where systems learn patterns from data."
    },
    {
        "question": "What is the capital of India?",
        "answer":   "The capital of India is New Delhi."
    }
]

prepare_finetune_dataset(training_pairs, output_file="finetune_training.jsonl")

Output:

 Saved 2 training examples → finetune_training.jsonl

Open the file and you'll see each line is a valid JSON object ready for fine-tuning! 📄

Orchestrating ETL — Run It Automatically! 🎼

Running an ETL script manually once is fine for testing. In production, you need it to run automatically every day, retry if it fails, and alert your team if something goes wrong. That's what orchestration tools do!

📌 Orchestration Tools — Pick Any One:

  • Apache Airflow → Industry standard, powerful, DAG-based (used in big enterprises)
  • Prefect → Modern, Pythonic, easier to get started with
  • Dagster → Great for data asset-centric pipelines
  • Mage → Visual editor + code, beginner-friendly
  • Kestra → YAML-based, great for DevOps teams

Let's see how to set up our ETL pipeline in both Airflow and Prefect!

Option A: Apache Airflow DAG

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

# Default settings applied to all tasks
default_args = {
    "owner":            "mlops-team",
    "email_on_failure": True,
    "email":            ["mlops-alerts@yourcompany.com"],
    "retries":          3,                           # retry 3 times on failure
    "retry_delay":      timedelta(minutes=5),        # wait 5 min between retries
}

with DAG(
    dag_id           = "etl_oracle_ml_pipeline",
    default_args     = default_args,
    description      = "Daily ETL from Oracle DB → Cloud Storage → Vector DB",
    schedule_interval = "@daily",                   # runs every day at midnight
    start_date       = datetime(2025, 1, 1),
    catchup          = False,
    tags             = ["mlops", "etl", "oracle"],
) as dag:

    # Task 1: Extract
    task_extract = PythonOperator(
        task_id         = "extract_from_oracle",
        python_callable = extract_from_oracle,
        op_kwargs       = {"last_run_time": "{{ ds }} 00:00:00"},
    )

    # Task 2: Transform
    task_transform = PythonOperator(
        task_id         = "transform_data",
        python_callable = engineer_features,
    )

    # Task 3: Load
    task_load = PythonOperator(
        task_id         = "load_to_storage",
        python_callable = load_to_cloud_storage,
        op_kwargs       = {"bucket_name": "my-company-ml-data"},
    )

    # Define the order: Extract → Transform → Load
    task_extract >> task_transform >> task_load

Airflow will now run this automatically every day! If any task fails, it retries 3 times before alerting your team. 🛡️

Option B: Prefect Flow (Easier for Beginners!)

from prefect import flow, task
from datetime import timedelta

# Decorate each function as a Prefect task
@task(name="Extract from Oracle", retries=3, retry_delay_seconds=60)
def prefect_extract(last_run_time: str):
    return extract_from_oracle(last_run_time)

@task(name="Transform Data", retries=2)
def prefect_transform(df):
    df = clean_data(df)
    df = remove_outliers(df, "order_value")
    df = engineer_features(df)
    return df

@task(name="Load to Storage", retries=2)
def prefect_load(df):
    return load_to_cloud_storage(df, bucket_name="my-company-ml-data")


# Compose all tasks into a flow
@flow(name="Oracle ETL Pipeline")
def run_etl_pipeline(last_run_time: str = "2025-06-15 00:00:00"):
    raw        = prefect_extract(last_run_time)
    clean      = prefect_transform(raw)
    output_uri = prefect_load(clean)
    print(f"🎉 Pipeline done! Data at: {output_uri}")
    return output_uri


# Run it!
run_etl_pipeline(last_run_time="2025-06-15 00:00:00")

Output:

 Extracted 1,243 new records since 2025-06-15 00:00:00
🔵 Starting clean. Input shape: (1243, 4)
   Removed 12 duplicate rows.
   Outliers removed from 'order_value': 8 rows dropped.
   New features added: day_of_week, hour_of_day, is_weekend, value_tier, order_type_encoded
 Saved 1,223 rows → s3://my-company-ml-data/features/customer_orders/...
🎉 Pipeline done! Data at: s3://my-company-ml-data/features/customer_orders/...

The whole ETL pipeline ran automatically from start to finish! 🚀

💡 Prefect vs Airflow: Airflow is the industry standard for complex enterprise pipelines but has a steeper learning curve. Prefect feels like plain Python and is much easier for beginners. Both are excellent — choose based on your team's experience!

Data Quality Checks — Never Load Bad Data! 🔬

A pipeline that runs silently and loads bad data is worse than one that fails loudly. Always validate your data before loading it!

Five Quality Rules to Always Check:

  • Row count → Did we get the expected number of records?
  • No missing values → Are critical columns fully populated?
  • No duplicates → Is each record unique where it should be?
  • Value ranges → Are numbers within expected boundaries?
  • Required columns → Do all expected columns exist?
def run_quality_checks(df: pd.DataFrame) -> bool:
    """
    Runs essential data quality checks before loading.
    If any check fails, raises an error and halts the pipeline.
    This prevents bad data from ever reaching your ML model!
    """

    print("\n📋 Running Data Quality Checks...")
    all_passed = True

    # Check 1: Enough rows?
    if len(df) < 100:
        print(f"   FAIL — Too few rows: {len(df)} (expected at least 100)")
        all_passed = False
    else:
        print(f"   PASS — Row count: {len(df):,}")

    # Check 2: Required columns exist?
    required = ["customer_id", "order_value", "created_at"]
    missing  = [c for c in required if c not in df.columns]
    if missing:
        print(f"   FAIL — Missing columns: {missing}")
        all_passed = False
    else:
        print(f"   PASS — All required columns present")

    # Check 3: No nulls in critical columns?
    for col in ["customer_id", "order_value"]:
        if col in df.columns:
            null_count = df[col].isna().sum()
            if null_count > 0:
                print(f"   FAIL — {null_count} nulls in '{col}'")
                all_passed = False
            else:
                print(f"  PASS — No nulls in '{col}'")

    # Check 4: Order values are positive?
    if "order_value" in df.columns:
        negative = (df["order_value"] < 0).sum()
        if negative > 0:
            print(f"   FAIL — {negative} negative order values found!")
            all_passed = False
        else:
            print(f"   PASS — All order values are positive")

    # Final result
    print("=" * 45)
    if not all_passed:
        raise ValueError("❌ Pipeline halted! Fix the data quality issues above.")

    print("🎉 All checks passed! Safe to load.\n")
    return True


run_quality_checks(enriched_df)

Output:

📋 Running Data Quality Checks...
  ✅ PASS — Row count: 1,223
  ✅ PASS — All required columns present
  ✅ PASS — No nulls in 'customer_id'
  ✅ PASS — No nulls in 'order_value'
  ✅ PASS — All order values are positive
=============================================
🎉 All checks passed! Safe to load.

Now we only load data we trust! ✅

Putting It All Together — Full Pipeline! 🏗️

Here is the complete end-to-end ETL pipeline in one clean script. This is what a production MLOps pipeline looks like!

📌 Stack used in this example (swap any of these for your preferred tools):

  • Database: Oracle DB — alternatives: PostgreSQL, MySQL, Snowflake, BigQuery
  • Cloud Storage: AWS S3 — alternatives: Google Cloud Storage, Azure Blob, MinIO
  • Vector DB: Qdrant — alternatives: Pinecone, Weaviate, Chroma, pgvector
  • Orchestrator: Prefect — alternatives: Airflow, Dagster, Mage, Kestra
  • Embeddings: HuggingFace Sentence Transformers — alternatives: OpenAI, Cohere, Google
"""
full_etl_pipeline.py

Complete MLOps ETL Pipeline:
Oracle DB → Clean → Quality Check → Cloud Storage + Vector DB
"""

import json
from datetime import datetime

def run_full_pipeline(config: dict) -> dict:
    """
    Runs the complete Extract → Transform → Quality Check → Load pipeline.
    Returns a summary report of what happened.
    """

    report = {"started_at": datetime.utcnow().isoformat(), "status": "running"}

    try:
        # ── STEP 1: EXTRACT ─────────────────────────────────────────────────
        print("\n" + "="*55)
        print("📤 STEP 1: EXTRACT")
        print("="*55)

        raw_df = extract_from_oracle(config["last_run_time"])
        report["extracted_rows"] = len(raw_df)

        # ── STEP 2: TRANSFORM ────────────────────────────────────────────────
        print("\n" + "="*55)
        print("⚙️  STEP 2: TRANSFORM")
        print("="*55)

        clean_df    = clean_data(raw_df)
        clean_df    = remove_outliers(clean_df, "order_value")
        enriched_df = engineer_features(clean_df)
        report["transformed_rows"] = len(enriched_df)

        # ── STEP 3: DATA QUALITY ─────────────────────────────────────────────
        print("\n" + "="*55)
        print("🔬 STEP 3: DATA QUALITY CHECKS")
        print("="*55)

        run_quality_checks(enriched_df)

        # ── STEP 4: LOAD ─────────────────────────────────────────────────────
        print("\n" + "="*55)
        print("📥 STEP 4: LOAD")
        print("="*55)

        storage_path = load_to_cloud_storage(
            df          = enriched_df,
            bucket_name = config["bucket_name"]
        )
        report["storage_path"] = storage_path

        # Also write back to Oracle DB for reporting
        load_to_oracle(df=enriched_df, table_name="ML_FEATURES_ORDERS")

        report["status"]      = "success"
        report["finished_at"] = datetime.utcnow().isoformat()

    except Exception as e:
        report["status"] = "failed"
        report["error"]  = str(e)
        print(f"\n❌ Pipeline FAILED: {e}")
        raise

    finally:
        print("\n📋 Pipeline Report:")
        print(json.dumps(report, indent=2))

    return report


# ── Run it! ──────────────────────────────────────────────────────────────────
pipeline_config = {
    "last_run_time": "2025-06-15 00:00:00",
    "bucket_name":   "my-company-ml-data"     # swap for your cloud provider's bucket
}

run_full_pipeline(pipeline_config)

Output:

=====================================================
📤 STEP 1: EXTRACT
=====================================================
✅ Extracted 1,243 new records since 2025-06-15 00:00:00

=====================================================
⚙️  STEP 2: TRANSFORM
=====================================================
🔵 Starting clean. Input shape: (1243, 4)
  ✅ Removed 12 duplicate rows.
  ✅ Outliers removed from 'order_value': 8 rows dropped.
  ✅ New features added: day_of_week, hour_of_day, is_weekend, value_tier, order_type_encoded

=====================================================
🔬 STEP 3: DATA QUALITY CHECKS
=====================================================
📋 Running Data Quality Checks...
  ✅ PASS — Row count: 1,223
  ✅ PASS — All required columns present
  ✅ PASS — No nulls in critical columns
  ✅ PASS — All order values are positive
🎉 All checks passed! Safe to load.

=====================================================
📥 STEP 4: LOAD
=====================================================
✅ Saved 1,223 rows → s3://my-company-ml-data/features/...
✅ Inserted 1,223 rows into Oracle table 'ML_FEATURES_ORDERS'

📋 Pipeline Report:
{
  "started_at": "2025-06-15T11:05:30",
  "status": "success",
  "extracted_rows": 1243,
  "transformed_rows": 1223,
  "storage_path": "s3://my-company-ml-data/features/...",
  "finished_at": "2025-06-15T11:06:12"
}

End-to-end pipeline complete in under 2 minutes! 🎉

Best Practices ✅ and Common Mistakes ❌

Always do these things:

  • ✅ Save raw data before transforming → You need a recovery point if your cleaning code has bugs
  • ✅ Use data quality checks → Run them between Transform and Load, every single time
  • ✅ Make pipelines idempotent → Running twice should produce the same result, no duplicates
  • ✅ Write small, testable functions → One function = one responsibility
  • ✅ Log row counts at every step → If 1,243 rows go in but 3 come out, you need to know immediately
  • ✅ Use Parquet format, not CSV → Much faster and smaller for ML workloads
  • ✅ Version your data → Use date-stamped paths so you can always trace what trained which model

Never do these things:

  • ❌ Hard-code passwords in your code → Use environment variables or a secrets manager
  • ❌ Transform the raw data in place → You'll never be able to recover the original
  • ❌ Skip schema validation → Your source DB might add a column tomorrow and silently break everything
  • ❌ Write one giant 1000-line script → You'll never debug it when it fails at 3am
  • ❌ Assume the source always sends clean data → It never does. Always validate.

Quick Summary 📝

What we learned today:

  • ETL Basics → Extract, Transform, Load — the backbone of every ML pipeline
  • Extract → Pull data from Oracle DB, APIs, files using incremental strategies
  • Transform → Clean, remove outliers, encode categories, engineer features
  • Load → Store in cloud storage (Parquet), Oracle DB, or vector databases
  • LLM ETL → Chunking documents and preparing JSONL fine-tuning datasets
  • Orchestration → Automate with Airflow or Prefect so it runs on a schedule
  • Data Quality → Validate before loading — never let bad data reach your model
  • Best Practices → Idempotency, versioning, small functions, never hard-code secrets

Happy learning! 🐼✨

Comments