Skip to content

Full-Stack AI Engineering

End-to-end knowledge for building, deploying, and scaling AI systems — from data to production.

[Data Collection] --> [Data Processing] --> [Feature Engineering]
|
[Model Training]
|
[Evaluation & Validation]
|
[Model Serving]
|
[Monitoring & Feedback]
|
[User-Facing Application (Web/Mobile/API)]

A full-stack AI engineer can work at every layer. This document covers what you need to know across all of them.

PatternUse CaseTools
Batch ETLDaily model retrainingSpark, Airflow, dbt
Stream processingReal-time featuresKafka, Flink, Kinesis
Change Data CaptureDatabase -> feature store syncDebezium, DynamoDB Streams
Lambda architectureCombined batch + real-timeBatch layer + speed layer
class DataValidator:
"""Schema and quality validation for ML datasets."""
def __init__(self):
self.checks = []
def add_check(self, name: str, fn: callable, severity: str = "error"):
self.checks.append((name, fn, severity))
def validate(self, df) -> list[dict]:
results = []
for name, fn, severity in self.checks:
passed, details = fn(df)
results.append({
"check": name,
"passed": passed,
"severity": severity,
"details": details
})
return results
# Common checks
def no_nulls(column):
def check(df):
null_count = df[column].isnull().sum()
return null_count == 0, f"{null_count} nulls in {column}"
return check
def value_range(column, min_val, max_val):
def check(df):
out_of_range = ((df[column] < min_val) | (df[column] > max_val)).sum()
return out_of_range == 0, f"{out_of_range} values out of [{min_val}, {max_val}]"
return check
def no_duplicates(columns):
def check(df):
dup_count = df.duplicated(subset=columns).sum()
return dup_count == 0, f"{dup_count} duplicate rows on {columns}"
return check
  • DVC (Data Version Control): Git-like versioning for datasets
  • Delta Lake / Apache Iceberg: Table format with time-travel and versioning
  • Importance: Reproducibility of training runs requires pinned data versions
[Raw Data] --> [Feature Pipeline] --> [Feature Store]
|
+-------------+-------------+
| |
[Online Store] [Offline Store]
(Redis/DynamoDB) (S3/BigQuery)
| |
[Real-time serving] [Training data]
TypeExampleLatencyUpdate Frequency
StaticUser age, locationN/ARarely
Batch30-day purchase countHoursDaily
Near-real-timeLast 5 items viewedSecondsOn event
Real-timeCurrent session durationMillisecondsContinuous
from dataclasses import dataclass
from typing import Any
import time
@dataclass
class Feature:
name: str
value: Any
timestamp: float
version: int
class FeatureStore:
"""Simplified online feature store."""
def __init__(self):
self.store: dict[str, dict[str, Feature]] = {} # entity_id -> {feature_name: Feature}
def set_feature(self, entity_id: str, name: str, value: Any, version: int = 1):
if entity_id not in self.store:
self.store[entity_id] = {}
self.store[entity_id][name] = Feature(name, value, time.time(), version)
def get_features(self, entity_id: str, feature_names: list[str]) -> dict[str, Any]:
"""Get multiple features for an entity. Returns dict of name -> value."""
entity = self.store.get(entity_id, {})
return {
name: entity[name].value if name in entity else None
for name in feature_names
}
def get_training_data(self, entity_ids: list[str], feature_names: list[str],
as_of: float = None) -> list[dict]:
"""Point-in-time correct feature retrieval for training."""
rows = []
for eid in entity_ids:
entity = self.store.get(eid, {})
row = {"entity_id": eid}
for name in feature_names:
if name in entity:
feat = entity[name]
if as_of is None or feat.timestamp <= as_of:
row[name] = feat.value
else:
row[name] = None
else:
row[name] = None
rows.append(row)
return rows
[Data Loader] --> [Preprocessing] --> [Model Training]
|
[Hyperparameter Search]
|
[Evaluation]
|
[Model Registry]

