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.

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.
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 = "<"
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("<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("<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<string, number | boolean>;
eventTimestamp: number;
sourceTable: string;
}
interface MaterializationConfig {
featureGroup: string;
sourceTopic: string;
transformFn: (raw: Record<string, unknown>) => Record<string, number | boolean>;
ttlSeconds: number;
deduplicationWindowMs: number;
}
class StreamingMaterializer {
private kafka: Kafka;
private redis: RedisClientType;
private configs: Map<string, MaterializationConfig>;
private lastSeen: Map<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<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<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 < 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<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:
| Metric | Value |
|---|---|
| p50 latency (3 feature groups) | 1.2ms |
| p95 latency | 2.8ms |
| p99 latency | 4.1ms |
| p99.9 latency | 7.2ms |
| Throughput (sustained) | 2.1M req/sec |
| Cache hit rate | 99.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 method | Coverage | Detection latency |
|---|---|---|
| Statistical distribution test | 94% | 4 hours |
| Value range monitoring | 100% | Real-time |
| Training-serving comparison | 87% | 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.
Recommended reading

The State of Agentic AI in 2026: Capabilities, Limitations, and Production Readiness
Comprehensive analysis of agentic AI in 2026 covering production capabilities, current limitations, and enterprise readiness benchmarks with real deployment data.

Observability for AI Agents: Tracing Multi-Step Reasoning Chains in Production
How to implement production observability for AI agents including distributed tracing, reasoning chain analysis, and debugging multi-step failures.

Measuring and Reducing AI Workload Carbon Emissions: A Practical Engineering Guide
Building a carbon-aware scheduling system for ML training and inference workloads that reduced our AI infrastructure emissions by 42% while maintaining SLA commitments.

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