Building a Real-Time AI Personalization Engine

Architecture for sub-100ms personalized content delivery using feature stores, contextual bandits, and streaming ML inference at scale

#personalization#real-time-ml#feature-stores#contextual-bandits
Cover image for the article: Building a Real-Time AI Personalization Engine

Real-time personalization is the difference between a product that feels generic and one that feels tailor-made. Netflix, Spotify, and Amazon generate billions in revenue from personalization engines that adapt to each user within milliseconds. Building these systems requires combining streaming data infrastructure, low-latency feature computation, and real-time model inference.

This article covers the architecture for a production personalization engine delivering personalized experiences to millions of users with sub-100ms latency.

System Requirements

RequirementTargetChallenge
Latency (P95)< 100msIncluding feature lookup + inference
Throughput50K+ requests/secPeak traffic handling
Freshness< 30 secondsUser actions reflected immediately
Cold startReasonable qualityNew users with no history
DiversityAvoid filter bubblesBalance relevance and exploration
A/B testablePer-user assignmentClean experiment isolation

Chart

Architecture Overview

The system operates across three time horizons:

Time HorizonExamplesUpdate SpeedStorage
Real-time (< 30s)Current session actions, cart stateStreamingRedis
Near-time (minutes)Recent interactions, daily patternsMicro-batchFeature store
Batch (hours/days)User embeddings, long-term preferencesScheduledData warehouse

Feature Store Design

The feature store provides pre-computed features at inference time:

from dataclasses import dataclass
from typing import Dict, Any, Optional
import redis
import json
import time

@dataclass
class UserFeatures:
    user_id: str
    embedding: list  # 64-dim user representation
    segments: list  # Behavioral segments
    preferences: dict  # Category/topic affinities
    recency: dict  # Time since last interaction per category
    frequency: dict  # Interaction counts (7d, 30d)
    lifetime_value: float
    churn_risk: float

@dataclass
class ContextFeatures:
    time_of_day: int  # Hour (0-23)
    day_of_week: int  # 0-6
    device_type: str
    location_region: str
    session_depth: int  # Pages viewed in current session
    referral_source: str

class FeatureStore:
    def __init__(self, redis_client: redis.Redis):
        self.redis = redis_client

    def get_user_features(self, user_id: str) -> Optional[UserFeatures]:
        """Retrieve user features with &#x3C; 2ms latency."""
        data = self.redis.get(f"features:user:{user_id}")
        if data:
            return UserFeatures(**json.loads(data))
        return None

    def update_real_time_features(self, user_id: str, event: dict):
        """Update features from streaming events."""
        pipe = self.redis.pipeline()

        # Update session-level features
        pipe.hincrby(f"session:{user_id}", "page_views", 1)
        pipe.hset(f"session:{user_id}", "last_action", event["action"])
        pipe.hset(f"session:{user_id}", "last_action_ts", time.time())

        # Update interaction counts
        category = event.get("category")
        if category:
            pipe.hincrby(f"interactions:{user_id}:7d", category, 1)

        pipe.execute()

    def compute_context_features(self, request: dict) -> ContextFeatures:
        """Compute request-time context features."""
        from datetime import datetime
        now = datetime.utcnow()

        return ContextFeatures(
            time_of_day=now.hour,
            day_of_week=now.weekday(),
            device_type=request.get("device", "desktop"),
            location_region=request.get("region", "unknown"),
            session_depth=int(request.get("session_depth", 0)),
            referral_source=request.get("referrer", "direct"),
        )

Contextual Bandit for Exploration-Exploitation

Balance showing users what they like (exploitation) with discovering new preferences (exploration):

import numpy as np
from typing import List, Tuple

class ContextualBandit:
    """Thompson Sampling with contextual features."""

    def __init__(self, n_arms: int, feature_dim: int, alpha: float = 1.0):
        self.n_arms = n_arms
        self.feature_dim = feature_dim
        self.alpha = alpha

        # Bayesian linear regression parameters per arm
        self.B = [np.eye(feature_dim) for _ in range(n_arms)]
        self.mu = [np.zeros(feature_dim) for _ in range(n_arms)]
        self.f = [np.zeros(feature_dim) for _ in range(n_arms)]

    def select_arms(self, context: np.ndarray, n: int = 5,
                   candidates: List[int] = None) -> List[Tuple[int, float]]:
        """Select top-N arms using Thompson Sampling."""
        if candidates is None:
            candidates = list(range(self.n_arms))

        sampled_rewards = []
        for arm in candidates:
            # Sample from posterior
            B_inv = np.linalg.inv(self.B[arm])
            theta_sample = np.random.multivariate_normal(
                self.mu[arm], self.alpha * B_inv
            )
            expected_reward = context @ theta_sample
            sampled_rewards.append((arm, expected_reward))

        # Return top-N by sampled reward
        sampled_rewards.sort(key=lambda x: x[1], reverse=True)
        return sampled_rewards[:n]

    def update(self, arm: int, context: np.ndarray, reward: float):
        """Update posterior with observed reward."""
        self.B[arm] += np.outer(context, context)
        self.f[arm] += reward * context
        self.mu[arm] = np.linalg.inv(self.B[arm]) @ self.f[arm]

Real-Time Scoring Service

import torch
import torch.nn as nn
from typing import List

