Skip to content
Rafe Uddaraj

The Ultimate System Design Handbook: Global Scale Architecture Through Real-World Case Studies

25 min readEnglishRead in Bangla
On this page

When we start a new software project, we usually focus our energy on fast feature development, clean code, and building a polished user interface. However, the real test of our code and architecture only begins when actual user traffic starts hitting our production servers. A standard Monolith application can easily handle a load of 1,000 users, but when that user base scales to 1 million or 10 million, bottlenecks begin to emerge across every layer of the system. Servers crash, database tables lock up, and page load times stretch past 3 to 4 seconds-a delay that directly translates to lost revenue in the modern web era.

System design is not about memorizing theoretical definitions; it is the practice of real-world engineering trade-off analysis. The core philosophy of this handbook is to identify bottlenecks through real-world case studies and engineer the correct architectural solutions. Rather than reading through standard dictionary definitions, we will analyze the scaling journey of a fictional startup platform called "FoodFast" alongside the actual system design patterns used by global tech giants. If you are a software engineer looking to take your skills to the next level and build production-grade systems, this handbook will serve as your complete, end-to-end roadmap.


Foundational Mental Models and Engineering Analogies

Before writing any code, we need to internalize a few real-world mental models to understand how large-scale systems fit together. When we design complex architectures, the same fundamental patterns appear over and over again. To understand these patterns easily, we can map them to familiar analogies from our daily lives. Below, we discuss four core engineering concepts alongside their real-world counterparts.

  • Vertical vs Horizontal Scaling (The Elevator vs. New Stairs Analogy): When the foot traffic in your building increases, upgrading the elevator's motor to carry more weight is Vertical Scaling, or scaling up. However, there is a physical limit to how large that motor can be before you hit mechanical constraints. Instead of overworking that single motor, building a new staircase or installing additional elevators right next to it represents Horizontal Scaling, or scaling out. In software engineering, rather than buying an infinitely large CPU or adding endless RAM to a single machine, we connect multiple standard servers across a network to share the workload.
  • Load Balancer (The Traffic Officer Analogy): When thousands of cars suddenly flood an intersection, a traffic officer stands in the center directing vehicles into different lanes to keep traffic moving smoothly. In our software systems, a Load Balancer performs that exact same duty. When millions of client requests hit our infrastructure, the load balancer uses round-robin or other routing algorithms to distribute that pressure evenly across our backend servers.
  • Monolith vs Microservices (The Single Kitchen vs. Food Court Analogy): A Monolith is like running an entire restaurant out of one massive kitchen where all prepping, cooking, and plating happen together. If a major gas pipe leak occurs in that kitchen, the entire restaurant shuts down. On the other hand, Microservices operate like a large food court where the pizza booth, burger stand, and coffee shop are completely separated. If the espresso machine breaks down at the coffee shop, the pizza booth continues selling slices without interruption, keeping the overall business operational.
  • Database Indexing (The Book Index Analogy): If you need to find a specific topic in a 1,000-page book, flipping through every single page from cover to cover is a Full Table Scan, which runs at an time complexity. Looking at the index at the back of the book to jump directly to the correct page is Indexing. In our databases, we implement indexing using B-Tree data structures, allowing us to pinpoint the exact data we need from billions of records in time.

Case Study 1: FoodFast MVP Launch and the First 10,000 Users (Monolith & Separation)

Imagine you are the Chief Architect of a new food delivery startup called "FoodFast." Your primary goal is to build a Minimum Viable Product (MVP) rapidly to validate your business idea in the market. At this stage, your budget is tight, and your engineering team is small. Therefore, instead of building a complex distributed architecture right away, you launch using a single-server Monolith application connected to a single relational database. For the first few months, everything runs smoothly, and your daily active user count hovers around a few hundred. However, after a successful digital marketing campaign, your user base explodes rapidly, reaching 10,000 active users.