Track every training run with:

  • Hyperparameters
  • Dataset version
  • Metrics (loss, accuracy, custom metrics)
  • Model artifacts
  • Code version (git hash)
  • Environment (library versions)

Tools: MLflow, Weights & Biases, Neptune, Comet

PatternWhat’s DistributedUse When
Data ParallelTraining data across GPUsModel fits on one GPU
Model ParallelModel layers across GPUsModel too large for one GPU
Pipeline ParallelModel stages as pipelineVery deep models
ZeROOptimizer state, gradients, parametersMemory-efficient data parallel
StrategyDescriptionWhen to Use
Grid SearchExhaustive search over combinationsFew parameters, small search space
Random SearchRandom samplingMany parameters (often better than grid)
Bayesian OptimizationModel the objective, sample promising pointsExpensive evaluations
Population-Based TrainingEvolutionary approach across workersVery large scale
from sklearn.metrics import precision_score, recall_score, f1_score
import numpy as np
class ModelEvaluator:
def __init__(self):
self.metrics = {}
def evaluate_classification(self, y_true, y_pred, y_prob=None):
self.metrics = {
"precision": precision_score(y_true, y_pred, average="weighted"),
"recall": recall_score(y_true, y_pred, average="weighted"),
"f1": f1_score(y_true, y_pred, average="weighted"),
}
if y_prob is not None:
self.metrics["auc_roc"] = self._auc(y_true, y_prob)
return self.metrics
def evaluate_regression(self, y_true, y_pred):
self.metrics = {
"mse": np.mean((y_true - y_pred) ** 2),
"mae": np.mean(np.abs(y_true - y_pred)),
"r2": 1 - np.sum((y_true - y_pred) ** 2) / np.sum((y_true - np.mean(y_true)) ** 2),
}
return self.metrics
def _auc(self, y_true, y_prob):
# Simplified AUC calculation
from sklearn.metrics import roc_auc_score
return roc_auc_score(y_true, y_prob, multi_class="ovr", average="weighted")
  • Shadow mode: Run new model alongside production, compare outputs without serving
  • Canary deployment: Serve 1-5% of traffic with new model
  • Interleaving: Mix results from old and new model in the same response
  • Guardrail metrics: Monitor safety/quality metrics that must not degrade
