I spent some time looking at how database connection scaling degrades at massive scale, specifically the transition from direct-client routing to proxy-based architectures. When an application fleet scales to hundreds of thousands of ephemeral containers, direct client-to-storage connections grow quadratically. If every container needs to talk to every storage shard, the connection footprint on the storage nodes quickly becomes unsustainable. Each TLS connection consumes megabytes of user-space memory for buffers and session caches, starving the actual database block caches. 🔌 To solve this, Meta introduced ZGateway, an asynchronous Layer 7 proxy built on a thread-per-core execution model. Instead of cross-thread synchronization, each physical CPU core runs its own event loop using epoll or io_uring, pinning connections to specific threads to avoid cache-line bouncing. • Upstream connection pooling consolidates thousands of client connections into a small, stable pool of persistent TCP connections to each storage node. • Request multiplexing assigns unique IDs to pipelined requests, letting the proxy write to upstream sockets sequentially without waiting for responses. • Zero-copy parsing passes the actual payload directly from downstream to upstream sockets using custom buffer chains to minimize CPU overhead. It is a reminder that sometimes adding an extra network hop actually improves overall latency by freeing up database nodes from the brutal overhead of connection management. https://lnkd.in/g-ZrScaj #SystemDesign #DatabaseEngineering #DistributedSystems
Scaling Database Connections with ZGateway Proxy Architecture
More Relevant Posts
-
🔥 Most microservices don't break because of bad application code—they collapse because server-side state consistency fails when networks partition. If your backend architecture relies on naive database locks or simple sync calls across nodes, you're sitting on a ticking time bomb. Here is an in-depth, master-class breakdown on mastering Distributed Consensus, solving split-brain anomalies, and engineering bulletproof server-side architectures that scale effortlessly under heavy load. Read the full deep dive on A to Z of Software Engineering: https://lnkd.in/gPhedcYj #SoftwareEngineering #BackendEngineering #DistributedSystems #ServerSide #SystemDesign #TechLeadership #CloudComputing #Microservices
To view or add a comment, sign in
-
Our VLDB paper for 2026 is now online "Hermes at Scale: Powering Distributed Queries with a Unified Memory Fabric". New hardware technologies like CXL and UnifiedBus bring opportunities for innovation in software architecture. https://lnkd.in/edanXaNr
To view or add a comment, sign in
-
Go Observability with OpenTelemetry A practical guide to instrumenting Go applications with OpenTelemetry, covering traces, metrics, and logs to improve visibility across distributed systems. Recommended for Platform Engineers focused on building observable and reliable applications. https://lnkd.in/e2d8QdYM
To view or add a comment, sign in
-
System design is less about naming technologies and more about reasoning through scalability, consistency models, failure modes, performance bottlenecks, and reliability under load System Design Scenarios/questions 1. Real-time Routing Design routing for 1B+ daily users, with traffic changing every few seconds across millions of road segments. How would you recompute routes efficiently? 2. Video View Count A video shows “1.2M views”, but doesn’t update after every view. How would you balance accuracy, consistency, latency, and cost? 3. Global Autocomplete Design autocomplete with <100ms global latency, while handling sudden regional trending spikes. How would you keep suggestions fresh without overwhelming the system? 4. Offline Document Sync Two users edit the same document offline and reconnect together. How do you resolve conflicts without losing changes? 5. Large-Scale Search Design search across billions of emails for millions of users. How would you partition and index data while maintaining fast, isolated searches? 6. Video Conferencing Design a platform supporting 100K+ concurrent calls. How would you handle packet loss, bandwidth changes, TURN load, and quality degradation? Backend / Deep Technical 7. PostgreSQL Regression A query goes from 80ms → 14 seconds after a deployment. What do you investigate first, and how do you find the root cause? 8. Redis Cache Collapse Cache hit rate drops from 92% → 31% in 40 seconds, overwhelming the database. What happened, and how would you prevent it? 9. Distributed Transaction Failure Payment succeeds, but Inventory times out. You don’t know whether it committed. How do you recover safely without 2PC? 10. Duplicate Events A service starts producing duplicates under heavy load. How would you design reliable event IDs and deduplication? 11. Database Write Regression Adding one PostgreSQL index causes write throughput to drop 40% with no application errors. Why can this happen, and how would you prove the cause? 12. gRPC Tail Latency p50: 42ms | p99: 1.8s CPU and memory look normal. What metrics, traces, and hypotheses would you investigate first? The goal isn’t just to know the architecture. It’s being able to explain why you chose it, what trade-offs you’re making, what happens when things fail, and how the system behaves at 10× scale. #cpp #systemdesign #c #design #softwarearchitecture #engineeringdesign #coreconcept #LLD #HLD
To view or add a comment, sign in
-
Timeline architectures break under sustained load unless you split the distribution model: push writes for normal users, pull reads for high-follower accounts. A pure push model forces your database to execute millions of write operations for a single post from a high-profile account. Your write queues back up, blocking standard users. Conversely, a pure pull model forces active users to constantly rebuild their timelines on-demand, wasting massive network bandwidth and CPU cycles on stale queries. To design a stable 𝗵𝘆𝗯𝗿𝗶𝗱 𝗳𝗮𝗻𝗼𝘂𝘁 system, run this evaluation: → Segment your user base by follower count to isolate high-traffic accounts. → Route posts from standard accounts to a write-time pipeline that pre-computes feeds in a memory cache. → Route posts from high-follower accounts to a separate database store, bypassing standard pre-computation caches. → Configure the timeline retrieval service to fetch pre-computed cache data first, then dynamically pull and merge updates from the high-follower store. → Monitor write-queue latency to adjust your threshold criteria dynamically. Relying on a single-model architecture means a high-traffic surge forces a choice between two bad states. Your cache storage runs out of space from pre-computing feeds that are never read, or read latencies spike into seconds as servers struggle to rebuild timelines on every API call. 𝗣𝗿𝗲-𝗰𝗼𝗺𝗽𝘂𝘁𝗲 𝘁𝗵𝗲 𝗺𝗮𝗷𝗼𝗿𝗶𝘁𝘆, 𝗱𝘆𝗻𝗮𝗺𝗶𝗰𝗮𝗹𝗹𝘆 𝗽𝘂𝗹𝗹 𝘁𝗵𝗲 𝗲𝘅𝘁𝗿𝗲𝗺𝗲𝘀. What follower count threshold triggers a pull-on-read path in your current feed architecture? #SystemDesign #DistributedSystems #SoftwareArchitecture #Scalability
To view or add a comment, sign in
-
-
Day 5/100: Zooming out from Algorithms to Architecture Decided to touch some system design concepts. Writing efficient code is only half the battle. The other half is understanding how to keep a backend standing when traffic spikes. I covered a massive amount of ground today. To make sense of it all, I mentally mapped the concepts into 4 core pillars: Networking & Traffic: • TCP/IP Models & Network Protocols • Reverse Proxies & Load Balancing techniques Architecture & State: • Microservices vs. Monoliths (Trade-offs) • Stateless vs. Stateful systems • Sticky Sessions Databases & Scaling: • SQL vs. NoSQL • Sharding & Consistent Hashing • Eventual vs. Strong Consistency • Failover strategies Caching Mechanisms: • Cache hits, misses, and TTL • Active invalidation • The dreaded "Cache Stampede" 💡 My biggest takeaway today: Learning about Cache Stampedes. The idea that your system can completely crash just because a single highly-requested cache key expires, forcing thousands of concurrent requests to hit the database at the exact same time is terrifying, but fascinating to design around. System design feels like a completely different muscle compared to algorithmic problem-solving. #100DaysOfCode #SystemDesign #SoftwareEngineering #Backend #Architecture #NeetCode
To view or add a comment, sign in
-
Scaling Read Replicas is Not a Substitute for Proper Indexing Engineering teams often treat read replicas as a silver bullet for database latency. This approach miscalculates the trade-offs of distributed systems. Adding replicas increases architectural complexity and introduces the risk of stale data via replication lag. If your query execution plan reveals a sequential scan on a high-cardinality table, the bottleneck is algorithmic. Throwing more compute at O(n) complexity is an expensive way to mask inefficient code. Hardware upgrades provide temporary relief, but they do not solve the underlying I/O saturation. True scalability starts with optimizing the data access layer. Precise indexing and selective projection reduce the IOPS required for each transaction. This preserves headroom on your primary instance without the overhead of managing cross-region synchronization or consistency models. Build for efficiency before you build for scale. Audit your slow query logs and execution plans before you expand your infrastructure footprint. #databaseengineering #backenddevelopment #scalability #startuparchitecture #cloudcomputing
To view or add a comment, sign in
-
5,000 requests hit your API simultaneously. What breaks first? Day 13/30: Engineering Challenge At 5,000 concurrent requests, every layer can become a bottleneck. 1️⃣ Load Balancer Too many connections or insufficient capacity can cause connection queues and increased latency. Solution: Horizontal scaling, connection limits, health checks, and proper LB capacity planning. 2️⃣ Application Instances CPU can hit 100%, memory pressure can increase, and instances may become unhealthy. Solution: Use autoscaling based on meaningful metrics such as CPU, latency, or request rate. 3️⃣ Thread Pool Suppose only 200 worker threads are available: 5,000 requests ↓ 200 worker threads ↓ 4,800 waiting ↓ Latency ↑ ↓ Timeouts If threads are blocked waiting for DB/API calls, the problem gets worse. Solution: Bounded thread pools, avoid unnecessary blocking, appropriate timeouts, and virtual threads where they fit the workload. 4️⃣ DB Connection Pool Imagine: 5 API instances×20 DB connections=100 DB connections 5,000 requests cannot simultaneously query the database. The rest wait for a connection. Solution: Tune the connection pool, optimize queries, cache hot data, and use read replicas where appropriate. 5️⃣ Database 🔥 The database itself may become the real bottleneck: Solution: Optimize Query ↓ Add Index ↓ Cache ↓ Read Replica ↓ Partition / Shard if necessary Don't increase the DB connection pool blindly. More connections can actually make an overloaded database worse. 6️⃣ Downstream Services If every request calls another service: 5,000 requests ↓ Service A ↓ Service B 🔥 Service B can become saturated and cause cascading latency/failures. Solution: Timeouts + retries with backoff + circuit breakers + bulkheads + rate limiting. 🛠️ So how would I handle 5,000 simultaneous requests? Think in layers: Load Balancer / | \ API API API ↓ ↓ ↓ Thread Pools / Virtual Threads ↓ Cache ↙ ↘ DB Queue ↓ ↓ Read Replica Workers ↓ Downstream Services Use: ✅ Horizontal scaling ✅ Load balancing ✅ Caching ✅ Rate limiting ✅ Backpressure ✅ Bounded thread pools ✅ Connection pooling ✅ Async processing ✅ Queues for bursty workloads ✅ Autoscaling ✅ Timeouts & circuit breakers ✅ Database optimization 🧠 The most important lesson 5,000 requests don't necessarily require 5,000 threads, 5,000 DB connections, or 5,000 database operations at once. The goal is to control concurrency at every layer. And remember: Scaling the application doesn't help if the database remains the bottleneck. Before adding more servers, ask: “Where is the queue forming?” That's usually where your real bottleneck is. 🔥
To view or add a comment, sign in
-
-
Hedged requests are a latency-reduction technique where a client issues duplicate RPC calls when an initial request exceeds a latency threshold. By betting on faster replicas, systems can dramatically cut tail latency (p99, p999) at the cost of increased load. This post covers the architecture, trade-offs, and production patterns that make hedged requests a cornerstone of resilient distributed systems. In any distributed system, the average request latency tells only half the story. The other half lives in the tail — the p99, p999, and p9999 percentile latencies that define the worst experiences your users actually endure. A service that averages 10ms but spikes to 2 seconds at the p99 is a service that feels broken to a subset of users, no matter how fast the median is. The root cause is almost always straggler nodes: individual replicas that fall behind due to garbage collection pauses, network jitter, disk contention, or noisy neighbors on shared hardware. A single slow host among a hundred can dominate your tail latency. Research from Google's "The Tail at Scale" paper demonstrated that even modest tail latency at the server level compounds into unacceptable user-facing latency when multiple downstream RPCs are chained in a request path. Each hop multiplies the probability of hitting a straggler. Architecting Hedged Requests for Resilient Low-Latency Distributed Systems Read the full guide: https://lnkd.in/d5GDUfH3 #distributedsystems #latency #hedgedrequests #resilience #architecture
To view or add a comment, sign in
-
𝗧𝗵𝗲 𝗽𝗿𝗼𝗯𝗹𝗲𝗺 𝗶𝘀𝗻’𝘁 𝗴𝗿𝗼𝘄𝗶𝗻𝗴 𝟭𝟬×. 𝗜𝘁’𝘀 𝗴𝗿𝗼𝘄𝗶𝗻𝗴 𝟭𝟬× 𝗳𝗮𝘀𝘁𝗲𝗿 𝘁𝗵𝗮𝗻 𝘆𝗼𝘂𝗿 𝗮𝗿𝗰𝗵𝗶𝘁𝗲𝗰𝘁𝘂𝗿𝗲 𝘄𝗮𝘀 𝗱𝗲𝘀𝗶𝗴𝗻𝗲𝗱 𝗳𝗼𝗿. From 1,000 to 1,000,000 users took 6 months. From 1,000,000 to 10,000,000 took 15 days. That’s the 𝗻𝗼𝗻-𝗹𝗶𝗻𝗲𝗮𝗿 𝗴𝗿𝗼𝘄𝘁𝗵 𝘁𝗿𝗮𝗽 in software architecture. At 1K users, simple assumptions work. At 1M users, those assumptions get expensive. And when growth accelerates suddenly, yesterday’s architecture can become tomorrow’s bottleneck. Here is where systems start breaking: 𝟭. 𝗗𝗮𝘁𝗮𝗯𝗮𝘀𝗲 𝗕𝗼𝘁𝘁𝗹𝗲𝗻𝗲𝗰𝗸𝘀 Synchronous DB operations can work perfectly well at smaller scale. At higher scale, connection limits, lock contention, disk I/O, query volume, and write throughput can become bottlenecks. The fix isn’t just “𝗯𝘂𝘆 𝗮 𝗯𝗶𝗴𝗴𝗲𝗿 𝗱𝗮𝘁𝗮𝗯𝗮𝘀𝗲.” You may need a combination of: • Caching • Read replicas • Connection pooling • Async queue processing • Partitioning / sharding 𝗚𝗼𝗮𝗹: Stop every request from competing for the same database resources. 𝟮. 𝗩𝗲𝗿𝘁𝗶𝗰𝗮𝗹 𝗦𝗰𝗮𝗹𝗶𝗻𝗴 𝗖𝗲𝗶𝗹𝗶𝗻𝗴 Upgrading instances with more CPU, RAM, and bandwidth works—until a single machine becomes the constraint. That’s when horizontal scaling becomes important: 𝟭 𝗹𝗮𝗿𝗴𝗲 𝗶𝗻𝘀𝘁𝗮𝗻𝗰𝗲 ➔ 𝗔 𝗳𝗹𝗲𝗲𝘁 𝗼𝗳 𝘀𝘁𝗮𝘁𝗲𝗹𝗲𝘀𝘀, 𝗶𝗻𝗱𝗲𝗽𝗲𝗻𝗱𝗲𝗻𝘁𝗹𝘆 𝘀𝗰𝗮𝗹𝗮𝗯𝗹𝗲 𝗶𝗻𝘀𝘁𝗮𝗻𝗰𝗲𝘀 Traffic and background tasks can now be distributed across multiple nodes instead of depending on one machine. 𝟯. 𝗗𝗮𝘁𝗮 𝗗𝗶𝘀𝘁𝗿𝗶𝗯𝘂𝘁𝗶𝗼𝗻 When a single DB node becomes a bottleneck because of workload size, throughput, or contention, data may need to be distributed through 𝗽𝗮𝗿𝘁𝗶𝘁𝗶𝗼𝗻𝗶𝗻𝗴 𝗼𝗿 𝘀𝗵𝗮𝗿𝗱𝗶𝗻𝗴. For multi-tenant systems, 𝘁𝗲𝗻𝗮𝗻𝘁 𝗶𝘀𝗼𝗹𝗮𝘁𝗶𝗼𝗻 can also reduce contention and provide clearer workload boundaries. The key rule: 𝗗𝗼𝗻’𝘁 𝗷𝘂𝗺𝗽 𝘁𝗼 𝘀𝗵𝗮𝗿𝗱𝗶𝗻𝗴 𝗼𝗻 𝗗𝗮𝘆 𝟭. Introduce it when real performance and capacity metrics show that you need it. That’s the core lesson of non-linear growth. Scaling isn’t just about handling today’s traffic. It’s about predicting 𝘄𝗵𝗶𝗰𝗵 𝗮𝘀𝘀𝘂𝗺𝗽𝘁𝗶𝗼𝗻 𝘄𝗶𝗹𝗹 𝗯𝗿𝗲𝗮𝗸 𝘄𝗵𝗲𝗻 𝘁𝗼𝗺𝗼𝗿𝗿𝗼𝘄’𝘀 𝘁𝗿𝗮𝗳𝗳𝗶𝗰 𝗷𝘂𝗺𝗽𝘀 𝟭𝟬×. 𝗦𝗲𝗻𝗶𝗼𝗿 𝗮𝗿𝗰𝗵𝗶𝘁𝗲𝗰𝘁𝘂𝗿𝗲 𝗶𝘀𝗻’𝘁 𝗮𝗯𝗼𝘂𝘁 𝗼𝘃𝗲𝗿-𝗲𝗻𝗴𝗶𝗻𝗲𝗲𝗿𝗶𝗻𝗴 𝗳𝗼𝗿 𝟭𝟬𝗠 𝘂𝘀𝗲𝗿𝘀 𝗼𝗻 𝗗𝗮𝘆 𝟭. It’s about building 𝗯𝗼𝘂𝗻𝗱𝗮𝗿𝗶𝗲𝘀 𝘁𝗵𝗮𝘁 𝗹𝗲𝘁 𝘆𝗼𝘂𝗿 𝗶𝗻𝗳𝗿𝗮𝘀𝘁𝗿𝘂𝗰𝘁𝘂𝗿𝗲 𝗲𝘃𝗼𝗹𝘃𝗲 𝗮𝘀 𝘀𝗰𝗮𝗹𝗲 𝗰𝗵𝗮𝗻𝗴𝗲𝘀. When your system hit its first massive traffic surge, what broke first? 𝗖𝗼𝗻𝗻𝗲𝗰𝘁𝗶𝗼𝗻 𝗽𝗼𝗼𝗹𝘀, 𝗖𝗣𝗨 𝗹𝗶𝗺𝗶𝘁𝘀, 𝗱𝗮𝘁𝗮𝗯𝗮𝘀𝗲 𝗜/𝗢, 𝗼𝗿 𝘀𝗼𝗺𝗲𝘁𝗵𝗶𝗻𝗴 𝗲𝗹𝘀𝗲? #SystemDesign #SoftwareArchitecture #BackendEngineering #DistributedSystems #Scalability
To view or add a comment, sign in
-