Real-Time Feature Serving with Sub-5ms P99 Latency

Architecture and implementation patterns for feature stores that serve ML features in real-time with consistent sub-5ms p99 latency at scale.

#feature-store#machine-learning#real-time#mlops
Cover image for the article: Real-Time Feature Serving with Sub-5ms P99 Latency

The gap between offline model training and online model serving often comes down to features. Your model trains on carefully engineered features computed in batch, but at inference time those same features need to arrive in under 5 milliseconds. This is the feature serving problem, and getting it wrong means either stale features degrading model accuracy or latency spikes killing your user experience.

After building feature stores serving 2M+ requests per second with p99 latency under 5ms, here's what actually works.

The Problem: Feature-Serving Latency at Scale

Consider a fraud detection model that needs 47 features at inference time: user transaction history, device fingerprints, merchant risk scores, velocity checks, and graph-based features. Computing these on-the-fly would take 200-800ms. Pre-computing and caching them should take <5ms. The engineering challenge is maintaining freshness while meeting latency SLAs.

Common failure modes I've seen:

  • Cold cache misses causing p99 spikes to 50-100ms during traffic bursts
  • Feature skew where online features diverge from training-time features
  • Consistency gaps where dependent features update at different times
  • Memory pressure from caching too many entities for too few features

Architecture: Dual-Write with Streaming Materialization

The architecture separates feature computation from feature serving, connected by a streaming layer that keeps the online store fresh.

Feature Store Architecture

Online Store: Redis Cluster with Feature Encoding

The online store optimizes for single-digit millisecond reads. I use Redis Cluster with a custom binary encoding that minimizes serialization overhead.

import struct
import time
from dataclasses import dataclass
from typing import Any
import redis
import numpy as np
from redis.cluster import RedisCluster

@dataclass
class FeatureVector:
    entity_id: str
    feature_group: str
    features: dict[str, Any]
    event_timestamp: float
    created_timestamp: float

class OnlineFeatureStore:
    FEATURE_SCHEMA = {
        "float32": ("f", 4),
        "float64": ("d", 8),
        "int32": ("i", 4),
        "int64": ("q", 8),
        "bool": ("?", 1),
    }

    def __init__(self, redis_nodes: list[dict], feature_schemas: dict):
        self.redis = RedisCluster(
            startup_nodes=redis_nodes,
            decode_responses=False,
            socket_timeout=0.005,
            socket_connect_timeout=0.01,
            retry_on_timeout=True,
            health_check_interval=30,
        )
        self.schemas = feature_schemas
        self._compile_schemas()

    def _compile_schemas(self):
        self._encoders = {}
        for group, schema in self.schemas.items():
            fmt = "&#x3C;"
            for feat_name, feat_type in schema.items():
                fmt += self.FEATURE_SCHEMA[feat_type][0]
            self._encoders[group] = struct.Struct(fmt)

    def get_features(self, entity_id: str, feature_groups: list[str]) -> dict:
        pipeline = self.redis.pipeline(transaction=False)
        keys = [f"feat:{group}:{entity_id}" for group in feature_groups]

        for key in keys:
            pipeline.get(key)

        results = pipeline.execute()
        features = {}

        for group, raw_bytes in zip(feature_groups, results):
            if raw_bytes is None:
                features[group] = self._get_default_features(group)
                continue
            features[group] = self._decode_features(group, raw_bytes)

        return features

    def _decode_features(self, group: str, raw_bytes: bytes) -> dict:
        encoder = self._encoders[group]
        timestamp_bytes = 8
        event_ts = struct.unpack("&#x3C;d", raw_bytes[:timestamp_bytes])[0]
        values = encoder.unpack(raw_bytes[timestamp_bytes:timestamp_bytes + encoder.size])

        feature_names = list(self.schemas[group].keys())
        return {
            "_event_timestamp": event_ts,
            **dict(zip(feature_names, values)),
        }

    def write_features(self, vector: FeatureVector, ttl_seconds: int = 86400):
        encoder = self._encoders[vector.feature_group]
        schema = self.schemas[vector.feature_group]
        values = [vector.features[name] for name in schema.keys()]

        payload = struct.pack("&#x3C;d", vector.event_timestamp) + encoder.pack(*values)
        key = f"feat:{vector.feature_group}:{vector.entity_id}"
        self.redis.setex(key, ttl_seconds, payload)

    def _get_default_features(self, group: str) -> dict:
        return {name: 0 for name in self.schemas[group].keys()}

Streaming Materialization Pipeline

Features are materialized from event streams into the online store. This ensures freshness without requiring batch recomputation.

import { Kafka, Consumer, EachMessagePayload } from 'kafkajs';
import { createClient, RedisClientType } from 'redis';

interface FeatureEvent {
  entityId: string;
  featureGroup: string;
  features: Record&#x3C;string, number | boolean>;
  eventTimestamp: number;
  sourceTable: string;
}

