How Databricks Feature Store serves features with sub-second freshness¶
Summary¶
Databricks Feature Store lets a data scientist author a feature once and have the same definition drive both large-scale offline batch flows and highly-fresh online streaming pipelines. The framework hides the infrastructure by orchestrating three finely-tuned components: Spark Real-Time Mode (RTM) for continuous stream processing, Lakebase (streaming-optimized Postgres) as online storage, and Model Serving for high-QPS retrieval. The headline result is end-to-end p99 latency of 200 ms from an event landing in Kafka to the computed feature being available in the online store — enabling fraud and personalization models to consume signals seconds old rather than the minutes-to-hours lag of scheduled Spark batch jobs.
Key takeaways¶
-
Author-once, serve-everywhere. A single declarative feature definition (e.g. "sum of a user's transaction amount over the last 10 minutes") compiles into both an offline batch pipeline for historical baselines and an online streaming pipeline for fresh signals — removing the pressure on data scientists to hand-write streaming-specific aggregation logic and stand up custom hosted infra. (Source: this page)
-
The 200 ms path has four hops. (1) Events land in Kafka (card transactions, ad impressions, clickstream). (2) A Spark RTM pipeline on serverless Lakeflow Spark Delta Pipelines continuously computes rolling aggregations. (3) Updated aggregates are written to Lakebase via a new streaming JDBC sink. (4) Model Serving retrieves the latest features from Lakebase at inference time and feeds them into the model automatically. End-to-end p99 = 200 ms.
-
Rolling windows are the natural fit for real-time serving. The Feature Store supports three window types — tumbling (wall-clock aligned, disjoint; fresh only at interval boundaries), sliding (wall-clock aligned, overlapping), and rolling (not wall-clock aligned; looks backward from each event's timestamp with millisecond resolution). Because "the sum of transactions in the last 10 minutes as of now" moves with every new event,
RollingWindowis always up-to-date — trading the efficiency of fewer emitted updates for maximum freshness. -
RTM runs stages concurrently, not sequentially like micro-batch. Traditional microbatch mode (MBM) collects events over an interval, processes them through each stage in sequence, checkpoints, then starts the next batch — putting a floor of seconds-to-minutes on stateful aggregations. RTM runs stages concurrently: aggregation operators eagerly process each row the moment it's available. For rolling aggregations there are two stages — (1) data processing (schema validation, coalescing, type casting) and (2) per-entity aggregation. Each incoming row immediately updates the aggregate in a local RocksDB state store and emits the new value downstream. Window expiration is also per-row: when a window elapses for a given event, the pipeline removes that event's contribution and emits the corrected aggregate. RocksDB runs locally on each executor, allowing state sizes that exceed cluster memory capacity.
-
RTM amortizes checkpoint cost over longer intervals without sacrificing exactly-once. MBM checkpoints at every batch boundary, and each checkpoint adds latency because it touches cloud object stores. RTM spreads planning/checkpoint cost across all rows in a longer interval instead of blocking the pipeline at each batch boundary. Exactly-once guarantees are maintained — on failure the pipeline replays at most 5 minutes of data from the Kafka source. The tradeoff: a modest increase in replay volume for a significant reduction in steady-state latency.
-
Serverless RTM on SDP eliminates cluster management, with near-zero-interruption restarts. Feature Store runs serverless RTM pipelines on Lakeflow Spark Delta Pipelines (SDP) — no machine provisioning, executor tuning, or cluster maintenance. When infra updates require a restart, SDP provisions the new serverless cluster fully-ready before stopping the old one, synchronized at the 5-minute checkpoint intervals — near-zero interruption to feature freshness.
-
Lakebase minimizes streaming write amplification vs stock Postgres. Streaming writes are a large number of small upserts — one fresh rolling-window value per Kafka row. In standard Postgres, the first modification to a page after a checkpoint writes the full 8 KB page image into the WAL (full-page write), not just the logical change; for hot entity rows updated frequently, this WAL amplification becomes the bottleneck for write throughput, replication, and recovery. Lakebase leverages compute/storage separation to write small, compact change records instead of full 8 KB snapshots; durability is protected because those compact records are acknowledged by a quorum of distributed safekeeper nodes. Full page snapshots are still generated later in the storage layer for recovery, rather than bloating the write path.
-
Scale of the serving tier. Lakebase's compute/storage separation autoscales the Online Feature Store to 10s of thousands of reads/sec at 10s-of-ms latency. Model Serving is fully horizontally scalable — inference server, auth layer, proxy, and rate limiter scale independently — sustaining 100K+ QPS on CPU endpoints, with fast elastic scaling to track traffic spikes without over-provisioning.
-
Feature lookup at inference is automatic via MLflow lineage. When a model is logged with MLflow, its feature dependencies are recorded. At inference time Model Serving automatically looks up the required features from Lakebase — no custom lookup code, no manual plumbing — and joins them with the inference request transparently.
-
Training data for stream features solved via an offline Kafka copy. Short stream retention makes training-data generation hard, so Feature Store stores an offline copy of the ingested Kafka data and computes the same feature values for historical points with point-in-time-accurate joins. The same mechanism backfills online streaming features to allow fast launch to production.
-
Governance via Unity Catalog. Features are first-class objects in Unity Catalog — discoverable, access-controlled, and tracked with full lineage; MLflow captures which features a model used, and deployment lineage connects models to feature dependencies.
Architecture / numbers¶
- End-to-end p99 latency: 200 ms (Kafka event → available in online feature store).
- Kafka replay on failure: at most 5 minutes (exactly-once bound).
- Online Feature Store read scale: 10s of thousands reads/sec at 10s-of-ms latency.
- Model Serving throughput: 100K+ QPS on CPU endpoints.
- Restart coordination point: the 5-minute checkpoint interval.
- Window types: tumbling (disjoint, boundary-fresh), sliding (overlapping), rolling (per-event, millisecond-resolution, always fresh).
Caveats¶
- The 200 ms p99 and throughput/QPS numbers are Databricks first-party, from a product blog; no independent benchmark.
- The article describes the fraud "sum over last 10 minutes" as the worked example; concrete numbers for state size, executor count, or cost are not given.
- Rolling windows trade efficiency (they emit an update per event) for freshness; tumbling/sliding remain preferred for features that don't change frequently and fit scheduled pipelines.
Source¶
- Original: https://www.databricks.com/blog/how-databricks-feature-store-serves-features-sub-second-freshness
- Raw markdown:
raw/databricks/2026-08-17-how-databricks-feature-store-serves-features-with-sub-second-ccf1ae1d.md
Related¶
- concepts/feature-store — Databricks Feature Store is a platform-native realization; this source adds the streaming/online path.
- rolling-window-aggregation — the window model central to fresh streaming features.
- systems/spark-streaming — Real-Time Mode is the compute engine.
- systems/apache-spark — RTM is a new execution mode for Structured Streaming.
- systems/rocksdb — per-executor local state store for rolling aggregates.
- systems/lakebase — streaming-optimized online feature store.
- systems/databricks-model-serving — inference-time feature retrieval.
- systems/mlflow — records feature dependencies for automatic lookup.
- systems/unity-catalog — features as governed first-class objects.
- postgres-full-page-write — the WAL-amplification failure mode Lakebase avoids for hot streaming upserts.
- author-once-serve-batch-and-streaming
- streaming-jdbc-sink-to-online-store
- backfill-online-features-from-offline-copy