Right at this milestone, your system hits its first major engineering bottleneck. You notice that during peak lunch and dinner hours, the application's page load time increases dramatically, and the server's CPU utilization gets pegged at 100%. The root cause is simple: your web server code and your database queries are competing for the exact same CPU cores and memory resources on a single machine. As thousands of users simultaneously browse restaurant menus and place orders, memory exhaustion and disk I/O contention take over. Your application's thread pool becomes completely blocked while waiting for database queries to return, preventing the server from accepting any new incoming client requests.

To solve this initial crisis, our first architectural move is App and DB Separation. We physically or virtually decouple the web server and the database into two separate machines. This allows the application server's CPU and thread pool to focus entirely on handling incoming HTTP requests, while the database server utilizes its own dedicated memory and disk resources to process queries. At the same time, we apply Vertical Scaling by upgrading our server hardware. We bump the application server from a 4-core CPU to a 16-core CPU, and we upgrade the database server's RAM from 16 GB to 64 GB. This immediate hardware upgrade restores system speed and stability. Below is an example of a production-ready database connection pool configuration after separating the app and database layers:

JavaScript
// db-pool-config.js
// Production-ready PostgreSQL connection pool setup using 'pg' module
const { Pool } = require('pg');
const pool = new Pool({
host: process.env.DB_HOST || 'db-primary.foodfast.internal',
port: parseInt(process.env.DB_PORT || '5432', 10),
database: process.env.DB_NAME || 'foodfast_production',
user: process.env.DB_USER || 'admin_user',
password: process.env.DB_PASSWORD,
// Limit the maximum number of clients in the pool to prevent memory exhaustion
max: 25,
// Close idle clients after 30 seconds of inactivity
idleTimeoutMillis: 30000,
// Return an error after 2 seconds if connection could not be established
connectionTimeoutMillis: 2000,
});
pool.on('error', (err, client) => {
console.error('Unexpected error on idle database client in FoodFast cluster', err);
process.exit(-1);
});
module.exports = {
query: (text, params) => pool.query(text, params),
getPool: () => pool,
};

Connection Pooling Efficiency

Opening a new database connection for every incoming request is an expensive operation because it requires a fresh TCP handshake and authentication step each time. Using a connection pool keeps a set of reusable connections open in memory, which significantly reduces query latency.


Case Study 2: FoodFast Flash Sale Crisis and 1 Million Users (Horizontal Scaling & Load Balancing)

As FoodFast's active user count hits 100,000, your marketing team announces a "Midnight Flash Sale" to drive further business growth. At exactly 12:00 AM, partner restaurants will offer a 50% discount on popular menu items. At 11:59 PM, 100,000 users log into the platform simultaneously, and at the stroke of midnight, thousands of order requests flood the API at once. Within seconds, your powerful 16-core web server crashes from an Out of Memory (OOM) error. The entire platform goes dark, and frustrated users immediately begin posting negative reviews across social media.

This failure proves that vertical scaling has strict physical and financial limits. You cannot simply install unlimited RAM or CPUs into a single server box, and relying on one machine introduces a Single Point of Failure (SPOF) into your architecture. The only sustainable way out of this crisis is to implement Horizontal Scaling. Instead of relying on one expensive super-server, we deploy 10 standard commodity servers running in parallel across our network. However, this introduces a new routing challenge: how do our clients' mobile apps and web browsers know which specific server IP to send their traffic to?

To solve this routing problem, we place a Load Balancing Layer directly in front of our application server cluster using tools like Nginx or HAProxy. This load balancer handles Layer 4 (Transport Layer) or Layer 7 (Application Layer) routing. We configure Nginx to use the Least Connections algorithm, which inspects the cluster and routes incoming requests to whichever server currently has the fewest active workloads. However, running across multiple servers requires us to make our application tier completely Stateless. Previously, user login sessions were stored in the local memory of our single server-a pattern known as sticky sessions. Now that requests can land on any server in the cluster, we move our session management out of local memory and into a centralized Redis session store, while adopting stateless JSON Web Tokens (JWT) on the client side.

