A production-style, end-to-end recommendation platform for data engineering, machine learning, deployment, serving, governance, and observability workflows on Kubernetes.
This project is an end-to-end recommendation platform for e-commerce. It turns catalog, user, session, impression, behavior, and order data into batch and real-time features, trains a Behavior Sequence Transformer (BST), and serves personalized Top-K product recommendations through a production-style MLOps workflow.
-
Data and analytics platform: Generates configurable historical and real-time e-commerce events in PostgreSQL and MinIO, then streams CDC records through Debezium and Kafka. Spark builds batch features and Iceberg Bronze/Silver/Gold tables, while Flink handles event-time processing, deduplication, watermarking, streaming quality windows, and online feature updates. Airflow orchestrates ingestion, validation, compaction, materialization, drift, and analytics workflows; Feast serves PostgreSQL offline features and Redis online features; Hudi, DataHub, Trino, dbt, Superset, and Evidently provide dataset versioning, lineage, governed analytics, data quality, and drift monitoring.
-
ML training and retraining platform: Trains a PyTorch Behavior Sequence Transformer with time-aware datasets, negative sampling, ranking metrics, checkpointing, and ONNX/Triton model packaging. Kubeflow Pipelines coordinates data preparation, KubeRay/Ray Tune hyperparameter search and distributed training, evaluation, and promotion. MLflow uses PostgreSQL for tracking and registry metadata and MinIO for artifacts and versioned models; offline NDCG gates, feature-drift checks, and online candidate error/latency gates control promotion and drift-triggered retraining.
-
Serving, infrastructure, and delivery: FastAPI retrieves Feast online features, calls the Triton V2 inference API, ranks candidates, and returns personalized Top-K recommendations through NGINX. KServe manages stable and candidate Triton deployments, while KEDA HTTP/resource scalers and HPA policies autoscale API and inference workloads. Terraform and Helm provision GCP/GKE and Kubernetes resources; Jenkins validates the 15-image catalog and automates testing, immutable image publishing, deployment, shadow traffic, sticky progressive A/B rollout, model promotion, champion fallback, Helm rollback, and candidate cleanup.
-
Web UI module: Provides a React, TypeScript, Vite, and TanStack Query storefront served by a non-root NGINX container, backed by a same-origin FastAPI API. The backend uses a bounded PostgreSQL connection pool for transactional user, event, and order writes, and calls the feature and recommendation services to exercise the complete
PostgreSQL β Debezium β Kafka β Flink β Redis/Feast β Tritonreal-time path. The frontend and backend are released atomically with Helm and include ingress routing, PDBs, External Secrets, Prometheus/OpenTelemetry instrumentation, CI security checks, deployment smoke tests, and revision-based rollback. -
Security and observability: Vault and External Secrets Operator manage runtime credentials; Istio mTLS, authorization policies, and Kubernetes NetworkPolicies secure service-to-service communication. Prometheus and Pushgateway collect infrastructure, pipeline, quality, drift, API, and model-rollout metrics; Grafana provides dashboards and alerts, Loki/Promtail centralize logs, and Tempo/OpenTelemetry provide distributed tracing.
The demo shows the production web flow from user interactions and streaming feature updates through personalized recommendation serving.
- ποΈ Business Domain
- π System Overview
- π¬ Recommendation Web Demo
- ποΈ Architecture
- π Repository Main Folder Structure
- π Code Documentation Standards
- ποΈ Coursework Documentation
The following diagram presents the End-to-End Platform architecture documented in high-level system design.
The serving module retrieves fresh online features, routes stable, candidate, or shadow traffic, scores candidates with KServe/Triton, and returns Top-K recommendations.
flowchart LR
subgraph UX2["User / Client"]
direction TB
EndUser2["End User"]
Client2["Client Application"]
EndUser2 --> Client2
end
subgraph API2["API Serving"]
direction TB
Gateway2["NGINX HTTPS Gateway"]
RecAPI2["Recommendation FastAPI"]
FeatureAPI2["Online Feature FastAPI"]
Router2{"Stable / A-B / Shadow Router"}
Gateway2 --> RecAPI2
RecAPI2 --> FeatureAPI2 --> RecAPI2
RecAPI2 --> Router2
end
subgraph FS2["Online Feature Store"]
direction TB
Redis2[("Feast Online Store / Redis<br/>user, item and candidate features")]
end
FeatureAPI2 -->|Feast SDK get_online_features| Redis2
Redis2 -->|online features| FeatureAPI2
subgraph MS2["Model Serving"]
direction TB
Stable2["Stable KServe + Triton"]
Candidate2["Candidate KServe + Triton"]
Scores2["Candidate Scores"]
TopK2["Top-K Ranking"]
Response2["Recommendation Response<br/>items + model + experiment metadata"]
Stable2 --> Scores2
Candidate2 --> Scores2
Scores2 --> TopK2 --> Response2
end
Router2 -->|control| Stable2
Router2 -->|candidate| Candidate2
Router2 -.->|shadow copy| Candidate2
Client2 -->|recommendation request| Gateway2
Response2 -->|recommendations| Client2
Client2 --> EndUser2
subgraph MD2["Model Delivery"]
direction TB
ModelStore2[("Versioned Model Store")]
ModelCD2["Jenkins Model CD<br/>shadow, A/B, promote, fallback"]
ModelStore2 --> ModelCD2
end
ModelCD2 -.-> Stable2
ModelCD2 -.-> Candidate2
subgraph OBS2["Serving Observability"]
direction TB
Prometheus2["Prometheus Metrics"]
Loki2["Promtail + Loki Logs"]
OTel2["OpenTelemetry OTLP"]
Tempo2["Tempo Traces"]
Grafana2["Grafana Dashboards"]
Prometheus2 --> Grafana2
Loki2 --> Grafana2
OTel2 --> Tempo2 --> Grafana2
end
RecAPI2 -.-> Prometheus2
RecAPI2 -.-> Loki2
RecAPI2 -.-> OTel2
FeatureAPI2 -.-> Prometheus2
FeatureAPI2 -.-> Loki2
FeatureAPI2 -.-> OTel2
Stable2 -.-> Prometheus2
Candidate2 -.-> Prometheus2
classDef edge fill:#fff3e0,stroke:#ef6c00,color:#4e342e;
classDef service fill:#e3f2fd,stroke:#1976d2,color:#0d47a1;
classDef store fill:#e8f5e9,stroke:#388e3c,color:#1b5e20;
classDef model fill:#f3e5f5,stroke:#7b1fa2,color:#4a148c;
classDef result fill:#fce4ec,stroke:#c2185b,color:#880e4f;
class EndUser2,Client2,Gateway2 edge;
class RecAPI2,FeatureAPI2,ModelCD2 service;
class Redis2,ModelStore2 store;
class Router2,Stable2,Candidate2 model;
class Scores2,TopK2,Response2,Prometheus2,Loki2,OTel2,Tempo2,Grafana2 result;
The model-delivery controller turns an MLflow candidate into a shadow deployment, then progressively exposes sticky user traffic at 10% β 25% β 50%. At every step, Jenkins evaluates fresh control/candidate samples from Prometheus; a regression immediately restores champion-only traffic, while a final pass promotes the candidate and removes the temporary Triton service. See A/B testing and progressive rollout details.
flowchart TD
Candidate["Kubeflow registers versioned model<br/>MLflow candidate = test"] --> Watcher["Rollout watcher claims candidate<br/>and triggers Jenkins Model CD"]
Watcher --> Shadow["Deploy candidate KServe/Triton<br/>shadow traffic, user response from champion"]
Shadow --> ShadowGate{"Candidate Ready and<br/>shadow health gates pass?"}
ShadowGate -->|No| Rollback["Automatic rollback<br/>A/B weight = 0, shadow disabled"]
ShadowGate -->|Yes| AB10["Sticky A/B rollout<br/>candidate weight = 10%"]
AB10 --> Samples["Collect fresh control/candidate samples<br/>Prometheus + Grafana"]
AB25["Increase candidate weight to 25%"] --> Samples
AB50["Increase candidate weight to 50%"] --> Samples
Samples --> SampleGate{"At least 100 fresh samples<br/>for both variants?"}
SampleGate -->|No: HOLD| Samples
SampleGate -->|Yes| OnlineGate{"Online gates pass?<br/>error delta <= 0.02<br/>p95 latency <= 1.5x control<br/>confidence >= 0.95x control"}
OnlineGate -->|Fail| Rollback
OnlineGate -->|Pass at 10%| AB25
OnlineGate -->|Pass at 25%| AB50
OnlineGate -->|Pass at 50%| Promote["Promote candidate<br/>update stable manifest and MLflow champion"]
Rollback --> Stable["Delete candidate Triton<br/>serve previous champion only"]
Promote --> NewStable["Delete temporary candidate Triton<br/>serve new champion only"]
classDef control fill:#e3f2fd,stroke:#1976d2,color:#0d47a1;
classDef candidate fill:#fff3e0,stroke:#ef6c00,color:#4e342e;
classDef decision fill:#fff8e1,stroke:#f9a825,color:#5d4037;
classDef safe fill:#e8f5e9,stroke:#388e3c,color:#1b5e20;
classDef rollback fill:#ffebee,stroke:#c62828,color:#b71c1c;
class Candidate,Watcher control;
class Shadow,AB10,AB25,AB50,Samples candidate;
class ShadowGate,SampleGate,OnlineGate decision;
class Promote,Stable,NewStable safe;
class Rollback rollback;
The data platform combines batch and CDC ingestion, Spark and Flink processing, Airflow orchestration, data-quality checks, DataHub lineage, and Feast offline/online feature stores.
flowchart LR
subgraph Sources["Source Simulation"]
Generator["Historical + Realtime<br/>Data Generator"]
SourceDB[("Operational PostgreSQL")]
Raw[("Raw Data / MinIO")]
Generator -->|historical files| Raw
Generator -->|realtime events| SourceDB
end
subgraph Processing["Ingestion & Processing"]
Debezium["Debezium CDC"]
Kafka["Kafka"]
BronzeLake[("Iceberg Bronze Tables")]
Spark["DP2 + DP3 Spark<br/>Transform & Features"]
SilverGoldLake[("Iceberg Silver / Gold Tables")]
Flink["Flink Realtime<br/>Feature Engineering"]
Raw --> BronzeLake --> Spark --> SilverGoldLake
SourceDB --> Debezium --> Kafka --> Flink
end
subgraph Features["Feast Feature Store"]
Offline[("Offline Store<br/>PostgreSQL")]
Materialize["Feast Materialization"]
Online[("Online Store<br/>Redis")]
Offline --> Materialize --> Online
end
Spark -->|batch features| Offline
Flink -->|stream features| Offline
Flink -->|fresh online features| Online
subgraph Control["Orchestration, Quality & Governance"]
Airflow["Airflow Orchestration"]
DataChecks["Data Quality + Pipeline Health"]
DataHub["DataHub Catalog + Lineage"]
end
Airflow -.-> Spark
Airflow -.-> Flink
Airflow -.-> Materialize
Spark -.-> DataChecks
Flink -.-> DataChecks
Kafka -.-> DataChecks
BronzeLake -.-> DataHub
SilverGoldLake -.-> DataHub
Kafka -.-> DataHub
The ranking model follows the architecture in Behavior Sequence Transformer for E-commerce Recommendation in Alibaba: positional and item features represent the ordered behavior sequence, a Transformer captures dependencies between interactions, and its target-item representation is combined with user, item, context, and cross features for CTR prediction.
flowchart LR
subgraph Inputs["Input Features"]
Other["Other Features<br/>user, item, context, cross"]
SequenceItems["Behavior Sequence + Target Item<br/>item_id, category_id"]
Position["Positional Features<br/>relative interaction time"]
end
Other --> OtherEmbedding["Other Feature Embeddings"]
SequenceItems --> SequenceEmbedding["Sequence Item Embeddings"]
Position --> SequenceEmbedding
SequenceEmbedding --> Transformer["Transformer Block<br/>Multi-Head Self-Attention + FFN"]
Transformer --> TargetRepresentation["Target-Item Sequence Representation"]
OtherEmbedding --> Concatenate["Concatenate"]
TargetRepresentation --> Concatenate
Concatenate --> MLP["Three-Layer MLP"]
MLP --> Sigmoid["Sigmoid"]
Sigmoid --> CTR["Click-Through Probability / Ranking Score"]
βββ apps/ # Deployable product and data/ML workloads
β βββ analytics/ # Analytics models and dashboard bootstrap
β βββ api-serving/ # Online feature and recommendation APIs
β βββ data-platform/ # Ingestion, processing, orchestration, feature store, and governance
β βββ ml-system/ # Training, experimentation, model promotion, and serving packaging
βββ images/ # Catalog and Dockerfiles for all 15 runtime images
βββ configs/ # Data-platform and ML-system runtime configuration
βββ docs/ # Architecture, design, and coursework documentation
βββ infra/ # Production infrastructure definitions
β βββ helm/ # Kubernetes application charts
β βββ terraform/gcp/ # GCP infrastructure as code
βββ jenkins/ # CI/CD jobs, model rollout, and deployment automation
βββ ops/ # GCP, Kubernetes, and production validation operations
βββ pipelines/kubeflow/ # Generated Kubeflow artifact location
βββ notebooks/ # Tracked exploration and ML workflow notebooks
βββ tests/ # Unit, contract, integration, end-to-end, and load testsThe tables below convert the major sections from Coursework Tracking (Public).xlsx into navigable documentation indexes.
Source: tab rubic (mini-coursework).
| Rubric area | Coverage |
|---|---|
| README and high-level design | Business domain, repository structure, table of contents, and deployable-unit architecture. |
| Engineering Fundamentals | Docker, Docker Compose, multi-stage builds, and image-size optimization. |
| Implement Data Generator | Offline skew, high cardinality, schema evolution, duplicates, streaming burst/late events, configuration, and raw storage. |
| Processing Jobs | Spark offline processing, Flink streaming processing, optimization evidence, pipeline integration, and window processing. |
| Data Storage | Lakehouse compaction/partitioning and data-warehouse indexing. |
| Data Pipeline Orchestration | Airflow DP1, DP2, and DP3 ingest/validate stages. |
| Data Governance | DataHub lineage, validation, and data contracts for DP1, DP2, and DP3. |
| Schema Design | Zone schemas, SCD2 dimensions, feature timestamps, table relationships, and naming conventions. |
| Novel Ideas | Grafana-based data-quality monitoring and analytics-platform extensions. |
Source: tab rubic final-coursework (final -.
| Rubric area | Coverage |
|---|---|
| High-Level System Design | End-to-end deployment, serving, model, infrastructure, security, and delivery architecture. |
| Web API: Pull Online Features | FastAPI, Pydantic validation, async feature retrieval, health checks, Helm rollout, and fallback. |
| Web API: Model Prediction | Online features, Triton request construction, inference, ranking, and response validation. |
| Real-Time Drift Detection and ML Telemetry | Drift telemetry, scheduled comparison, dashboards, and Kubeflow retraining trigger. |
| Autoscale | KEDA/HPA autoscaling for APIs and Triton with load-test evidence. |
| Validation & Verification | Coverage, fixtures/mocks, equivalence partitions, boundary values, mutation/property-based tests, and load tests. |
| Improve the Data Generator | Configurable data drift and ID-label generation for training joins. |
| Feature Store | Incremental materialization, streaming writes to offline/online stores, and TTL design. |
| ML | Feast training-data retrieval, train/validation split, BST training, evaluation, and model saving. |
| ML Pipelines | Kubeflow pipeline stages, Ray Tune, distributed training, evaluation, and promotion. |
| Versioning | MLflow model versioning and incremental data versioning. |
| CI/CD | CI/CD for materialization, training, DP1βDP3, APIs, inference, drift detection, and streaming jobs. |
| Routing & Gateway | NGINX gateway, hidden services, authentication, rate limits, domains, and HTTPS. |
| Infrastructure as Code | Terraform-managed GCP/GKE services and infrastructure layout. |
| Observability | API and infrastructure metrics, logs, traces, Grafana dashboards, and drift monitoring. |
| A/B Testing | Stable/candidate traffic split and per-version monitoring. |
| Security | Centralized secret management, service-mesh authentication, mTLS, and authorization. |
| Repository Design | Clean repository boundaries, clean code, and design-pattern evidence. |
| Low-Level ML Design | Five key service classes and their implementation mappings. |
| Novel Ideas | Automated shadow deployment, progressive A/B gates, promotion, fallback, and cleanup. |
Source: tab rubic final-coursework (final -llm).
| Rubric area | Coverage |
|---|---|
| README and High-Level System Design | Business domain, repository structure, table of contents, and deployable-unit architecture. |
| Deploy LLM Inference Platform + Setup Custom Model | llama.cpp custom model serving, llm-d Agent Gateway routing, benchmark comparison, and load-aware optimization evidence. |
| Deploy a Global Model Config | Shared kagent ModelConfig, Agent Gateway routing, Secret reference, applied resource evidence, and end-to-end Agent inference. |
| Deploy Agent Registry | Vault-backed Agent Registry 0.4.0, persistent pgvector, namespace-scoped kagent deployment RBAC, and live UI/API proof. |
| RAG | Work in progress. |
| User and Chunk Retrieval MCP Tool + Agent | Work in progress. |
| Real-Time Drift Detection MCP Tool + Agent | Work in progress. |
| Demonstrate Basic Understanding of Agents | Work in progress. |
| Deploy a Coordinator Agent | Work in progress. |
| Agent Warm-Up | Work in progress. |
| Validation & Verification | Work in progress. |
| Improve the Data Generator | Work in progress. |
| CI/CD | Work in progress. |
| Routing & Gateway (NGINX Ingress Controller) | Work in progress. |
| IaC | Work in progress. |
| Observability | Work in progress. |
| A/B Testing | Work in progress. |
| Security | Work in progress. |
| Repository Design | Work in progress. |
| Low-Level ML Design | Work in progress. |
| Novel Ideas | Work in progress. |

