Skip to content

DATABRICKS 2026-08-12

Read original ↗

Databricks Network Configuration delivery to Tens of Millions of Serverless VMs

Summary

Databricks' serverless platform launches tens of millions of VMs per day across AWS, Azure, and GCP, and every VM must fetch its network configuration — allowed storage destinations, private-link endpoints, Unity Catalog grants, Delta Sharing destinations — before it can run a customer workload. With each node fetching config at startup and then polling for updates over its lifetime, this adds up to billions of network-config requests per day. The original design served this on the critical path of cluster creation by synchronously fanning out to every upstream service and aggregating their responses; the compound availability and 5,000 ms p99 latency of that chain became the bottleneck. Databricks re-architected delivery around an event-driven pipeline feeding a pre-computed snapshot store, so the serving path collapses to a single storage read. Results: RPC p99 latency 5,000 ms → 125 ms (97.5% reduction), service availability 99.8% → 99.99%, and 86% fewer upstream calls. (Source: sources/2026-08-12-databricks-network-configuration-delivery-to-tens-of-millions-of-serverless-vms)

Key takeaways

  • The problem is a synchronous aggregation on the launch critical path. Every cluster start triggered the network-config service to synchronously call all upstream services, aggregate, compute the per-workspace config, and return it. Each synchronous call re-did expensive, often duplicated computation across all workspaces, and load grew proportionally with tenants and their configured resources. (Source: sources/2026-08-12-databricks-network-configuration-delivery-to-tens-of-millions-of-serverless-vms)
  • Compound availability was the silent killer. With several upstream services in series on the hot path, their individual availabilities multiply, so the aggregate drops fast — translating directly into more serverless cluster-launch failures per year.
  • Pre-computation decouples the critical path — described in the post as "the single most impactful architectural decision." Moving expensive aggregation to the background turns a multi-service dependency chain into a single storage read. This is the snapshot pre-computation idea and an instance of async-projected read model / critical-path dependency minimization.
  • The system splits cleanly into a management path and a serving path — a control / data plane separation. The management path runs async: upstream services emit change events to a message queue → an event processor resolves which workspaces are affected and fans out per-workspace update notifications → a local event manager fetches the relevant upstream details, recomputes that workspace's config, and writes it to the snapshot store with a new version mark. The serving path is critical and fast: a starting cluster reads its config directly from the snapshot store with one read and no upstream calls.
  • A periodic reconciler is the safety net for missed events. A low-frequency reconciler re-syncs all workspaces in the background, giving eventual consistency even when events are dropped — "the reliability of a synchronous framework with the efficiency of push." This is reconciliation as a safety net.
  • Computation is partition-colocated with the workspaces it serves. Configs are computed and stored locally within each service partition, co-located with the workspaces they serve. This distributes compute, reduces blast radius during incidents, and eliminates cross-partition dependencies on the serving path. See partition-colocated-precomputation.
  • Events carry only identifiers, never payload. Change events carry just workspace and resource IDs. This keeps them lightweight, makes them idempotent (replayable in any order), and avoids moving sensitive customer data through the messaging pipeline.
  • Static stability falls out for free. Because the serving path reads a materialized snapshot rather than recomputing, an upstream outage doesn't break launches — clusters keep getting served the last-known-good static config.
  • Extensibility was designed in from day one. The modular, stage-based pipeline means adding a new upstream data source is a new stage implementation with zero changes to the core pipeline.

Architecture

Old (synchronous, on critical path):

cluster start ─► network-config service ─► [upstream A]
                    (sync aggregate)     ─► [upstream B]   ─► compute ─► return
                                         ─► [upstream C]
                 p99 ≈ 5,000 ms · compound availability 99.8%

New (event-driven precompute + thin serving path):

MANAGEMENT PATH (async, background)
  upstream services ─► message queue ─► event processor
     (emit change                        (resolve affected
      events: IDs only)                   workspaces, fan out
                                          per-workspace notices)
                                              │
                                              ▼
                                    per-partition event manager
                                    (fetch upstream details,
                                     recompute workspace config,
                                     write versioned snapshot)
                                              │
        periodic reconciler ──────────────────┤  (re-sync all
        (low-frequency, catches                │   workspaces)
         missed events)                        ▼
                                       [ pre-computed snapshot store ]
                                              ▲
SERVING PATH (critical, fast)                 │ single read, no upstream calls
  cluster start ─────────────────────────────┘
     p99 ≈ 125 ms · availability 99.99%

Operational numbers

Metric Before (old) After (new) Improvement
RPC latency (p99) ~5,000 ms 125 ms 97.5% reduction
Service success rate 99.8% 99.99% reduced downtime
Upstream call volume baseline −86% called only on change events
Scale tens of millions of VMs/day; billions of network-config requests/day

How events flow (worked example)

A customer creates a new Unity Catalog connection → Unity Catalog emits a change event to the message queue → the event processor determines which workspaces are attached to the affected metastore and fans out a per-workspace update notification → in that workspace's partition, the event manager fetches the updated connection details, recomputes the workspace's network config, and stores it with a new version mark. From then on, cluster requests are served directly from the snapshot store with no upstream calls.

Caveats / notes

  • The design trades strong consistency for scalability on the common path; freshness relies on events plus the reconciler backstop. The post reports improved config freshness overall (change-driven updates beat the old poll-and-recompute cadence), but the model is fundamentally eventually consistent between the event and the reconciler.
  • The post gives improvement percentages and p99/availability figures but not the reconciler cadence, partition count, or the concrete message-queue technology.
  • The legacy synchronous framework was fully deprecated after rollout.

Source

Last updated · 766 distilled / 2,225 read