Nginx
# nginx-load-balancer.conf
# Layer 7 Load Balancing configuration for FoodFast backend cluster
upstream foodfast_backend_cluster {
# Routelessly distribute load to the server with the least active connections
least_conn;
server app-node-01.foodfast.internal:3000 max_fails=3 fail_timeout=10s;
server app-node-02.foodfast.internal:3000 max_fails=3 fail_timeout=10s;
server app-node-03.foodfast.internal:3000 max_fails=3 fail_timeout=10s;
server app-node-04.foodfast.internal:3000 max_fails=3 fail_timeout=10s;
server app-node-05.foodfast.internal:3000 max_fails=3 fail_timeout=10s;
# Maintain open keepalive connections between Nginx and backend servers
keepalive 64;
}
server {
listen 80;
server_name api.foodfast.com;
location / {
proxy_pass http://foodfast_backend_cluster;
proxy_http_version 1.1;
proxy_set_header Connection "";
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header Host $http_host;
# Strict timeout configurations to prevent hanging requests during spikes
proxy_connect_timeout 3s;
proxy_read_timeout 7s;
}
}
Monolith to Horizontal Scaling Evolution
Monolith to Horizontal Scaling Evolution

Stateless Architecture Requirement

When using a load balancer, never store user session data, temporary file uploads, or local state inside the memory of your application servers. If your servers are not strictly stateless, users will be logged out randomly or lose their data whenever the load balancer routes their next request to a different server node.


Case Study 3: When the Database Becomes the Bottleneck & Global Expansion (Replication, Sharding & NoSQL)

With horizontal scaling and load balancing in place, FoodFast's web server cluster can easily handle hundreds of thousands of requests per second. Furthermore, your business has expanded beyond a single city; you now operate nationwide and in several neighboring countries, pushing your total user base past 1 million. Right at this stage of global expansion, a new severe bottleneck surfaces. While your 10 application servers process HTTP requests independently, they all hammer the exact same relational database to read and write data.

This heavy traffic load causes severe database read and write contention. Your database connection pool is frequently exhausted, and complex transaction queries begin throwing deadlock errors. When we analyze the database slow query logs, we discover that restaurant search requests are performing full table scans. Searching through a table of 100,000 restaurants to find open locations in a specific neighborhood is pegging the database CPU at 100%. Our first remediation step is to apply B-Tree Indexing to the most frequently searched database columns. This simple indexing step reduces our search query time complexity from down to , dropping execution times from 500 milliseconds down to just 5 milliseconds.

However, indexing only optimizes read speeds. Every time a user places a new order (a write operation), the database must update and rebuild its index trees, which adds storage space complexity and slows down write performance. To fix this, we move to a distributed database architecture by implementing Database Replication. We set up a single Leader (write) database paired with 4 Follower (read) databases. Since 90% of our application traffic consists of read operations-such as browsing menus and searching for restaurants-we route all read traffic across our read-replica followers. We reserve the leader database exclusively for write operations like checkout and payment processing, allowing the leader to asynchronously replicate its data changes down to the followers.

Eventually, our core orders table surpasses tens of millions of rows, making it impossible for a single leader database to hold all that write traffic efficiently. This is when we implement Sharding, or horizontal partitioning. We choose the user's geographic location or city ID as our Sharding Key, splitting our monolithic database into smaller, independent regional databases. For example, all users in Dhaka are routed to Shard-A, while users in Chittagong are routed to Shard-B. At the same time, we realize that a relational SQL database is not the best tool for every data type. While we must maintain strict ACID (Atomicity, Consistency, Isolation, Durability) compliance for orders and payments using PostgreSQL, we migrate our high-volume, unstructured data-such as real-time GPS delivery driver tracking and customer support chat logs-over to NoSQL databases like MongoDB and Cassandra. Our design choices here are heavily guided by the CAP Theorem (Consistency, Availability, Partition Tolerance); we prioritize strict consistency for our billing systems while prioritizing high availability for our live location tracking feeds.

Distributed Database Replication and Sharding
Distributed Database Replication and Sharding

Database Sharding Complexity

Think carefully before implementing database sharding. Once you shard a database across multiple physical machines, running cross-shard join queries or maintaining multi-table transactions becomes extremely complex. Always exhaust your indexing, caching, and read-replication options before partitioning your primary database.


