🤖 Category 04 • Intelligent Data Products & MLOps

AI & Machine Learning: Features, Pipelines & Products

Scaling artificial intelligence from raw data to production serving. Master large-scale feature engineering with PySpark MLlib, automated model evaluation, MLOps feature stores (Feast), model drift monitoring, and modern Generative AI / RAG architectures.

1. The Data Engineer's Role in Modern AI / ML

Foundational Shift

80% of real-world machine learning effort is data engineering: cleaning irregular datasets, preventing data leakage, transforming raw records into reproducible training features, and orchestrating low-latency feature serving. As Generative AI matures, data engineers are now responsible for building embeddings pipelines, vector indexing, and high-context Retrieval-Augmented Generation (RAG) platforms.

Data Leakage Prevention

Ensuring test/validation features do not leak into training sets (e.g. calculating scaling stats strictly on train splits).

Feature Stores (Feast / Hopsworks)

Unified storage providing dual interfaces: Low-latency Key-Value store (Redis) for online inference and Parquet/S3 for batch training.

Model & Data Drift

Detecting concept drift and statistical covariate shift in incoming production features using Kolmogorov-Smirnov tests.

Generative AI & RAG

Chunking documents, generating dense embeddings with OpenAI/HuggingFace, and querying vector databases with hybrid search.

Topic 4.1

Machine Learning Fundamentals & Evaluation Metrics

Statistical Foundation

Understanding core model families and rigorous evaluation metrics is vital when validating automated training pipelines.

Supervised Learning

Regression (Linear, Ridge, XGBoost for price/demand prediction) and Classification (Logistic, Random Forest for churn/fraud detection).

Unsupervised Learning

Clustering (K-Means, DBSCAN for customer segmentation) and Anomaly Detection (Isolation Forests for transaction fraud).

Classification Metrics

Precision (True Positives / Predicted Positives), Recall (True Positives / Actual Positives), F1-Score, and ROC-AUC curve area.

Regression Metrics

Mean Absolute Error (MAE), Root Mean Squared Error (RMSE • heavily penalizes large outliers), and R-squared coefficient.

Topic 4.2

Feature Engineering at Scale with PySpark

Distributed Feature Engineering

Transforming raw event streams into high-signal numerical representations that maximize model predictive power.

Imputation & Outlier Clipping

Handling missing values using median/mode strategy and clipping extreme financial tails using IQR (Interquartile Range) bounds.

Categorical Encoding

One-Hot Encoding for low cardinality (< 10 values), Target Encoding with smoothing, and StringIndexer / VectorIndexer in Spark MLlib.

Temporal & Rolling Lag Features

Creating rolling 7-day, 30-day, and 90-day transaction totals, velocity ratios, and time-since-last-transaction deltas.

Feature Selection

Variance Threshold (removing zero-variance constant features), Mutual Information, and Tree Feature Importances.

📄 pyspark_mllib_pipeline.py
from pyspark.sql import SparkSession
from pyspark.ml import Pipeline
from pyspark.ml.feature import StringIndexer, OneHotEncoder, VectorAssembler, StandardScaler
from pyspark.ml.classification import RandomForestClassifier
from pyspark.ml.evaluation import BinaryClassificationEvaluator

spark = SparkSession.builder.appName("MLFeaturePipeline").getOrCreate()

def build_churn_prediction_pipeline(df):
    """
    Builds an end-to-end distributed ML pipeline:
    Categorical indexing -> Encoding -> Vector Assembly -> Scaling -> Model Training
    """
    # 1. Index categorical strings
    indexer = StringIndexer(
        inputCols=["subscription_plan", "payment_method"],
        outputCols=["plan_index", "payment_index"],
        handleInvalid="keep"
    )

    # 2. One-hot encode indices
    encoder = OneHotEncoder(
        inputCols=["plan_index", "payment_index"],
        outputCols=["plan_vec", "payment_vec"]
    )

    # 3. Assemble all feature vectors into a single features column
    feature_cols = ["plan_vec", "payment_vec", "tenure_months", "monthly_spend", "support_tickets_30d"]
    assembler = VectorAssembler(inputCols=feature_cols, outputCol="raw_features")

    # 4. Standardise features
    scaler = StandardScaler(inputCol="raw_features", outputCol="features", withMean=True, withStd=True)

    # 5. Classifier
    rf = RandomForestClassifier(labelCol="churned", featuresCol="features", numTrees=50, maxDepth=6)

    # Chain into an atomic, reproducible pipeline
    pipeline = Pipeline(stages=[indexer, encoder, assembler, scaler, rf])

    # Train/Test Split
    train_df, test_df = df.randomSplit([0.8, 0.2], seed=42)
    model = pipeline.fit(train_df)

    # Evaluate
    predictions = model.transform(test_df)
    evaluator = BinaryClassificationEvaluator(labelCol="churned", metricName="areaUnderROC")
    roc_auc = evaluator.evaluate(predictions)
    print(f"Test ROC-AUC: {roc_auc:.4f}")

    return model
