ETL in MLOps Pipelines: Data Extraction, Transformation, and Loading for Machine Learning
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="">00:00>
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
Post a Comment