Case Study 4: Boosting Performance to Rocket Speed (Caching Strategies & Edge Computing)

FoodFast's database layer is now distributed and horizontally scaled. Your servers are stable, and orders process without crashing the platform. However, when reviewing your application analytics reports, you notice that the mobile app still takes roughly 2.5 seconds to load the home screen. In today's fast-paced internet ecosystem, a 2.5-second load time is slow enough to cause significant user drop-off as customers switch to faster competitor apps. When we trace the network request waterfall, we see that every time a user opens the app, our application servers query the database from scratch to fetch restaurant lists and popular menu items. Reading this data from physical disk storage and dealing with network round-trip delays is wasting valuable milliseconds.

The most effective engineering tool for dropping latency from milliseconds down to microseconds is In-Memory Caching. We integrate in-memory data stores like Redis or Memcached directly into our architecture. Any data that is read frequently but rarely updated-such as restaurant menus, item prices, and store categories-is cached directly in system RAM. We implement the Cache-Aside pattern for our read-heavy endpoints. When a client request arrives, the application first checks the Redis cache (a cache hit). If the data is found, it is returned immediately without touching the primary database. If the data is missing (a cache miss), the application queries the PostgreSQL database, returns the payload to the user, and writes a copy of that data into Redis to serve subsequent requests.

Because RAM is an expensive hardware resource with limited storage capacity, we cannot keep billions of food items cached in memory indefinitely. To manage this footprint, we configure proper Cache Eviction Policies. We set our Redis cluster to use the Least Recently Used (LRU) algorithm. When memory reaches its capacity limit, the system automatically purges data keys that have not been accessed recently, freeing up RAM for active, hot data. Furthermore, to maintain data consistency across our platform, we assign a logical Time To Live (TTL) expiration value to every cache key. This ensures that when a restaurant manager updates an item's price, users do not see stale pricing data lingering in the cache.

Additionally, our platform hosts thousands of high-resolution food images and promotional videos uploaded daily. When a user scrolls through the app, loading these heavy static media files directly from our primary origin servers consumes massive network bandwidth and slows down rendering time. To solve this delivery bottleneck, we integrate a Content Delivery Network (CDN) such as Cloudflare or AWS CloudFront. The CDN caches copies of our media files across globally distributed edge servers located close to our end users. Now, when a customer opens the app, food images load directly from the nearest edge location rather than traveling across the world to our origin data center. This drastically reduces Round Trip Time (RTT), making the application feel instantaneous.

TypeScript
// redis-cache-aside.ts
// Production implementation of the Cache-Aside pattern using Redis and Node.js
import { createClient } from 'redis';
import { getPool } from './db-pool-config';
const redisClient = createClient({ url: process.env.REDIS_URL });
redisClient.connect();
export async function getRestaurantMenu(restaurantId: string): Promise<any> {
const cacheKey = `foodfast:menu:${restaurantId}`;
try {
// Step 1: Check if the menu exists in the Redis in-memory cache
const cachedData = await redisClient.get(cacheKey);
if (cachedData) {
console.info(`[Cache Hit] Serving menu from Redis for restaurant: ${restaurantId}`);
return JSON.parse(cachedData);
}
// Step 2: Cache Miss - Query the primary PostgreSQL database
console.warn(`[Cache Miss] Fetching menu from DB for restaurant: ${restaurantId}`);
const dbPool = getPool();
const queryText = `
SELECT id, name, description, price, category, is_available
FROM menu_items
WHERE restaurant_id = $1 AND is_deleted = false
`;
const result = await dbPool.query(queryText, [restaurantId]);
if (result.rows.length === 0) {
return null;
}
const menuData = result.rows;
// Step 3: Write data to Redis Cache with a TTL of 3600 seconds (1 hour)
await redisClient.setEx(cacheKey, 3600, JSON.stringify(menuData));
return menuData;
} catch (error) {
console.error('Error in getRestaurantMenu execution:', error);
// Fallback: If Redis fails, gracefully degrade and fetch from DB directly
const dbPool = getPool();
const fallbackResult = await dbPool.query('SELECT * FROM menu_items WHERE restaurant_id = $1', [restaurantId]);
return fallbackResult.rows;
}
}

