Skip to content

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

  1. 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)

  2. 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.

  3. 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, RollingWindow is always up-to-date — trading the efficiency of fewer emitted updates for maximum freshness.

  4. 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.

  5. 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.

  6. 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.

  7. 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.

  8. 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.

  9. 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.

  10. 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.

  11. 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

  • 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
Last updated · 766 distilled / 2,225 read