0%
Data-Intensive Applications
Foundations of Data Systems
Distributed Data
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Connection Pooling and Resource Management
Every time your application opens a database connection, a surprising amount of work happens before a single query executes. Understanding this overhead explains why naive connection handling breaks at scale and why connection pooling exists.
The Connection Lifecycle
Opening a database connection involves multiple network round trips:
TCP Handshake: The client and server exchange three packets (SYN, SYN-ACK, ACK) to establish a reliable connection. On a cross-AZ network with 1ms round-trip latency, this takes roughly 1.5ms (one and a half round trips). Across availability zones in the same region, latency is typically 1-3ms. Across regions (US East to EU West), round-trip time jumps to 80-120ms, making the handshake alone cost 120-180ms. This is why cross-region database connections without pooling are catastrophically slow.
TLS Negotiation: If encryption is enabled (and it should be in production), the client and server perform a TLS handshake: exchanging certificates, agreeing on cipher suites, and deriving session keys. This adds 2-4 more round trips and CPU time for cryptographic operations. On a typical cloud setup, TLS adds 5-15ms.
Authentication: The database verifies credentials. PostgreSQL runs its authentication protocol (SCRAM-SHA-256, certificate-based, or others). MySQL performs its own handshake with challenge-response authentication. This is another round trip plus server-side computation. SCRAM-SHA-256, which is the default in modern PostgreSQL, involves multiple rounds of cryptographic computation that add measurable CPU cost when thousands of connections authenticate per minute.
Connection Setup: The database allocates memory for the connection's session state: query parser buffers, sort buffers, result set buffers, transaction context, and prepared statement caches. In PostgreSQL, this means forking a new backend process (not a thread, a full process), which involves copying page tables and initializing per-process data structures.
Query Execution: Finally, your actual SQL runs. The database parses the query, generates an execution plan, executes it against the data, and returns the result set. For simple key lookups, this might take under 1ms.
Connection Teardown: The connection is closed, TCP FIN packets exchanged, and all server-side memory freed. In PostgreSQL, the backend process terminates. In MySQL, the thread is cleaned up and its buffers deallocated.
For a query that takes 2ms to execute, the setup and teardown can take 20-50ms. You spend 90% of the time on overhead and 10% on actual work.
The math gets worse when you consider that most web requests execute multiple queries. A typical page load might run 3-5 queries: one to authenticate the user, one to load the main data, one or two for related data, and one for feature flags or configuration. If each query opens and closes its own connection, the page spends 100-250ms on connection overhead for queries that execute in 10ms total. Users perceive the application as slow, but the database is barely working: the network is doing all the waiting.
Memory Cost Per Connection
Each open connection consumes memory on the database server, and the amount is larger than most developers expect.
PostgreSQL allocates roughly 10MB per connection. Each backend process gets its own memory for work_mem (sort/hash operations), temp_buffers (temporary tables), and maintenance_work_mem. With 1,000 connections, your database server needs 10GB of RAM just for connection overhead, before any data caching.
MySQL is lighter at roughly 1MB per connection by default, but this grows with sort_buffer_size, join_buffer_size, and read_buffer_size settings. At 1,000 connections, that is still 1GB dedicated to connection state.
These numbers are per-connection overhead, not the memory used by actual queries. When a query runs a sort operation, PostgreSQL allocates work_mem (default 4MB) on top of the baseline connection memory. A complex query joining multiple tables can consume 50-100MB temporarily. Multiply that by hundreds of concurrent connections running complex queries and you see why production database servers need 128-256GB of RAM, most of which goes to connection and query overhead rather than data caching.
In system design interviews, when asked about scaling a database-backed service, always mention connection limits as a constraint. A database server with 64GB of RAM running PostgreSQL can realistically handle about 500 connections before memory pressure degrades query performance. This number surprises interviewers who assume databases handle thousands of connections natively.
Thread and File Descriptor Exhaustion
Memory is not the only resource that connections consume. Each TCP connection uses a file descriptor on both the client and server operating systems. Linux defaults to a per-process limit of 1,024 file descriptors. A server with 1,000 database connections, plus connections for HTTP clients, log files, and other I/O, can hit this limit and start failing to accept new connections with "too many open files" errors.
On the application side, many frameworks use one thread per connection. Opening 500 connections means 500 threads, each consuming stack space (typically 512KB to 1MB per thread). That is 250-500MB of memory just for thread stacks, plus the CPU cost of context switching between 500 threads. When the OS scheduler spends more time deciding which thread to run than actually running threads, you have entered the overhead spiral.
Why Open-Per-Request Breaks at Scale
Consider a web application with 100 concurrent users. Each request opens a connection, runs a query, and closes the connection. At 100 requests per second, you are creating and destroying 100 TCP connections every second. Each connection costs 20-50ms of setup overhead, consuming CPU on both the application server and the database server.
Now scale to 1,000 concurrent users. You need 1,000 simultaneous connections. On PostgreSQL, that is 10GB of RAM consumed by connection state alone. The database spends more time managing connections than executing queries. Response times spike, timeouts cascade, and the system collapses under its own overhead.
The failure mode is particularly nasty because it creates a positive feedback loop. Slow responses cause clients to retry, which creates more connections, which makes the database slower, which causes more timeouts, which causes more retries. Within minutes, a system handling 1,000 requests per second can accumulate 5,000 pending connections as retries pile up. The database becomes unresponsive, health checks fail, and the load balancer marks every instance as unhealthy. What started as a gradual performance degradation becomes a total outage.
Connection Limits in Cloud-Managed Databases
Cloud-managed databases impose hard connection limits that make this problem even more acute. AWS RDS instances have default connection limits based on instance size: a db.t3.micro allows 66 connections, a db.r5.large allows 901. Google Cloud SQL for PostgreSQL caps at 500 connections for most instance sizes. Azure Database for PostgreSQL Flexible Server limits vary but are typically in the hundreds.
These limits exist because the cloud provider has sized the instance's memory to handle that many connections without degrading performance. Exceeding them does not just cause memory pressure: the database outright refuses new connections with a "too many connections" error. Your application gets an exception instead of a connection, and the user sees an error page.
Scaling up the instance to increase the connection limit is expensive and has diminishing returns. Doubling the instance size from db.r5.large to db.r5.xlarge doubles the cost but does not double the practical connection capacity because each connection still consumes the same memory and the CPU must handle more context switches. The cost-effective approach is to reduce the number of connections reaching the database through pooling rather than increasing the database's ability to absorb connections.
The critical insight is that connection limits are a shared resource across all application instances. If you run 20 microservice instances against a database with a 500-connection limit, each instance can use at most 25 connections. Add a second microservice that also queries the same database, and the budget shrinks further. Connection pooling is not optional in this environment: it is the only way to fit your application's demand within the database's connection budget.
The problem compounds during autoscaling events. When traffic spikes and your container orchestrator scales from 5 instances to 20, the database suddenly sees 4x the connections. If each new instance opens a pool of 50 connections during startup, the database receives 750 new connection requests in seconds. Without an external connection pooler to absorb this burst, the database may reject connections or become unresponsive at exactly the moment your application is trying to handle increased traffic.
The Real Cost: Latency Under Load
At low traffic, the overhead of open-per-request is tolerable. A single user might not notice 30ms of connection setup on a 200ms page load. But at scale, the overhead stacks multiplicatively.
Consider a checkout flow that touches 4 database tables across 3 queries. With open-per-request:
- Query 1: 30ms setup + 5ms query + 5ms teardown = 40ms
- Query 2: 30ms setup + 3ms query + 5ms teardown = 38ms
- Query 3: 30ms setup + 8ms query + 5ms teardown = 43ms
- Total: 121ms (only 16ms is actual query work)
With a connection pool:
- Query 1: 0.01ms checkout + 5ms query + 0.01ms return = 5.02ms
- Query 2: 0.01ms checkout + 3ms query + 0.01ms return = 3.02ms
- Query 3: 0.01ms checkout + 8ms query + 0.01ms return = 8.02ms
- Total: 16.06ms (essentially all query work)
The pooled version is 7.5x faster because it eliminates the overhead entirely. For latency-sensitive applications like payment processing or real-time bidding, this difference is the gap between meeting and missing SLA targets.
The solution is obvious in hindsight: stop creating and destroying connections for every request. Keep a pool of pre-established connections and reuse them.