PatternLatencyThroughputUse Case
Real-time (online)< 100msPer-requestAPI calls, search ranking
Near-real-time< 1sMicro-batchContent moderation
BatchMinutes-hoursVery highRecommendations, reports
StreamingContinuousHighAnomaly detection, monitoring
from abc import ABC, abstractmethod
from typing import Any
import asyncio
class ModelServer(ABC):
"""Abstract model serving interface."""
@abstractmethod
async def predict(self, input_data: Any) -> Any:
pass
@abstractmethod
async def health(self) -> dict:
pass
class BatchingModelServer(ModelServer):
"""Server with dynamic batching for GPU efficiency."""
def __init__(self, model, max_batch_size: int = 32, max_wait_ms: float = 10.0):
self.model = model
self.max_batch_size = max_batch_size
self.max_wait_ms = max_wait_ms
self.queue: asyncio.Queue = asyncio.Queue()
self._running = False
async def predict(self, input_data: Any) -> Any:
future = asyncio.get_event_loop().create_future()
await self.queue.put((input_data, future))
return await future
async def _batch_loop(self):
while self._running:
batch_inputs = []
batch_futures = []
# Wait for at least one item
input_data, future = await self.queue.get()
batch_inputs.append(input_data)
batch_futures.append(future)
# Collect more items up to batch size or timeout
deadline = asyncio.get_event_loop().time() + self.max_wait_ms / 1000
while len(batch_inputs) < self.max_batch_size:
remaining = deadline - asyncio.get_event_loop().time()
if remaining <= 0:
break
try:
input_data, future = await asyncio.wait_for(
self.queue.get(), timeout=remaining
)
batch_inputs.append(input_data)
batch_futures.append(future)
except asyncio.TimeoutError:
break
# Run batch inference
try:
results = self.model.predict_batch(batch_inputs)
for future, result in zip(batch_futures, results):
future.set_result(result)
except Exception as e:
for future in batch_futures:
future.set_exception(e)
async def start(self):
self._running = True
asyncio.create_task(self._batch_loop())
async def health(self) -> dict:
return {"status": "healthy", "queue_size": self.queue.qsize()}
TechniqueSpeed ImprovementQuality Impact
Quantization (FP16 -> INT8)2-4xMinor (<1% accuracy loss)
Knowledge Distillation5-10x (smaller model)Moderate (1-5% loss)
Pruning2-5xMinor-moderate
ONNX Runtime1.5-3xNone
TensorRT2-5xNone-minor
Batching3-10x throughputNone
PatternDescriptionExample
Streaming responseShow results as they generateChatGPT, Claude
Progressive disclosureShow confidence levelsSearch suggestions
Human-in-the-loopRequire approval for actionsGitHub Copilot suggestions
Feedback collectionThumbs up/down on outputsTraining data for improvement
ExplanationShow reasoning or sourcesRAG with citations
# Good ML API design principles
# 1. Async by default (inference can be slow)
# 2. Include request IDs for tracing
# 3. Return confidence scores
# 4. Support streaming for long-running inference
# 5. Version your models in the API
# POST /v1/predict
{
"model": "text-classifier-v2",
"input": {"text": "This movie was great!"},
"parameters": {
"temperature": 0.0,
"top_k": 5
}
}
# Response
{
"request_id": "req_abc123",
"model": "text-classifier-v2",
"output": {
"label": "positive",
"confidence": 0.94,
"alternatives": [
{"label": "neutral", "confidence": 0.04},
{"label": "negative", "confidence": 0.02}
]
},
"usage": {
"input_tokens": 12,
"processing_time_ms": 45
}
}
CategoryMetricsAlerting Threshold
Latencyp50, p95, p99 inference timep99 > 2x baseline
ThroughputRequests/sec, tokens/secDrops below expected
ErrorsError rate by type> 0.1%
Model qualityPrediction distribution shiftKL divergence > threshold
Data qualityFeature null rates, value distributionsOut of expected range
ResourcesGPU utilization, memory, CPU> 90% sustained
import numpy as np
from collections import deque
class DriftDetector:
"""Detect distribution shift in model predictions."""
def __init__(self, window_size: int = 1000, threshold: float = 0.1):
self.reference_distribution = None
self.window = deque(maxlen=window_size)
self.threshold = threshold
def set_reference(self, predictions: list[float]):
"""Set the reference distribution from validation data."""
self.reference_distribution = np.histogram(predictions, bins=50, density=True)
def observe(self, prediction: float) -> bool:
"""Observe a new prediction. Returns True if drift detected."""
self.window.append(prediction)
if len(self.window) < self.window.maxlen // 2:
return False # Not enough data
current_hist = np.histogram(list(self.window),
bins=self.reference_distribution[1],
density=True)
# KL divergence
ref = self.reference_distribution[0] + 1e-10
cur = current_hist[0] + 1e-10
kl_div = np.sum(ref * np.log(ref / cur))
return kl_div > self.threshold
ComponentLatencyNotes
Feature store lookup (online)1-10msRedis/DynamoDB
Embedding search (ANN)5-50msDepends on index size
Classification model inference5-50msDepends on model size
LLM inference (first token)100-500msDepends on input length
LLM inference (per token)10-30msAfter first token
Image model inference50-200msResNet/ViT class
Batch training jobMinutes-daysDepends on data/model size

This knowledge maps to interviews at:

  • Google: Recommendation system design, ML serving platform
  • OpenAI/Anthropic: LLM inference, training infrastructure
  • Netflix: Recommendation ML, A/B testing platform
  • Amazon: SageMaker-like platform design, product ranking
  • Apple: On-device ML, privacy-preserving ML
  • NVIDIA: GPU-optimized training/inference
  • Palantir: Enterprise ML platform, data pipeline
  • Jane Street: Signal generation, online learning