Cache Stampede Prevention

When a popular cache key expires (its TTL drops to zero), thousands of concurrent requests can hit the primary database at the exact same millisecond-a failure mode known as a cache stampede. To prevent this, implement mutex locks or distributed locking mechanisms so that only a single background thread queries the database and repopulates the cache while other requests wait briefly for the updated payload.


Case Study 5: The Monolith Breakdown and 50 Developers (Microservices & Event-Driven Flow)

FoodFast has now evolved into a major enterprise technology enterprise. Your engineering department has grown to over 50 software developers, QA engineers, and DevOps specialists working concurrently. Despite this growth in headcount, you notice that releasing new features to production takes much longer than it used to. The root cause is your legacy monolithic codebase. Every piece of business logic-user authentication, restaurant onboarding, shopping cart management, payment processing, delivery rider matching, and push notifications-is tightly trapped inside a single, massive Git repository.

This monolithic architecture creates severe tight coupling across your engineering teams. When a junior developer attempts to modify an email template inside the notification module, a minor syntax error causes the entire payment processing pipeline to crash in production. With 50 engineers pushing code to the same repository daily, merge conflicts and broken builds become routine. Fixing a small bug requires compiling, testing, and deploying the entire monolith, which takes several hours per release cycle. Furthermore, when our payment system needs more CPU during peak hours, we cannot scale that specific function alone; we are forced to scale the entire monolithic server footprint, resulting in massive hardware resource waste.

To escape this maintenance nightmare, we migrate our entire platform to a Microservices Architecture. Following Domain-Driven Design (DDD) principles, we break our monolith down into smaller, independent services aligned with specific business capabilities. We create a dedicated Auth Service, Restaurant Service, Order Service, Payment Service, Rider Matching Service, and Notification Service. Each microservice is assigned its own independent database and CI/CD deployment pipeline. Now, the billing engineering team can deploy updates to the payment service 10 times a day without coordinating with or risking stability for the rest of the company.

To manage communication between our client applications and these numerous backend services, we deploy an API Gateway as a single entry point. Our mobile apps no longer need to track the internal IP addresses or URLs of dozens of individual microservices. The API Gateway receives all incoming client traffic, handles SSL termination and user authentication, and routes the requests to the appropriate downstream service. For inter-service communication, we adopt two distinct architectural patterns. When a service requires an immediate response-such as checking if a user is authenticated during checkout-we use synchronous communication via REST or gRPC protocols. The gRPC protocol uses Protocol Buffers to transmit data in a compact binary format, running up to 10 times faster than standard JSON over HTTP.

For operations where an immediate client response is unnecessary, we implement asynchronous communication. For example, when a customer successfully places a food order, the Order Service does not make synchronous HTTP calls to the Notification Service or the Analytics Service and wait for them to respond. Instead, the Order Service publishes an OrderCreatedEvent message to a distributed message broker or event bus, such as Apache Kafka or RabbitMQ. Independent consumer services-including Payment, Notification, and Analytics-subscribe to that Kafka topic, ingest the message, and execute their respective background tasks (such as generating invoices, sending SMS alerts, and dispatching delivery riders) independently. This Event-Driven Architecture keeps our system loosely coupled, highly responsive, and effortlessly scalable.

Microservices Architecture and Event-Driven Kafka Flow
Microservices Architecture and Event-Driven Kafka Flow

Case Study 6: Building an Bulletproof System for Disaster Recovery (Fault Tolerance, Rate Limiting & Observability)

In a distributed microservices ecosystem, one absolute law holds true: anything can fail at any time. On a busy Friday night during peak dinner hours, our third-party banking payment gateway partner suffered an internal system outage. Their servers began taking over 30 seconds to respond to payment authorization requests. Our FoodFast Payment Service kept waiting for those responses, which quickly exhausted its internal thread pools and open outbound connection limits. Because our Order Service made synchronous calls to the Payment Service, it also ran out of available threads and crashed. Within minutes, this cascading failure brought down our entire microservices ecosystem like a falling stack of dominoes.