interface MaterializationConfig {
  featureGroup: string;
  sourceTopic: string;
  transformFn: (raw: Record&#x3C;string, unknown>) => Record&#x3C;string, number | boolean>;
  ttlSeconds: number;
  deduplicationWindowMs: number;
}

class StreamingMaterializer {
  private kafka: Kafka;
  private redis: RedisClientType;
  private configs: Map&#x3C;string, MaterializationConfig>;
  private lastSeen: Map&#x3C;string, number> = new Map();

  constructor(kafkaBrokers: string[], redisUrl: string) {
    this.kafka = new Kafka({
      brokers: kafkaBrokers,
      clientId: 'feature-materializer',
    });
    this.redis = createClient({ url: redisUrl });
    this.configs = new Map();
  }

  async start(configs: MaterializationConfig[]): Promise&#x3C;void> {
    await this.redis.connect();

    for (const config of configs) {
      this.configs.set(config.sourceTopic, config);
    }

    const consumer = this.kafka.consumer({ groupId: 'feature-store-materializer' });
    await consumer.connect();

    for (const config of configs) {
      await consumer.subscribe({ topic: config.sourceTopic, fromBeginning: false });
    }

    await consumer.run({
      eachMessage: async (payload: EachMessagePayload) => {
        await this.processMessage(payload);
      },
    });
  }

  private async processMessage({ topic, message }: EachMessagePayload): Promise&#x3C;void> {
    const config = this.configs.get(topic);
    if (!config || !message.value) return;

    const raw = JSON.parse(message.value.toString());
    const entityId = raw.entity_id || raw.entityId || raw.id;
    const eventTimestamp = raw.event_timestamp || Date.now();

    const dedupeKey = `${config.featureGroup}:${entityId}`;
    const lastProcessed = this.lastSeen.get(dedupeKey) || 0;

    if (eventTimestamp - lastProcessed &#x3C; config.deduplicationWindowMs) {
      return;
    }

    const features = config.transformFn(raw);
    const key = `feat:${config.featureGroup}:${entityId}`;

    const payload = this.encodeFeatures(features, eventTimestamp);
    await this.redis.setEx(key, config.ttlSeconds, payload);

    this.lastSeen.set(dedupeKey, eventTimestamp);
  }

  private encodeFeatures(features: Record&#x3C;string, number | boolean>, timestamp: number): string {
    return JSON.stringify({ _ts: timestamp, ...features });
  }
}

Performance Optimizations

Connection Pooling and Pipelining

Never issue single Redis commands in a feature retrieval path. Always use pipelines to batch multiple feature group lookups into a single round trip.

Feature Group Co-location

Features commonly requested together should share a Redis hash slot. Use hash tags in keys: feat:{user_123}:transaction_velocity and feat:{user_123}:device_fingerprint will route to the same shard.

Prewarming and Background Refresh

For entities with predictable access patterns (e.g., users who log in daily), schedule background refresh jobs that update the cache before the TTL expires. This eliminates cold-start latency for active entities.

Binary Encoding vs JSON

JSON encoding adds 3-5x storage overhead and 2-4x serialization cost versus compact binary encoding. For a feature vector of 50 float32 values, JSON uses ~800 bytes while binary uses 200 bytes. At 10M cached entities, that's 6GB of savings.

Benchmarks: Real-World Performance

Tested with 50 features across 3 feature groups, 10M entities cached, Redis Cluster with 6 shards:

MetricValue
p50 latency (3 feature groups)1.2ms
p95 latency2.8ms
p99 latency4.1ms
p99.9 latency7.2ms
Throughput (sustained)2.1M req/sec
Cache hit rate99.4%
Feature freshness (median)1.8 seconds
Memory per entity (50 features)412 bytes

Feature Skew Monitoring

Feature skew — where online features diverge from offline features — is the silent killer of model performance. We log a sample of online feature values and compare distributions against the training dataset weekly.

Skew detection methodCoverageDetection latency
Statistical distribution test94%4 hours
Value range monitoring100%Real-time
Training-serving comparison87%24 hours

Handling Edge Cases

Cache misses: Return default feature values (typically zeros or global means) rather than blocking on a database lookup. Log the miss for async backfill. Models trained with feature dropout are naturally resilient to occasional default values.

Feature group partial failures: If one feature group times out, serve the request with available features rather than failing entirely. Track partial-feature inference rates and alert above thresholds.

Schema evolution: Version your feature schemas. When adding features to a group, use a new encoding version prefix in the binary payload. Old readers skip unknown fields; new readers handle missing fields with defaults.

Conclusion

Sub-5ms feature serving at scale requires treating features as a data product with its own storage layer, streaming pipeline, and monitoring. The key principles are binary encoding for minimal serialization cost, streaming materialization for freshness, and pipeline-based reads for network efficiency. Start with Redis Cluster for the online store, add streaming materialization for freshness, and invest in feature skew monitoring before it costs you model accuracy silently.

Comments

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