class PersonalizationModel(nn.Module):
    """Two-tower model for real-time personalization scoring."""

    def __init__(self, user_dim: int = 64, item_dim: int = 128,
                 context_dim: int = 32, hidden_dim: int = 128):
        super().__init__()
        # User tower
        self.user_tower = nn.Sequential(
            nn.Linear(user_dim + context_dim, hidden_dim),
            nn.ReLU(),
            nn.Linear(hidden_dim, 64),
            nn.LayerNorm(64),
        )
        # Item tower
        self.item_tower = nn.Sequential(
            nn.Linear(item_dim, hidden_dim),
            nn.ReLU(),
            nn.Linear(hidden_dim, 64),
            nn.LayerNorm(64),
        )

    def forward(self, user_features, context_features, item_features):
        user_input = torch.cat([user_features, context_features], dim=-1)
        user_repr = self.user_tower(user_input)
        item_repr = self.item_tower(item_features)
        # Dot product scoring
        scores = (user_repr * item_repr).sum(dim=-1)
        return scores


class ScoringService:
    def __init__(self, model: PersonalizationModel, feature_store: FeatureStore):
        self.model = model.eval()
        self.features = feature_store

    @torch.inference_mode()
    def score(self, user_id: str, candidate_items: List[dict],
             request_context: dict) -> List[dict]:
        """Score candidate items for a user in real-time."""
        # Get features
        user_feats = self.features.get_user_features(user_id)
        context_feats = self.features.compute_context_features(request_context)

        # Convert to tensors
        user_tensor = torch.tensor(user_feats.embedding, dtype=torch.float32)
        context_tensor = torch.tensor([
            context_feats.time_of_day / 24.0,
            context_feats.day_of_week / 7.0,
            context_feats.session_depth / 20.0,
        ] + [0.0] * 29, dtype=torch.float32)  # Pad to context_dim

        item_tensors = torch.stack([
            torch.tensor(item["embedding"], dtype=torch.float32)
            for item in candidate_items
        ])

        # Batch score
        user_batch = user_tensor.unsqueeze(0).expand(len(candidate_items), -1)
        context_batch = context_tensor.unsqueeze(0).expand(len(candidate_items), -1)

        scores = self.model(user_batch, context_batch, item_tensors)

        # Combine with business rules
        results = []
        for item, score in zip(candidate_items, scores.tolist()):
            final_score = score * item.get("quality_score", 1.0)
            # Freshness boost
            if item.get("published_hours_ago", 999) &#x3C; 24:
                final_score *= 1.2
            results.append({**item, "personalization_score": final_score})

        return sorted(results, key=lambda x: x["personalization_score"], reverse=True)

Cold Start Handling

New users with no history require special treatment:

StrategyWhen to UseQuality
Popular itemsFirst session, no signalsBaseline
Segment-basedDevice, location, referrer knownGood
Onboarding quizHigh-value productVery Good
Explore-heavy banditAfter 3+ interactionsImproving
Full personalization10+ interactionsOptimal
class ColdStartHandler:
    def __init__(self, feature_store: FeatureStore):
        self.features = feature_store

    def get_strategy(self, user_id: str, request: dict) -> str:
        """Determine cold start strategy."""
        user_features = self.features.get_user_features(user_id)

        if user_features is None:
            return "popular"
        elif sum(user_features.frequency.values()) &#x3C; 3:
            return "segment_popular"
        elif sum(user_features.frequency.values()) &#x3C; 10:
            return "explore_heavy"
        else:
            return "full_personalization"

Performance Metrics

Production deployment serving 8M personalized page loads daily:

MetricValueTarget
Scoring latency (P50)22ms< 50ms
Scoring latency (P95)68ms< 100ms
Feature lookup latency3ms< 5ms
CTR lift vs non-personalized+34%> 20%
Revenue per user lift+18%> 10%
Engagement time lift+22%> 15%
Cold start CTR2.1% (vs 3.8% warm)> 1.5%
Model refresh latency4 hours< 6 hours
Infrastructure cost$12,400/month< $15K

A/B Testing Integration

class ExperimentRouter:
    """Route users to personalization variants."""

    def __init__(self, experiments_config: dict):
        self.experiments = experiments_config

    def get_variant(self, user_id: str, experiment_id: str) -> str:
        """Deterministic assignment based on user_id hash."""
        import hashlib
        hash_input = f"{user_id}:{experiment_id}"
        hash_value = int(hashlib.md5(hash_input.encode()).hexdigest(), 16)

        experiment = self.experiments[experiment_id]
        cumulative = 0
        for variant, weight in experiment["variants"].items():
            cumulative += weight
            if (hash_value % 100) &#x3C; cumulative:
                return variant

        return experiment["control"]

Key Takeaways

  • Real-time features are the biggest accuracy lever. A user's current session behavior predicts their immediate intent better than historical averages. Invest in streaming feature infrastructure.
  • Two-tower models enable sub-50ms scoring. Pre-compute user and item representations, then score at request time with a simple dot product.
  • Contextual bandits solve explore-exploit. Pure exploitation creates filter bubbles; pure exploration hurts user experience. Thompson Sampling balances both elegantly.
  • Cold start is a design problem, not a model problem. Segment-based defaults, progressive personalization, and exploration-heavy early strategies address it systematically.
  • Measure business impact, not model accuracy. CTR lift and revenue per user matter more than offline precision. Run A/B tests on every personalization change.

The best personalization systems are invisible to users. They simply experience a product that seems to understand what they want, exactly when they want it.

Comments

    No comments yet. Be the first to share your thoughts.