To prevent this catastrophic failure mode from happening again, we engineer strict Fault Tolerance into our architecture. We implement the Circuit Breaker Pattern across all inter-service network calls. Just as an electrical circuit breaker trips during a power surge to prevent your house from catching fire, a software circuit breaker monitors network failures to protect your system. Using patterns established by libraries like Resilience4j or Hystrix, we configure our circuit breakers to trip into an "Open" state if a downstream service fails or times out on 50% of its requests within a 10-second window. While the circuit is open, our application stops sending requests to the failing service entirely and immediately returns a graceful fallback response (for example: "Online card processing is temporarily unavailable; your order has been switched to Cash on Delivery"). This protects our primary thread pools from blocking and keeps the rest of the platform online.

Simultaneously, we protect our infrastructure against malicious actors, web scrapers, and Denial of Service (DDoS) attacks by enforcing Rate Limiting. We implement the Token Bucket or Leaky Bucket algorithm at the API Gateway layer, establishing strict traffic rules that limit individual IP addresses or user accounts to a maximum of 20 requests per second. If a script or scraper exceeds this limit-such as a competitor attempting to scrape our entire menu catalog-the API Gateway immediately blocks the traffic and returns an HTTP 429 Too Many Requests status code, shielding our backend cluster from synthetic load spikes.

Finally, managing 50 microservices running across thousands of containers requires complete internal visibility, so we build a comprehensive Observability platform. We structure this observability strategy around the three core pillars: Logging, Metrics, and Tracing. We aggregate application logs from every container and stream them into a centralized ELK Stack (Elasticsearch, Logstash, Kibana) for real-time querying. To track system health, we use Prometheus to scrape hardware metrics (CPU, memory, disk usage) and application metrics (request latency, error rates), visualizing them on live Grafana dashboards. To trace requests across our distributed network, we deploy Jaeger or Zipkin for Distributed Tracing. When a user complains about a slow checkout, we can inspect a single trace ID and immediately identify which specific database query inside which microservice is causing the latency spike.

TypeScript
// circuit-breaker-pattern.ts
// Production-ready implementation of the Circuit Breaker pattern using the Opossum library
import CircuitBreaker from 'opossum';
import axios from 'axios';
// The potentially unreliable external service call (e.g., Bank Payment Gateway)
async function callExternalPaymentGateway(payload: any): Promise<any> {
const response = await axios.post('https://bank-gateway.external.com/api/v1/charge', payload, {
timeout: 3000, // Strict 3-second timeout to prevent thread blocking
});
return response.data;
}
// Circuit Breaker configuration options for the FoodFast payment service
const breakerOptions = {
timeout: 3500, // If function execution takes longer than 3.5s, trigger a failure
errorThresholdPercentage: 50, // When 50% of requests fail, trip the circuit to OPEN
resetTimeout: 30000, // After 30 seconds, transition to HALF-OPEN state to test service health
};
const paymentCircuitBreaker = new CircuitBreaker(callExternalPaymentGateway, breakerOptions);
// Fallback logic executed immediately when Circuit is OPEN or service fails
paymentCircuitBreaker.fallback((payload: any, error: any) => {
console.warn(`[Circuit Breaker OPEN] Payment gateway failed for order ${payload.orderId}. Executing fallback.`);
return {
status: 'ACCEPTED_FALLBACK',
paymentMethod: 'CASH_ON_DELIVERY',
message: 'Online payment is temporarily unavailable. Order converted to Cash on Delivery.',
originalError: error.message,
};
});
paymentCircuitBreaker.on('open', () => console.error('CRITICAL ALERT: Payment Circuit Breaker TRIPPED to OPEN!'));
paymentCircuitBreaker.on('halfOpen', () => console.info('Payment Circuit Breaker is HALF-OPEN. Testing gateway health...'));
paymentCircuitBreaker.on('close', () => console.info('Payment Circuit Breaker CLOSED. Gateway operating normally.'));
export async function processOrderPayment(payload: any): Promise<any> {
return await paymentCircuitBreaker.fire(payload);
}