Topic 4.3

Generative AI Concepts, Vector Databases & RAG Architecture

GenAI Systems

Retrieval-Augmented Generation (RAG) grounds Large Language Models (LLMs) with private enterprise knowledge bases. Data engineers design the indexing pipeline converting documents into searchable vector embeddings.

Chunking & Tokenization

Splitting financial documents (10-K, earnings transcripts) into semantically coherent 500-token chunks with 50-token sliding overlap.

Dense Vector Embeddings

Transforming text chunks into high-dimensional vectors (e.g. 1536-dimensional OpenAI text-embedding-3) capturing semantic meaning.

Vector Stores & Indexing

Storing vectors in Pinecone, Milvus, Qdrant, or PostgreSQL pgvector using HNSW (Hierarchical Navigable Small World) for sub-10ms search.

Hybrid Search & Reranking

Combining BM25 keyword search with dense vector similarity, followed by a cross-encoder reranker (Cohere) to maximize context relevance.

2. Production Enterprise MLOps & RAG Architecture

Production Pattern
[Raw Data Lake / CRM / Market Feeds] ↓ [PySpark Distributed Feature Pipeline] • Imputation • Rolling Lag • VectorAssembler ↓ [Feast Feature Store] • S3 (Offline Training Parquet) • Redis (Online Low-Latency KV) ↓ ↓ [Batch Model Training / MLflow Registry] [Real-Time Scoring API / FastAPI] ↓ [RAG Ingestion: PDF/Filings Chunking] → [Embeddings Pipeline] → [Vector DB: pgvector]
Hands-On Lab

Lab 04: Customer Churn Feature Store & Inference Pipeline

Lab Guide
1

Feature Definition in Feast

Define an entity customer_id and feature view in Feast tracking 30-day transaction counts, payment failure rate, and session frequency.

2

Historical Feature Retrieval

Execute point-in-time correct historical joins in PySpark to generate training datasets without data leakage.

3

Model Training & MLflow Logging

Train an XGBoost classifier, log hyperparameters and ROC-AUC metrics to MLflow, and register the winning artifact.

4

Online Materialization & Serving

Materialize the latest features into Redis and deploy a FastAPI microservice returning churn risk probability in under 15ms.

3. AI / ML Engineering Interview Questions

MLOps Scenarios
Q1: What is data leakage and how do you design data pipelines to strictly prevent it? ▲

Answer: Data leakage occurs when information from outside the training dataset is inadvertently used to create the model, producing artificially inflated evaluation metrics that collapse in production.
Prevention Strategy:
1. Strict Temporal Splitting: Always split time-series or financial data chronologically (e.g. Train on Jan–Aug, Test on Sep–Dec). Never use random train/test splits.
2. Fit Transformers on Train Only: Scalers, imputers, and encoders must be fit() exclusively on the training split, then applied using transform() on test/validation data.
3. Point-in-Time Feature Stores: When generating training features, query features as-of the historical observation timestamp using as-of joins to prevent including future customer activity.

Q2: Explain the difference between Data Drift and Concept Drift in production systems. ▼

Answer:
• Data Drift (Covariate Shift): The statistical distribution of the input features P(X) changes over time, while the relationship with the target remains identical. (e.g. During inflation, average transaction amounts shift from $50 to $85, but purchasing behavior rules remain the same).
• Concept Drift: The statistical relationship between input features and the output target P(Y | X) changes. (e.g. Due to new macroeconomic policy or competitor actions, a customer with $10,000 balance who was previously safe now defaults).
We monitor data drift using Population Stability Index (PSI) or Wasserstein Distance, and concept drift using rolling production F1/Accuracy tracking.

4. Related Tracks

Explore Next

Explore Cloud Data Platforms & Lakehouse Infrastructure

Discover how AWS, Azure, S3, and Databricks provide the resilient compute backbone for modern data engineering.

☁️ Explore Cloud Track →
Open AI/ML & GenAI interview questions →