ZGateway: Learnings from Putting a Proxy in Front of ZippyDB¶
Summary¶
Meta introduced ZGateway, a stateless proxy tier that sits between ZippyDB clients and the ZippyDB database (ZServer) fleet. ZippyDB is Meta's most widely used key-value store, backing product metadata, counters, and configuration and serving billions of operations per second across a globally distributed fleet. In the original direct-access model every client connected to every database host it needed, producing a dense many-to-many mesh of TLS connections — a typical client held tens of thousands of outbound connections, and a typical database host accepted tens of thousands of inbound ones. That mesh scaled fan-in with the client population, so every new client cohort degraded every database host, and reconnection storms (a cohort restart, a routing bug) exhausted file descriptors and drove hosts into reboot loops. ZGateway collapses that mesh into two bounded hops: clients hold a sticky pool to their regional gateway, and each ZServer sees connections only from the gateway fleet, whose size Meta controls. Beyond connection management, the shared vantage point of a proxy tier makes possible cross-client request batching/coalescing, per-tenant admission control (Discriminant Load Shedding), read caching with live invalidation, control-plane load balancing, cross-region failover, and consolidated transaction bookkeeping — each of which is easier, safer, or only possible in one managed tier rather than in a million client binaries.
Key takeaways¶
-
A proxy is a control point, not just a hop. Interposing between many callers and a shared backend buys three things: it bounds the problem (the backend sees a fleet its operators control, not the whole client population); it creates a home for shared work (pooling, retries, routing, caching, admission control solved once by the team that knows the backend); and it creates a control point — the only place to see the whole workload, attribute load, and change behavior in minutes instead of waiting out a fleet-wide client rollout. The tradeoff is an extra hop and one more tier to operate; it pays off when the client population is large, diverse, and not yours to change. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
-
Direct access made fan-in scale with adoption. A single client could touch tens of thousands of distinct shards across hundreds of thousands of database hosts, producing a many-to-many TLS mesh. Fan-in was
H_client · p— linear in the client population — so every new cohort made every database host worse. At a few thousand clients the mesh is an inefficiency; at scale it is a reliability limit. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb) -
The reduction changes scaling behavior, not just a one-time count. With ZGateway, per-host connection counts collapse ~97–98% (model-based estimate; 20 regions, 500k database hosts, 30k proxy hosts, 1M clients, 50k shards/client) and total persistent connections drop ~19x because each backend connection multiplexes many clients. The deeper win: database-host fan-in becomes
R · S_host(regions × shard density per host), independent of the client population — "an unbounded number driven by everyone else becomes a bounded number we control." (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb) -
Cross-client batching + coalescing is a capability no client library can have. A shared batcher on each gateway host groups requests headed for the same destination (keyed by use case and physical shard) into one backend RPC; a coalescer collapses simultaneous reads of the same key into a single fetch fanned out to all callers. A client-side batcher can only merge its own process's requests. Effects: fewer/larger backend requests, lower QPS and CPU, steadier load, stretched per-use-case rate-limit budgets — and a hot key can never become a stampede against one replica. This let Meta retire the fragile, CPU-hungry client-side batching libraries that lived in a million binaries. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
-
Batching ships with two safety valves against OOM. Requests are parked in an in-memory batch flushed when a linger window elapses, a size limit is crossed, or a request-count cap is hit, so added latency stays bounded. Idle eviction erases batch-map entries idle beyond a TTL on the next flush (slow-burn growth); an in-flight cap rejects new executions once flushed-batch coroutines pile up faster than they drain (acute overload). (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
-
Discriminant Load Shedding (DLS) makes isolation structural. Every request maps to a per-tenant bucket, keyed by use case and split by priority, and buckets drain round-robin. When one tenant floods the tier, its own bucket fills and its excess is shed while every other bucket keeps draining — isolation as a property of the structure, not of luck. In a controlled overload at >90% CPU across ~1,350 active tenant buckets, only 6 (the actual noisy neighbors) were shedding; the other ~1,344 executed 99.9% of requests with zero rejections, goodput held near 97–98%, and the machinery cost ~8% of CPU. A CPU concurrency controller (AIMD loop) sits in front of DLS to adjust the shared token bucket's admit rate; a memory handler guards OOM the same way. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
-
Migration onto the proxy is pure client-side configuration. Routing to ZGateway is controlled by client-side flags scoped per service and shard prefix: a percentage knob ramps eligible traffic, a region filter limits blast radius, a global kill switch gives instant rollback — no client code change, controllable in real time. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
-
A stateless tier enables control-plane load balancing over heterogeneous hosts. Because ZGateway is stateless, any request can go to any host in a regional tier. The tier mixes ~26-core to ~126-core hosts, so equal treatment of unequal hosts produces hot outliers (→ error spikes + ServiceRouter throttling). Since ServiceRouter routes by weighted consistent hashing, a control-plane balancer reads each host's recent CPU on a fixed cadence, normalizes the tier average to 1.0, and nudges each weight opposite to its load — damped, clamped, re-centered on a target median, with a change throttle that moves only the most-imbalanced hosts (shard reshuffling is costly on cache tiers, where moving a weight means moving keys). One fixed policy can't serve both a calm tier and a tier in shock, so the balancer is becoming adaptive — classifying tier state (steady drift, task churn, flat initial weights, bimodal load, hot outliers, regional skew) and applying a matching policy. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
-
Cross-region resilience rides on the same routing substrate. Historically ZGateway was strictly regional (great for latency, but a saturated regional tier queued and timed out while healthy capacity sat idle next door). Because it sits on ServiceRouter, routing can cross region boundaries in controlled ways: global routing (a routing table spanning regions so a saturated tier fails over to a healthy one), mega-regions (group geographically close regions so overflow spills nearby and keeps most latency benefit), and rings (declare exactly which regions back each other up and in what proportion). Each is enabled per tier/region behind a percentage knob. The failover signal mattered as much as the routing: a simple regional CPU average smooths over exactly the hot conditions to catch, so failover keys off a sharper measure tuned to fire before a region tips into overload. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
-
Read caching with live invalidation, on the same pipeline. ZGateway runs in two flavors sharing one pipeline: a pure proxy and a read-through cache. On a cache tier, hot reads are served from an in-process cache; on a miss, a per-key fill lock collapses a thundering herd for one key into a single backend fetch. Freshness comes from a change-data-capture stream of write and checkpoint events that invalidates or refills affected entries under an explicit bounded-staleness contract, and each host owns a slice of the keyspace by consistent hashing. (Source: sources/2026-09-03-meta-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb)
Systems, concepts, and patterns extracted¶
Systems: ZGateway (new), ZippyDB / ZServer, ServiceRouter (new — Meta's hyperscale service mesh), Thrift/fbthrift (TLS termination in the Thrift/ServiceRouter stack).
Concepts: persistent-connection-amplification, fan-out-amplification, discriminant-load-shedding (new), concepts/request-collapsing, concepts/thundering-herd, concepts/connection-pool-exhaustion, concepts/stateless-compute, concepts/noisy-neighbor, concepts/tenant-isolation, concepts/control-plane-data-plane-separation, concepts/batching-latency-tradeoff, concepts/change-data-capture.
Patterns: collapse-many-to-many-mesh-into-two-bounded-hops (new), shared-batcher-across-clients (new), patterns/central-proxy-choke-point, connection-pooling-amortizes-handshake-cost, patterns/staged-rollout.
Operational numbers¶
- ZGateway is stateless, handles >1 billion operations/second, and currently carries ~40% of all ZippyDB traffic (projected past 60%) while adding ~6% computational overhead to an average use case.
- Runs as regional tiers discovered through ServiceRouter, in two flavors sharing one pipeline (pure proxy + read-through cache), running Meta's thick C++ ZippyDB client as its engine, one internal client per use case.
- Direct-access mesh: a typical client held tens of thousands of outbound connections; a typical database host accepted tens of thousands of inbound ones. A single client could touch tens of thousands of distinct shards across hundreds of thousands of database hosts.
- Fan-in/fan-out model (round mock figures): 20 regions, 500,000 database hosts, 30,000 proxy hosts, 1,000,000 clients, 50,000 shards/client → per-host connection counts collapse ~97–98%; total persistent connections drop ~19x end-to-end.
- Fan-out expected distinct bins:
E(H,B) = H(1 − e^(−B/H)); per-bin hit probabilityp(B) = 1 − e^(−B/H). Post-proxy database-host fan-in reduces to approximatelyR · S_host(regions × shard density per host), independent of both fleets. - DLS controlled overload: >90% CPU across ~1,350 active tenant buckets → only 6 shedding, ~1,344 executed 99.9% of requests with zero rejections; goodput ~97–98%; DLS machinery cost ~8% of CPU.
- Transaction-store consolidation: rolled out behind a flag, in nine phases, up to the highest-volume regions, reaching 100% of transaction traffic with no reliability regression.
Caveats¶
- The connection-count collapse figures (~97–98% per host, ~19x total) are an explicit model-based estimate with round mock inputs; the per-pair TLS multiplier cancels in the ratio. They illustrate scaling behavior, not measured fleet numbers.
- The 2021 ZippyDB architecture post is referenced but out of scope here; this post is specifically about the proxy layer in front of ZippyDB, not ZippyDB internals.
- "What's Next" items — agent-operated heuristics (exposing control loops as a structured surface for AI agents), co-location (pushing part of the gateway down beside the ZServer host), and a multi-process gateway (fault isolation via separate connection/TLS, worker, cache, transaction processes) — are stated directions, not shipped systems.
Source¶
- Original: https://engineering.fb.com/2026/09/03/core-infra/zgateway-proxy-zippydb-meta/
- Raw markdown:
raw/meta/2026-09-03-zgateway-learnings-from-putting-a-proxy-in-front-of-zippydb-fbe0118d.md
Related¶
- systems/zgateway — the stateless proxy tier described here.
- systems/zippydb — the key-value store ZGateway fronts.
- systems/servicerouter — Meta's service mesh that ZGateway rides for discovery, weighted-consistent-hash routing, and cross-region failover.
- discriminant-load-shedding — per-tenant round-robin bucket shedding (DLS).
- collapse-many-to-many-mesh-into-two-bounded-hops — the fan-in/fan-out reduction.
- shared-batcher-across-clients — cross-client batching + coalescing.
- companies/meta — canonical company page.