Distributed Tracing Overhead

Never enable 100% request sampling for distributed tracing in a high-traffic production environment. Capturing a trace for every single HTTP request consumes massive network bandwidth and quickly fills up storage clusters. Instead, configure sample rates between 1% and 5%, which provides more than enough statistical data to monitor system performance and catch latency bottlenecks.


Case Study 7: End-to-End Blueprints of Global Tech Giants (URL Shortener, News Feed & Video Streaming)

Having explored the individual engineering layers of a scalable architecture through the FoodFast journey, we will now synthesize these concepts by analyzing how global tech companies design their core production systems. Below, we break down three real-world architectural blueprints from scratch.

Case Study 7A: URL Shortener Design (Bitly or TinyURL Style)

The primary responsibility of a URL shortening service is to take a long web URL, generate a compact unique short link, and reliably redirect users back to the original long URL whenever they click the short link.

  • Trade-Off and Ratio Analysis: In a URL shortening platform, the volume of users clicking short links is vastly higher than the volume of users creating new links. We design the system assuming a 10:1 read-to-write ratio. If the platform generates 100 new short links per second, it must handle 1,000 redirect requests per second. Therefore, our architectural focus centers on optimizing read performance and achieving low redirect latency.
  • Unique ID Generation and Encoding: To generate a compact short link, we need to create a 7-character unique key. We use Base62 encoding, which utilizes alphanumeric characters (a-z, A-Z, 0-9). A 7-character Base62 string provides , or roughly 3.5 trillion unique combinations-enough to support system growth for decades without running out of keys. Rather than relying on slow database auto-increment locks, we deploy a Distributed Sequence Generator like Twitter Snowflake or a dedicated Key Generation Service (KGS). The KGS pre-generates blocks of unique keys and holds them in memory, allowing application servers to grab new short keys instantly without database write latency.
  • Data Layer and Caching Strategy: We store the mappings between short keys and original long URLs inside a horizontally sharded NoSQL database such as DynamoDB or Apache Cassandra, as this data model requires no complex relational table joins. To drive redirect latency down to near-zero milliseconds, we deploy a massive Redis caching cluster in front of the database. When a user clicks a short link, the application checks Redis, grabs the original long URL from memory, and immediately returns an HTTP 301 Permanent Redirect or 302 Temporary Redirect status code to send the browser to its destination.

Case Study 7B: News Feed System Design (Facebook or Twitter Style)

The core engineering challenge of a social media news feed is generating a real-time, chronologically sorted timeline of posts from hundreds of friends and followed accounts for millions of concurrent users.

  • Fan-Out Architecture (Write vs. Read Push/Pull Models): When a user publishes a new post, the process of delivering that post to the timelines of all their followers is called Fan-Out. There are two primary architectural strategies here. The first is Fan-Out on Write (the Push Model). When a standard user publishes a post, a background service grabs their follower list and pushes the post ID directly into the in-memory timeline caches (Redis lists) of every individual follower. The main advantage here is that when a follower opens the app, their timeline is already pre-computed in RAM, loading instantly without database queries.
  • The Celebrity Problem and Hybrid Solutions: However, the push model breaks down when encountering the "Celebrity Problem." Imagine a famous user with 50 million followers publishes a new post. Pushing that single post into 50 million individual timeline caches simultaneously creates a massive write storm that can overwhelm servers and crash the caching tier. To prevent this, platforms use Fan-Out on Read (the Pull Model) for accounts with large follower counts. When a celebrity publishes a post, it is simply saved to their personal timeline database table without pushing to followers. When an average user opens their app, the Feed Generation Service merges their pre-computed push timeline from Redis with an on-the-fly pull query fetching recent posts from the celebrities they follow. This hybrid push/pull model is the standard pattern used in production by platforms like Facebook and Twitter.

