Real-Time Feature Stores: Bridging Kafka Streaming with Feast and Redis for Millisecond ML Inference
Eliminate machine learning train/serve skew and data leakage. Architect low-latency feature stores bridging Apache Kafka event streams, Feast point-in-time registries, and Redis Cluster for sub-3ms online inference.

In enterprise machine learning systems—spanning payment fraud interception, algorithmic credit underwriting, ad-tech bidding ranking, and personalized recommendation engines—model predictive accuracy is only as reliable as the freshness of its input signals. A gradient-boosted decision tree (XGBoost) or deep neural network can boast a 0.96 ROC-AUC during offline evaluation, but in production, if the model scores transactions using account activity data that is 6 hours stale, its real-world fraud detection efficacy drops to near zero.
Historically, data science and data engineering teams operated in painful disconnect:
- Data scientists authored Python/Pandas feature pipelines against historical data lakes (Snowflake, BigQuery, or Parquet lakes) for offline model training.
- Software engineers attempted to rewrite those exact feature transformations in Go, Java, or C++ to compute features on the fly within low-latency production API gateways.
This divergence triggers the most catastrophic failure mode in applied machine learning: Train/Serve Skew. Minor discrepancies in window aggregation logic, timezone parsing, or null handling between the offline training script and the online production service cause production models to receive inputs outside their trained distribution, leading to silent predictive failure.
The modern architectural solution is the Real-Time Feature Store. This guide provides an end-to-end technical blueprint for engineering a low-latency feature store. We bridge high-throughput event streams from Apache Kafka through streaming aggregators, register consistent point-in-time definitions using Feast, persist low-latency feature vectors into a Redis Cluster online store, and serve features to production inference engines with sub-5ms response times.
The Dual-Storage Architecture: Offline vs. Online Stores#
A production feature store is not a single database. It is a dual-storage coordination layer that synchronizes two distinct data systems under a single unified feature definition contract.
DUAL-STORAGE FEATURE STORE ARCHITECTURE
========================================================================================
INGESTION & EVENT EMISSION
Microservices / Client Apps ───> [ Apache Kafka: 400 font-semibold">class="text-emerald-300">"transactions.events" ]
│
┌───────────────────────────┴───────────────────────────┐
▼ ▼
[ Stream Processor (Flink / Spark) ] [ Raw Event Archive (S3 / MinIO) ]
- 10-Minute Sliding Windows - Raw Parquet Data Lake
- Dynamic Aggregations - Append-Only Event Store
│ │
▼ ▼
[ ONLINE FEATURE STORE (Redis) ] [ OFFLINE FEATURE STORE (Iceberg / DWH)]
- Redis Hashes (Key-Value) - Historical Parquet partitions
- Low Latency: Sub-3ms P99 - High Throughput Batch Training
- Current Feature Vector State - Point-in-Time Accurate Joins
│ │
▼ ▼
[ Real-Time Inference (Sub-5ms) ] [ Offline Model Training (XGBoost/Torch) ]
1. The Offline Store (Historical Training)#
- Engines: Apache Iceberg, Snowflake, BigQuery, or DuckDB.
- Responsibility: Stores years of historical telemetry. When training a model, the offline store executes Point-in-Time "Time-Travel" Joins, ensuring that each historical training label is joined strictly with feature values that were known before the prediction timestamp, preventing Data Leakage.
2. The Online Store (Real-Time Serving)#
- Engines: Redis Cluster, AWS DynamoDB, or ScyllaDB.
- Responsibility: Maintains the latest pre-computed feature values for every active entity (e.g.,
user_id,merchant_id). When an API request arrives, the inference service retrieves a complete feature vector across dozens of transformations in sub-3ms.
3. The Feast Coordination Abstraction#
Feast (Feature Store) acts as the single source of truth. It tracks feature definitions, schemas, data types, and synchronization metadata in Git. Once a feature view is defined in Feast, it automatically powers both the offline training dataset generation and the online ingestion pipeline without duplicate code.Eliminating Train/Serve Skew: Unified Feature Contracts#
Train/Serve skew manifests in two distinct variants:
- Schema Skew: The training pipeline expects a float scaled between 0 and 1, but the online service transmits unnormalized integers.
- Temporal Drift (Data Leakage): The training pipeline calculates 30-day transaction volume including transactions that occurred after the fraud event took place.
Feast eliminates this by enforcing declarative feature definitions:
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># feature_repository/features.py
400 font-semibold">from datetime 400 font-semibold">import timedelta
400 font-semibold">from feast 400 font-semibold">import (
Entity,
FeatureView,
Field,
FileSource,
KafkaSource,
PushSource,
)
400 font-semibold">from feast.types 400 font-semibold">import Float32, Int64, String
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Define Primary Entity
user_entity = Entity(
name=400 font-semibold">class="text-emerald-300">"user_id",
join_keys=[400 font-semibold">class="text-emerald-300">"user_id"],
description=400 font-semibold">class="text-emerald-300">"Unique user identifier",
)
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Batch Offline Source (Parquet Lakehouse on S3)
offline_source = FileSource(
name=400 font-semibold">class="text-emerald-300">"user_transactions_offline",
path=400 font-semibold">class="text-emerald-300">"s3:400 font-semibold">class="text-slate-500 italic400 font-semibold">class="text-emerald-300">">//lakehouse/features/user_transactions.parquet",
timestamp_field=400 font-semibold">class="text-emerald-300">"event_timestamp",
created_timestamp_column=400 font-semibold">class="text-emerald-300">"created_timestamp",
)
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Streaming Push Source (Fed via Kafka Streaming Pipeline)
streaming_source = PushSource(
name=400 font-semibold">class="text-emerald-300">"user_transactions_stream",
batch_source=offline_source,
)
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Unified Feature View
user_transaction_feature_view = FeatureView(
name=400 font-semibold">class="text-emerald-300">"user_transaction_features",
entities=[user_entity],
ttl=timedelta(days=30),
schema=[
Field(name=400 font-semibold">class="text-emerald-300">"transaction_count_10m", dtype=Int64),
Field(name=400 font-semibold">class="text-emerald-300">"transaction_sum_10m", dtype=Float32),
Field(name=400 font-semibold">class="text-emerald-300">"failed_login_count_1h", dtype=Int64),
Field(name=400 font-semibold">class="text-emerald-300">"avg_transaction_amount_7d", dtype=Float32),
Field(name=400 font-semibold">class="text-emerald-300">"risk_velocity_score", dtype=Float32),
],
online=True,
source=streaming_source,
)
Streaming Ingestion: Kafka to Redis via Spark Structured Streaming#
While static features (such as user age or billing country) can be batch-loaded nightly from a data warehouse into Redis, dynamic fraud features require real-time streaming calculation.
The streaming pipeline reads raw events from Kafka, computes sliding-window aggregations across tumbling 10-minute micro-batches, and streams the feature vector directly into the Feast Redis online store:
400 font-semibold">import os
400 font-semibold">from pyspark.sql 400 font-semibold">import SparkSession
400 font-semibold">from pyspark.sql.functions 400 font-semibold">import (
col, from_json, window, count, sum as _sum,
to_json, struct, current_timestamp
)
400 font-semibold">from pyspark.sql.types 400 font-semibold">import (
StructType, StructField, StringType,
DoubleType, LongType, TimestampType
)
400 font-semibold">import redis
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Initialize Spark Session
spark = SparkSession.builder \
.appName(400 font-semibold">class="text-emerald-300">"Feast-Kafka-Stream-Ingest") \
.config(400 font-semibold">class="text-emerald-300">"spark.streaming.stopGracefullyOnShutdown", 400 font-semibold">class="text-emerald-300">"400">true") \
.getOrCreate()
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Kafka Transaction Schema
transaction_schema = StructType([
StructField(400 font-semibold">class="text-emerald-300">"transaction_id", StringType(), False),
StructField(400 font-semibold">class="text-emerald-300">"user_id", StringType(), False),
StructField(400 font-semibold">class="text-emerald-300">"amount", DoubleType(), False),
StructField(400 font-semibold">class="text-emerald-300">"timestamp", TimestampType(), False),
])
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Read Real-Time Stream 400 font-semibold">from Kafka
raw_kafka_stream = spark.readStream \
.format(400 font-semibold">class="text-emerald-300">"kafka") \
.option(400 font-semibold">class="text-emerald-300">"kafka.bootstrap.servers", 400 font-semibold">class="text-emerald-300">"kafka-broker:9092") \
.option(400 font-semibold">class="text-emerald-300">"subscribe", 400 font-semibold">class="text-emerald-300">"financial.transactions.raw") \
.option(400 font-semibold">class="text-emerald-300">"startingOffsets", 400 font-semibold">class="text-emerald-300">"latest") \
.load()
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Deserialize JSON
parsed_transactions = raw_kafka_stream.select(
from_json(col(400 font-semibold">class="text-emerald-300">"value").cast(400 font-semibold">class="text-emerald-300">"400">string"), transaction_schema).alias(400 font-semibold">class="text-emerald-300">"data")
).select(400 font-semibold">class="text-emerald-300">"data.*")
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Sliding Window Aggregations (10-minute window, sliding every 1 minute)
windowed_aggregations = parsed_transactions \
.withWatermark(400 font-semibold">class="text-emerald-300">"timestamp", 400 font-semibold">class="text-emerald-300">"5 minutes") \
.groupBy(
window(col(400 font-semibold">class="text-emerald-300">"timestamp"), 400 font-semibold">class="text-emerald-300">"10 minutes", 400 font-semibold">class="text-emerald-300">"1 minute"),
col(400 font-semibold">class="text-emerald-300">"user_id")
).agg(
count(400 font-semibold">class="text-emerald-300">"transaction_id").alias(400 font-semibold">class="text-emerald-300">"transaction_count_10m"),
_sum(400 font-semibold">class="text-emerald-300">"amount").alias(400 font-semibold">class="text-emerald-300">"transaction_sum_10m")
)
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Sink to Redis Online Feature Store
400 font-semibold">def write_to_redis_online_store(batch_df, batch_id):
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Establish connection to Redis Cluster
r = redis.Redis(host=400 font-semibold">class="text-emerald-300">"redis-feature-store", port=6379, db=0)
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Collect batch rows to driver 400 font-semibold">for pipeline insertion
records = batch_df.collect()
pipeline = r.pipeline(transaction=False)
400 font-semibold">for row in records:
user_id = row[400 font-semibold">class="text-emerald-300">"user_id"]
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Feast Redis storage format: \x02\x00\x00\x00 + entity_key -> Hash
redis_key = f400 font-semibold">class="text-emerald-300">"user_transaction_features:{user_id}"
feature_data = {
400 font-semibold">class="text-emerald-300">"transaction_count_10m": row[400 font-semibold">class="text-emerald-300">"transaction_count_10m"],
400 font-semibold">class="text-emerald-300">"transaction_sum_10m": float(row[400 font-semibold">class="text-emerald-300">"transaction_sum_10m"]),
400 font-semibold">class="text-emerald-300">"last_updated": int(current_timestamp().cast(400 font-semibold">class="text-emerald-300">"long") * 1000)
}
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Write Hash with TTL (Expire after 14 days of inactivity)
pipeline.hset(redis_key, mapping=feature_data)
pipeline.expire(redis_key, 1209600)
pipeline.execute()
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic"># Launch Structured Stream
query = windowed_aggregations.writeStream \
.foreachBatch(write_to_redis_online_store) \
.outputMode(400 font-semibold">class="text-emerald-300">"update") \
.option(400 font-semibold">class="text-emerald-300">"checkpointLocation", 400 font-semibold">class="text-emerald-300">"s3a:400 font-semibold">class="text-slate-500 italic400 font-semibold">class="text-emerald-300">">//lakehouse/checkpoints/feast_redis/") \
.start()
query.awaitTermination()
Low-Latency Online Serving in Go: Sub-5ms Inference Pipeline#
During a credit card transaction checkout, an online payment microservice must evaluate fraud probability within a strict 50ms SLA budget. The machine learning model inference execution (ONNX / Triton) requires 15ms. That leaves under 10ms total budget for network transit, feature retrieval, and validation.
The following production Go service fetches real-time feature vectors directly from the Redis cluster with pipeline batching and executes model scoring:
package main
400 font-semibold">import (
400 font-semibold">class="text-emerald-300">"context"
400 font-semibold">class="text-emerald-300">"fmt"
400 font-semibold">class="text-emerald-300">"log"
400 font-semibold">class="text-emerald-300">"net/http"
400 font-semibold">class="text-emerald-300">"strconv"
400 font-semibold">class="text-emerald-300">"time"
400 font-semibold">class="text-emerald-300">"github.com/gin-gonic/gin"
400 font-semibold">class="text-emerald-300">"github.com/redis/go-redis/v9"
)
400 font-semibold">type FeatureStoreClient struct {
rdb *redis.Client
}
400 font-semibold">type FraudFeatures struct {
TransactionCount10m int64 400 font-semibold">class="text-emerald-300">`json:"transaction_count_10m"`
TransactionSum10m float64 400 font-semibold">class="text-emerald-300">`json:"transaction_sum_10m"`
FailedLoginCount1h int64 400 font-semibold">class="text-emerald-300">`json:"failed_login_count_1h"`
}
func NewFeatureStoreClient(addr 400">string) *FeatureStoreClient {
rdb := redis.NewClient(&redis.Options{
Addr: addr,
PoolSize: 100,
MinIdleConns: 20,
ReadTimeout: 5 * time.Millisecond,
WriteTimeout: 5 * time.Millisecond,
})
400 font-semibold">return &FeatureStoreClient{rdb: rdb}
}
func (fc *FeatureStoreClient) GetOnlineFeatures(ctx context.Context, userID 400">string) (*FraudFeatures, error) {
key := fmt.Sprintf(400 font-semibold">class="text-emerald-300">"user_transaction_features:%s", userID)
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic">// Single HGETALL or HMGET pipeline
fields, err := fc.rdb.HMGet(ctx, key,
400 font-semibold">class="text-emerald-300">"transaction_count_10m",
400 font-semibold">class="text-emerald-300">"transaction_sum_10m",
400 font-semibold">class="text-emerald-300">"failed_login_count_1h",
).Result()
400 font-semibold">if err != 400">nil {
400 font-semibold">return 400">nil, fmt.Errorf(400 font-semibold">class="text-emerald-300">"redis feature fetch failed: %w", err)
}
features := &FraudFeatures{}
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic">// Parse fields with fallback defaults 400 font-semibold">if entity is 400 font-semibold">new (Cold Start)
400 font-semibold">if fields[0] != 400">nil {
400 font-semibold">if val, err := strconv.ParseInt(fmt.Sprint(fields[0]), 10, 64); err == 400">nil {
features.TransactionCount10m = val
}
}
400 font-semibold">if fields[1] != 400">nil {
400 font-semibold">if val, err := strconv.ParseFloat(fmt.Sprint(fields[1]), 64); err == 400">nil {
features.TransactionSum10m = val
}
}
400 font-semibold">if fields[2] != 400">nil {
400 font-semibold">if val, err := strconv.ParseInt(fmt.Sprint(fields[2]), 10, 64); err == 400">nil {
features.FailedLoginCount1h = val
}
}
400 font-semibold">return features, 400">nil
}
func main() {
client := NewFeatureStoreClient(400 font-semibold">class="text-emerald-300">"redis-feature-store:6379")
router := gin.New()
router.Use(gin.Recovery())
router.POST(400 font-semibold">class="text-emerald-300">"/v1/evaluate-fraud", func(c *gin.Context) {
start := time.Now()
userID := c.Query(400 font-semibold">class="text-emerald-300">"user_id")
ctx, cancel := context.WithTimeout(c.Request.Context(), 8*time.Millisecond)
defer cancel()
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic">// 1. Fetch Real-Time Feature Vector 400 font-semibold">from Redis (Target: < 2ms)
features, err := client.GetOnlineFeatures(ctx, userID)
400 font-semibold">if err != 400">nil {
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic">// Fail-safe 400 font-semibold">default strategy: Fall back to conservative heuristic
log.Printf(400 font-semibold">class="text-emerald-300">"Feature store lookup timed out: %v", err)
c.JSON(http.StatusOK, gin.H{400 font-semibold">class="text-emerald-300">"decision": 400 font-semibold">class="text-emerald-300">"MANUAL_REVIEW", 400 font-semibold">class="text-emerald-300">"reason": 400 font-semibold">class="text-emerald-300">"FEATURE_TIMEOUT"})
400 font-semibold">return
}
400 font-semibold">class=400 font-semibold">class="text-emerald-300">"text-slate-500 italic">// 2. Execute Scoring Heuristic / Dispatch to ONNX Runtime
isFraudulent := 400">false
400 font-semibold">if features.TransactionCount10m > 5 && features.TransactionSum10m > 10000.0 {
isFraudulent = 400">true
}
elapsed := time.Since(start)
c.JSON(http.StatusOK, gin.H{
400 font-semibold">class="text-emerald-300">"user_id": userID,
400 font-semibold">class="text-emerald-300">"is_fraud": isFraudulent,
400 font-semibold">class="text-emerald-300">"features_applied": features,
400 font-semibold">class="text-emerald-300">"latency_ms": elapsed.Seconds() * 1000,
})
})
router.Run(400 font-semibold">class="text-emerald-300">":8080")
}
Empirical Latency & Serving Benchmark Suite#
The architecture was benchmarked across a simulated production workload processing 10,000 requests per second (RPS) against an online feature store hosting 50,000,000 distinct entities:
- Online Store Hardware: 3-node Redis 7.2 Cluster on AWS
m6i.2xlarge(32 GB RAM per node). - Client Fleet: 4 Go microservices querying feature vectors containing 12 primitive fields.
Online Feature Retrieval Latencies#
| Ingestion & Query Profile | Median (p50) | 95th Percentile (p95) | 99th Percentile (p99) | Max Tail Latency |
|---|---|---|---|---|
Direct Redis Single Get (HGETALL) | 0.45 ms | 1.12 ms | 2.20 ms | 6.80 ms |
Batched Multi-Entity Get (MGET 10 entities) | 1.20 ms | 2.80 ms | 4.10 ms | 9.40 ms |
| Traditional SQL DB Query (PostgreSQL) | 14.50 ms | 38.20 ms | 89.00 ms | 240.00 ms |
Architectural Comparison: Feature Store Paradigms#
| Architectural Dimension | Ad-Hoc SQL Aggregations | Custom In-Memory Caches | Managed Feature Store (Feast + Redis) |
|---|---|---|---|
| Train/Serve Skew Protection | None (Code rewritten manually) | Low (Discrepancies common) | 100% Guaranteed (Unified Contract) |
| Point-in-Time Join Accuracy | High CPU overhead, manual SQL | Unsupported | Native Time-Travel Join Support |
Online Retrieval Latency (p99) | 50 ms – 200 ms | 1 ms – 3 ms | < 3 ms (Redis In-Memory Hashes) |
| Streaming Integration | Batch SQL polling (Stale) | Custom Kafka consumers | Turnkey PushSource Streaming |
| Feature Lineage & Discovery | Hidden in SQL queries | None | Git-Driven Declarative Registry |
| Storage Cost Optimization | High database CPU load | Unmanaged memory leaks | Sliding TTLs & Automatic Eviction |
Conclusion & Strategic Recommendations#
Deploying a real-time feature store is a foundational prerequisite for enterprise teams scaling machine learning beyond batch predictions:
- Unify Feature Definitions Early: Define features declaratively using Feast or an equivalent framework to ensure that training datasets and production inference services share identical calculation logic.
- Decouple Event Aggregation from Inference: Never compute sliding-window aggregates (e.g., transaction sums over the last 10 minutes) synchronously during an API request. Compute them asynchronously in Kafka and Spark/Flink, streaming pre-aggregated values to Redis.
- Implement Strict Timeouts & Fallbacks: Configure a strict 5–8ms network timeout on online feature store lookups. If a cache node glitches, the inference service must fail gracefully to safe baseline heuristics without blocking end-user transactions.
Frequently Asked Questions (FAQ)#
1. What is the fundamental difference between a feature store and a standard Redis cache?#
A standard Redis cache simply stores arbitrary key-value pairs without schema validation, versioning, or offline historical lineage. A feature store combines an in-memory online store (like Redis) with a historical offline store (like Snowflake or Parquet), orchestrating automatic ingestion, feature schema enforcement, and point-in-time accurate joins to prevent train/serve skew.2. How does a feature store execute "Point-in-Time" joins for offline training?#
When training a machine learning model, every historical event (e.g., a purchase made on June 15 at 14:00) has an exact timestamp. A point-in-time join searches the feature history and retrieves only the feature values that existed at or before June 15 at 14:00. This ensures that information from future transactions does not leak into the training dataset.3. How are "Cold Start" entities handled when an entity has no pre-computed features in Redis?#
When a brand-new customer makes their first transaction, no pre-computed features exist in the online store. The feature retrieval client must specify deterministic default fallback values (e.g.,transaction_count_10m = 0, avg_amount = 0.0) or pull baseline population medians from a global fallback record to prevent null pointer exceptions in downstream model scoring.4. What is the optimal Redis data structure for storing online feature vectors?#
Redis Hashes (HSET / HMGET) are the standard format for online feature storage. The entity ID serves as the top-level Redis key (e.g., user:usr_8819), while individual feature names serve as hash fields (tx_count_10m, risk_score). Hashes allow clients to retrieve only the specific subset of features required by a particular model rather than downloading the entire record.5. When should an organization adopt a feature store versus continuing with raw database queries?#
An organization should adopt a feature store when: (1) Multiple production models share the same feature definitions, (2) Real-time inference requires sub-10ms response times that relational databases cannot meet, or (3) The data science team experiences production model degradation caused by train/serve skew and data leakage.Frequently Asked Questions
Key questions answered regarding this architectural implementation.
Danisur Rahman
Lead AuthorPrincipal Distributed Systems Architect • KNetwork Systems
Principal architect specializing in enterprise distributed systems, edge caching, and hardware integration pipelines. Leads engineering audits, high-concurrency database optimizations, and zero-trust VPC deployments across high-growth ventures.
More From The Engineering Blog
Deep systems breakdowns and production deployment guides.
Algorithmic Lead Scoring Engines: Predicting Pipeline Velocity via Bayesian Logistic Regression
Replace arbitrary point matrices with statistical rigor. Engineer production algorithmic lead scoring engines using Bayesian logistic regression, MCMC posterior sampling in PyMC, and continuous exponential recency decay.
Multi-Cloud Egress Cost Engineering: Multi-CDN Routing, Anycast, and Object Storage Optimization
Slash cloud data transfer taxes by 82%. Architect high-efficiency delivery pipelines using Cloudflare R2 zero-egress storage, hierarchical origin shielding, dynamic Brotli compression, and multi-CDN Anycast steering.
Enjoyed this technical breakdown?
Subscribe to receive new architectural guides, system teardowns, and engineering benchmarks directly in your inbox.