Skip to content
Reliable Data Engineering
Practice problem hard feature-storemlpoint-in-timeonline-servingstreaming
Practise with timer, notes and rubric

Design a Feature Platform for Recommendations

Problem

A streaming service (video or music) wants to rank content per user in real time. ML engineers need features like “genres watched in the last 7 days”, “items interacted with in this session”, “item popularity in the last hour”, both for training (on history) and inference (p99 < 30 ms for feature retrieval). Design the feature platform.


Clarifying questions

QuestionAssumed answer
Users / items?200 M users, 5 M items
Inference QPS?50k ranking requests/s at peak, each needs features for 1 user + 500 candidate items
Feature freshness?Batch features daily; session features within seconds
Number of features?~500 features across many teams
Training cadence?Daily retrain on 30 days of impressions with labels (clicked/played > 30 s)

1. Estimates

Online lookups: 50k req/s × (1 user + 500 items) → item features must be cached in the ranking service (5M items × 2 KB = 10 GB, fits in memory, refreshed hourly)
User features: 50k lookups/s × ~5 KB → 250 MB/s reads → Redis/Cassandra/DynamoDB cluster
Online store size: 200M users × 5 KB = 1 TB → sharded KV store
Training data: 2B impressions/day × 30 days × ~500 features → TB-scale point-in-time joins → Spark

2. Architecture

flowchart LR
    subgraph Define["Feature definitions (code, registry)"]
        REG[(Registry: name, owner, entity,<br/>source, transform, TTL, version)]
    end
    EV[[Kafka: user events]] --> STR[Streaming features<br/>Flink / Spark SS]
    LAKE[(Lakehouse: history)] --> BAT[Batch features<br/>Spark / SQL, daily]
    STR --> OFF[(Offline store<br/>Delta, time-stamped)]
    BAT --> OFF
    STR --> ON[(Online store<br/>KV, latest value)]
    BAT -->|materialise| ON
    OFF --> PIT[Point-in-time join<br/>training set builder]
    LBL[(Labels: impressions + outcomes)] --> PIT
    PIT --> TRAIN[Training]
    ON --> SERVE[Feature serving API]
    SERVE --> RANK[Ranking service]
    RANK -->|log features used| LOGS[[Kafka: feature logs]]
    LOGS --> OFF
    REG -.-> STR
    REG -.-> BAT
    REG -.-> SERVE

3. Deep dives

3.1 Point-in-time correctness (the #1 topic)

Training example: user U saw item I at t = 2026-09-10 20:05 and played it. The features must be the values as known at 20:05, not today’s values.

flowchart LR
    subgraph timeline["Feature values for user U"]
        A["09-08 02:00<br/>genre_7d = drama"] --> B["09-09 02:00<br/>genre_7d = drama"] --> C["09-10 02:00<br/>genre_7d = comedy"] --> D["09-11 02:00<br/>genre_7d = thriller"]
    end
    E["Label event<br/>09-10 20:05"] -.->|"AS OF join picks"| C
# Spark: AS-OF join via window, or Delta/Databricks feature engineering `timestamp_lookup_key`
from pyspark.sql import Window
joined = (labels.join(features, "user_id")
          .where(features.feature_ts <= labels.event_ts)
          .withColumn("rn", row_number().over(
              Window.partitionBy("impression_id").orderBy(col("feature_ts").desc())))
          .where("rn = 1"))

Leakage sources: joining current features, features computed with data after the event (e.g. a daily batch at 02:00 on 09-11 that includes 09-10 evening), label leakage via features derived from the outcome itself.

Feature logging (log the exact features used at serving time) eliminates skew for those features: train on what you served.

3.2 Training/serving skew

Same feature, two implementations (SQL for training, Java for serving) → different values → offline AUC great, online results poor. Fixes:

3.3 Streaming session features

Session features (“last 10 items played in this session”, “time since last play”) must be fresh within seconds. Flink keyed by user_id updates a small state and writes to the online store; the same events are appended to the offline store with timestamps so training can replay them point-in-time.

3.4 Online store design

3.5 Governance and discovery

Registry with owners, descriptions, lineage (source tables), freshness SLAs, usage (which models use it). Deprecation workflow. PII flags (no raw PII as features; use derived signals).

4. Trade-offs

DecisionChoiceAlternative
Build vs buyManaged feature store on the lakehouse (Databricks/Feast/Tecton)Custom: only at very large scale/specific needs
Online storeRedis (latency) / DynamoDB (ops) / Cassandra (scale, multi-DC)Lakehouse serving (too slow for 30 ms)
Item featuresIn-process cacheRemote lookup ×500 (latency)
Training dataLogged features + PIT joins for new featuresPIT joins only (skew risk)

5. Failure modes

FailureHandling
Online store partial outageDefaults + degraded-mode flag; model trained with missing-value handling
Stale batch features (job failed)Freshness monitor per feature group; serve previous values with staleness metadata; alert owner
Feature driftDistribution monitoring (PSI) offline vs online
Backfill of a new featureCompute historically with PIT semantics; parity-check vs streaming values on overlap

6. What separates a senior answer

7. Follow-up questions

How do you add a new feature and train on 30 days of history immediately?

Backfill the feature’s offline history using its batch definition with time-travel/AS-OF semantics (values as of each day), validate parity with the streaming definition on recent days, then build the training set with PIT joins. Logged features won’t contain it for the past, so PIT backfill is the only option.

Embeddings as features: how do you manage them?

Treat embedding tables as versioned feature groups (model version in the key/name); recompute on model retrain; serve via ANN index for candidate generation and KV for ranking features; never mix versions between user and item embeddings.


Self-assessment rubric