Case Study 7C: Video Streaming Platform Design (YouTube or Netflix Style)

Building a video streaming platform requires handling massive file uploads, optimizing network bandwidth utilization, and delivering smooth video playback across variable internet connection speeds without buffering.

  • Video Ingestion and Distributed Transcoding: When a content creator uploads a raw 4K video file, serving that huge file directly to end users on mobile devices is impractical. The uploaded video lands in an object storage bucket (such as AWS S3 or Google Cloud Storage), which immediately fires an event into an Apache Kafka message queue. A fleet of background Transcoding Workers consumes this job from the queue and uses tools like FFmpeg to convert the raw video into multiple standard display resolutions, including 1080p, 720p, 480p, and 360p.
  • Chunking and Adaptive Bitrate Streaming (ABR): During the transcoding process, each video resolution is sliced into short, 2- to 4-second file segments called chunks (saved as .ts or .m4s files), and a master playlist or manifest file (.m3u8 or .mpd) is generated to map them together. This segmented delivery approach is known as Adaptive Bitrate Streaming (running over HLS or DASH protocols).
  • CDN Edge Delivery and the Client Player: We distribute these transcoded video chunks and manifest files out to global Content Delivery Network (CDN) edge servers. When a user presses play, their video player downloads the manifest file and assesses their current network connection speed. If the user has a fast fiber connection, the player fetches high-definition 1080p chunks from the nearest CDN edge node. If their Wi-Fi signal drops or their cellular network degrades, the client player smoothly drops down to fetching 360p chunks on the fly, keeping the video playing continuously without freezing or buffering.
YouTube and Netflix Video Streaming Transcoding Pipeline
YouTube and Netflix Video Streaming Transcoding Pipeline

Developer Checklist and Production Guidelines

Successful system design is about choosing the right architectural trade-off at the right time. Whenever you design a new platform or refactor a legacy system, review this production-ready engineering checklist. Walking through these points will help you catch critical architectural flaws before they turn into production outages.

  • Scalability and Statelessness: Are your web servers and API tiers configured as completely stateless applications? If a running server node is abruptly terminated, will client sessions and pending requests fail gracefully without losing data?
  • Caching and Expiration Policies: Is your caching tier configured with an appropriate eviction policy like LRU or LFU? Have you assigned a logical Time To Live (TTL) value to every cache key to prevent memory leaks and eliminate stale data?
  • Database and Query Indexing: Have you applied B-Tree indexes to the database columns used most frequently in search WHERE clauses and JOIN conditions? Have you run EXPLAIN ANALYZE on your slow queries to verify that they are using index scans instead of full table scans?
  • Synchronous vs. Asynchronous Processing: Are background workloads that do not require an immediate user response (such as sending confirmation emails, generating PDF reports, or logging analytics) pushed to an asynchronous message queue like Kafka or RabbitMQ?
  • Fault Tolerance and Circuit Breakers: Have you implemented strict network timeout rules and the Circuit Breaker pattern on all outbound REST and gRPC calls to third-party APIs and internal microservices to prevent cascading system failures?
  • Rate Limiting and Security: Is your API Gateway configured with Token Bucket or Leaky Bucket rate limiting rules to protect your backend cluster against brute-force attacks, aggressive web scrapers, and DDoS traffic spikes?
  • Single Points of Failure (SPOF): Have you eliminated every Single Point of Failure across your infrastructure by deploying redundant standby nodes for your load balancers, database masters, and caching clusters?
  • Observability and Alerting: Do you have Prometheus and Grafana dashboards actively scraping CPU, memory, latency, and error rate metrics from your production containers? Are automated PagerDuty or Slack alerts configured to notify your on-call engineering team the moment an error rate threshold is breached?
All articles

Engineering Roadmap: Why 'Exactly Once' Delivery Is Mostly a Lie

A complete engineering guide on why exactly-once delivery is mathematically impossible at the network level, and how to build effectively-once architectures in production using idempotency, transactional outbox, and inbox patterns.

System DesignENBN

Get in touch

Questions about a video, an article, or working together.