System design notes

The notes for 190 lessons of System Design Simplified, by Gaurav Sen, in the order the course teaches them. Lessons marked free can be watched without an account.

Basics

How do I use this course?

  • What do we offer?Free

    This is a detailed walkthrough of the system design course at InterviewReady. If you are looking to prepare for interviews, this walkthrough will help. The walkthrough describes how high-level, low level and fundamental concepts of design are separated into sections inside the course. The platform includes note-taking, quizzes, architecture diagrams, and discussion forums. We also have monthly live zoom classes to discuss doubts and chat about design concepts!

Concepts

Databases Deep Dive

  • What are Databases?Free

    We dive into the fundamental characteristics of databases. Databases as a storage solution can support various types of structured data. They store, update, and retrieve information. Databases are expected to be durable. Different Database Solutions have trade-offs involved as discussed.

  • Storage and RetrievalFree

    We look at the intricate relationship between hardware and software in the context of database storage. The focus is on how data is stored in tables with specific data types, leveraging indexes and algorithms internally. We look at how SQL queries are run, breaking down components such as FROM, WHERE, GROUP BY, HAVING, ORDER BY, and LIMIT.

  • What is a NoSQL database?

    NoSQL databases encompass various types of data storage requirements. Key-value stores have risen in popularity due to their off-the-shelf solutions, like fault tolerance, sharding, etc... We also touch on the concept of database consistency, highlighting its dependence on the application's nature. Databases are generally considered a single source of truth for data-related operations, with an expectation of consistent data. Database Management System (DBMS) manages configurations, lead replicas, and system coordination, serving as a vital component in the overall database infrastructure.

  • Types of Databases: Graph, TimeSeries and Object

    We discuss three types of databases (less used) databases: Graph DB Timeseries DB Object DB

  • What database should you choose?

    The right database for your project depends your app needs, what you're comfy with, and the pros and cons of each choice. For example, Postgres is solid but isn't built for write heavy workloads, while Cassandra is great for fault tolerance but needs tweaking for consistency. Think of a database as a tool that handles data, letting you focus on your app's important parts. SQL queries are the usual way to connect with databases.

Consistency in Distributed Systems

  • What is Data Consistency?Free

    Consistency in a distributed system refers to how up-to-date a piece of data is. A highly consistent system reflects all updates to data, while an inconsistent system provides stale data. Consistency is important because highly consistent systems are easier to reason about and provide a better user experience. The video introduces the levels of consistency that can be provided to users in an application or service API.

  • Linearizable Consistency

    Data Consistency Levels - Linearizable Linearizable At this consistency level, we want to show all changes in the database until the current read request. It means all changes which have happened in the database before the read operation will be reflected in the read query. For example, suppose initially we had x = 10. update x to 13 update x to 17 read x --> Returns 17 update x to 1 read x --> Returns 1 To achieve this we use a single-threaded single server. So every read-and-write request will always be ordered. Using the above example, the first read x will be executed after updating the value to 17. This is useful when systems need perfect consistency. The video summarises this for you!

  • Eventual Consistency

    Eventual Consistency At this level of consistency, we can send stale data for a read request. Eventually, our systems catch up with all updates and return fresh data. To achieve this we can process read and write requests parallelly (using multiple servers) or concurrently (using multiple threads). Example: x = 10 update x to 13 update x to 17 read x --> Could return 10 or 13 or 17 update x to 1 read x --> Could return 10 or 13 or 17 or 1 .... ... Eventually, when all update operations have run through... read x --> returns 1. In the example above, the write request is made before the read request but the read request is processed first. So the client may get any possible value for X. But after a given time, the write request will be processed. Then all read requests will show the correct data i.e., X = 10. Note that eventual consistency is a very loose guarantee, and is often coupled with stronger guarantees like "read your writes". "Read Your Writes" consistency: This type of consistency ensures that a client that makes a write request will immediately see the updated value when it makes a read request, even if the update has not yet propagated to all servers. This type of consistency is often used in systems where low latency is more important than consistency. "Monotonic Read" consistency: This type of consistency ensures that a client will only see values that are as up-to-date or more up-to-date than the values it has previously seen. In other words, once a client has seen a value, it will never see an older value again. "Monotonic Write" consistency: This type of consistency ensures that a client will only see values that are as up-to-date or older than the values it has previously written. In other words, once a client has written a value, it will never see a newer value from another client.

  • Causal Consistency

    Causal Consistency At this consistency level, if a previous operation is related to the current operation, then the previous operation must be executed before the current operation. Causal consistency is stronger and slower than eventual consistency because operations for the same key are processed sequentially. It is looser and faster than the serializable consistency level because it does not wait for all previous operations to complete. Causal Consistency fails when performing aggregation operations. Dive into the video for a detailed explanation!!

  • Quorum

    Quorum Quorum is a mechanism by which we can get consistency guarantees based on our requirements. It defines how many readers must acknowledge a read operation = R. It also defines how many writers must acknowledge a read operation = W. For example, let's say we have three replicas: x = 20 , x = 20 , x = 20. Update second replica x = 40. The second node crashes. Read x => ?? (Depends on the following scenarios) Eventual consistency: W = 1 and R = 1. In this case, when there is a read operation after the second node crashing, we will get a result of X = 20 from the remaining nodes. This is stale data. It will eventually be consistent because the second replica is expected to come back up. When it is back online, we will get the correct result of X = 40. Serializable consistency: W = 2 and R = 2. The write operation must go to two replicas before an acknowledgment is sent to the client. This means, apart from the second replica, either the first or third replica must set X = 40 before a client gets a write acknowledgment. Let us assume one replicates the data. Now we have X = 40, X = 40, and X = 20. On the write operation, we need 2 nodes to reply successfully. Notice that it is impossible to miss the latest data now, because at least one of the nodes we hit will have the last operation reflected in it. The logic is based on the pigeon hole principle: https://en.wikipedia.org/wiki/Pigeonhole_principle The mathematical formula for quorum consistency: Consistency depends on the values R, W, and N. N here is the total number of nodes in the system. We can either have an eventually consistent system (R + W<=N) or a strongly consistent system (R+W > N) Disadvantages of using quorum We need multiple replicas chatting with each other, so costs and latency could be high. If we have an even number of replicas, then it can cause a split-brain problem: https://interviewready.io/learn/system-design-course/data_replication/the_split_brain_problem

  • Data Consistency Levels Tradeoffs

    Data Consistency Levels Tradeoffs ╔═════════════╦══════════════╦══════════╦═════════╦══════════════╗ ║ ║ Serializable ║ Eventual ║ Causal ║ Quorum ║ ╠═════════════╬══════════════╬══════════╬═════════╬══════════════╣ ║ Consistency ║ Highest ║ Lowest ║ Mid-way ║ Configurable ║ ╠═════════════╬══════════════╬══════════╬═════════╬══════════════╣ ║ Efficiency ║ Lowest ║ Highest ║ Mid-way ║ Configurable ║ ╚═════════════╩══════════════╩══════════╩═════════╩══════════════╝ Quest Caps on and dive into the video right now!!

  • Transaction Isolation Levels - Read Uncommitted Data

    Transaction Isolation Levels - Read Uncommitted Data Transactions are a collection of queries that perform one unit of work. They are atomic which means either all queries in a transaction are executed or none of it is executed. BEGIN Marks the start of the transaction. COMMIT Marks the end of the transaction, which results in the changes persisting to the database. ROLLBACK Marks the end of the transaction and undoes all the changes to the database. Isolation: If two transactions are running concurrently and queries in one transaction do not affect the other transaction then the two transactions are said to be isolated from each other. Read Uncommitted At this isolation level, even uncommitted data can be read by concurrent transactions. Although it is very fast, it leads to dirty reads.

  • Read Committed

    Read Uncommitted At this isolation level, one transaction can only read committed data from other transactions. An example would be a movie theater booking system. Reading committed data would be appropriate while booking seats since we ensure that a seat is only counted as booked once the transaction has been committed.

  • Repeatable Reads

    Repeatable Reads At this isolation level, we ensure that when a query reads a row, that row will remain unchanged for the entire transaction. Each transaction has its data copy that is updated. Once the transactions are committed, updates are persisted in the database. This is known as Snapshot Isolation since a transaction is isolated from another by taking different data snapshots. If two transactions concurrently change the same key to different values, we must roll back one transaction. This is called Optimistic Concurrency Control since we are optimistic about our changes until they are proven incompatible.

  • Serializable Isolation Level

    Serializable Isolation Level This is the highest isolation level. At this level, all transactions are executed serially. All of the operations are executed serially or every operation is ensured to not meddle with the operations of other transactions. We use this isolation level when we want to avoid Phantom Reads. A Phantom Read occurs when two identical read queries are executed and the collection of rows returned by the second query is different from the first. It is different from non-repeatable reads because in phantom reads the change in the value can be due to the insertion of a new row. The same is explained in the video with an example. Go ahead!

  • Transaction Level Implementations

    How are transaction levels implemented in the real world? We only need a single entry in the database and it is overwritten whenever there is an update operation If we are making an update to a key then the older value of the key stays in the database and the newer value is kept in the local copy till the commit finally goes through. We take the values that we care about but we are not changing and keep a version of them. For every key, we will store all the values that it has ever had in different transaction commits We use causal ordering here. If two transactions use queries for the same key then they must be ordered. Transactions that do not have any conflict can run concurrently.

  • Conclusion - Transaction Isolation

    Conclusion - Transaction Isolation For Efficiency Read Uncommitted > Read Committed > Repeatable Read > Serializable For Isolation Read Uncommitted < Read Committed < Repeatable Read < Serializable Voila! You have completed this lesson...Wanna test your learning? Go ahead with the quiz, pal.

  • PDF summaryPDF

Caches Deep Dive

  • Caching: Basics

    Benefits of a cache: Saves network calls Avoids repeated computations Reduces DB load Drawbacks of a cache: Can be expensive Potential thrashing Cache Policies: LRU LFU

  • Write Policies: Write Back Policy

    Write Policies: Write Back Policy How do you maintain the consistency of the cache? Let's have a look at write policies. Write Policy A write policy is triggered when there is a write operation in the cache. Keep in mind that it is different from the replacement policy. A replacement policy is triggered when there is no space for a new key and a key is evicted from the cache. A write request means some entry is added, updated or deleted in the cache. But because a cache is a source of truth each of the write requests will impact the entire system. Write-Back Policy If the key-value pair that is to be updated is present in the cache then it is updated. However, the key-value pair is not immediately updated in the database. So as long as the cache is alive, users will get consistent data. However, if the cache is not alive, the data will be stable. To avoid this problem we use: Timeout-Based Persistence Event-Based Write Bac Replacement Write Back Dive into the video right now, and let us help you get the most out of your next dive.

  • Write Through Policy

    Write-Through Policy In this policy, when there is a write request, we evict the key that is being updated, while simultaneously updating the database. The next time there is a read request, that is when the cache polls the database for the entry, persists the entry and sends the response to the user. However, we can run into problems when using this policy. For example, Initially, we have X = 10 There is a write request for X = 20 Then there is a read request for X, but the write request is not updated yet. So the read request returns X = 10 . So it can cause inconsistency. To avoid such problems, we can lock the data which is to be written and only unlock the data after the update operation is completed. This policy is useful When we need a high level of consistency When we need a high level of persistence However, it is less efficient compared to the write-back policy. Let's be real, reading the steps alone is a little bit boring (and sometimes frustrating) so we made this video for you.

  • Write Around Policy

    Write-Around Policy In this policy, when we get a write request, instead of updating the entry in the cache, we update the entry in the database. Now when we get a read request, we will send the stale value from the cache. And we will be getting stale values until the entry is evicted from the cache. This policy is useful When we need a high level of efficiency When we need a high level of persistence However, it makes our system eventually consistent. Go ahead and watch the video to understand it better!!

  • Replacement Policies: LRU, LFU and Segmented LRU

    Replacement Policies In this video, we learn to determine which key needs to be evicted. This is based on the replacement policy LRU Policy (Least Recently Used) In this policy, we evict the entry that has not been used for the longest. To implement this policy, we need to add another data point i.e., Last Used. To understand this policy, watch the video where you can understand it with examples. Time interval in which Write-Around Policy Replacement Policies LRU Policy (Least Recently Used) an entry was not used = current_timestamp - LastUsed for that entry. LFU Policy (Least Frequently Used) In this policy, we evict the entry that is used least frequently. Don't forget to watch the video for a better explanation using examples!!

  • PDF SummaryPDF

Networks Deep Dive

  • Breakdown: The physical layer

    Breakdown -The Physical Layer Wondering what networks are and what problems they solve? Let's break it down from the physical layer. The bonus waiting for you in this video is getting to know how to solve the most common interview problems. In this video you will learn: How messages are read in the networks? Why do we need delimiters? What is the role of the physical layer? Quest caps on for we are about to find out!!

  • Breakdown: The Routing Layer

    Breakdown-The Routing Layer Wondering what is the work of the routing layer? We do know that routers are a part of networks. Nowadays, most houses are equipped with Wi-Fi routers. You might want to know why we need routers, if so... In this video you will learn: How do we route messages in networks? The role of routers explained using a real-case scenario Why are identity and locations important in networks? How do we use logical and physical addressing to find the right receiver? Quest caps on for we are about to find out!!

  • Breakdown: The behavioral layer

    Breakdown-The Routing Layer Wondering what is the work of the routing layer? We do know that routers are a part of networks. Nowadays, most houses are equipped with Wi-Fi routers. You might want to know why we need routers, if so... In this video you will learn: How do we route messages in networks? The role of routers explained using a real-case scenario Why are identity and locations important in networks? And how we use logical and physical addressing to find the right receiver Quest caps on for we are about to find out!!

  • Connecting to the internet: ISPs, DNS and everything in between

    Connecting to the internet: ISPs, DNS and everything in between At times do you wonder how your computer connects to the internet? Now that you know the basics of networking, have you ever thought about where you get data from a server? What is the role of an IP address? Wanna know about the Internet backbone? Dive into the video for we are about to learn about Domain Name Servers and Internet Service Providers.

  • Internal routing: MAC addresses and NAT

    Internal routing: MAC addresses and NAT How does an internal router know where to send an external response to? Do you really wanna find out? Dig into the video right now as in this video we are going to learn: Content Delivery Network(CDN) Problems we face if we do not use CDN Decoding the responses from CDNs

  • HTTP, WebSockets, TCP and UDP

    HTTP, WebSockets, TCP and UDP This video is everything you need to know about Hypertext Transfer Protocols. Let me ask you something. Can a server send a client request to the client by itself? Do you know about the web sockets? Wanna know about Transmission Control Protocol(TCP)? Get diving straight into the video folks.

  • Communication Standards: REST, GraphQL and GRPC

    Communication Standards: REST, GraphQL and GRPC In this video you will learn about the following: In a microservice architecture, services need to communicate with each other. How do servers communicate internally? How do you communicate between services? Different services are written in different languages. But why do we need one common language to define the objects? How do you query the microservices? REST GraphQL gRPC (google Remote Procedure Call) is generally used by microservices to communicate internally. It is written over HTTP 2.0

  • Head of line blocking

    Head-of-line blocking in HTTP Head-Of-Line blocking occurs when the message/data packet at the head of the queue cannot move forward due to congestion even if other messages/packets behind this could. HTTP 2.0 solves this problem using Multiplexing. So what is multiplexing? Why HTTP 3.0 uses UDP instead of TCP? The video above answers all your questions!!

  • Video transmission: WebRTC and HTTP-DASH

    Video transmission: WebRTC and HTTP-DASH What do servers do while handling video data? Video files are large and cannot be sent in a single response. They have to be broken into packets. Why HTTP is not good for video transmission? HTTP is stateless. The client has to specify which chunk it wants because the server does not know about the previous requests. Do you want to learn more that includes the working of Conference Video Protocol? Dig into the video right now!

  • Summary PDF NetworksPDF

Database Replication and Migration

  • Primary Replica Architectures

    Data replication involves copying data from one location to another, while data migration is manually moving data to a new location. Replication is more common in mature companies and those using cloud solutions. In a replication scenario, clients are connected to a server with a database, but the database could fail due to various reasons. To address this, a replica (or backup) of the original database is created. Data consistency is maintained by sending write requests to the primary database, while read requests are directed to both the primary and the replica. This approach benefits read-intensive applications, as write operations are infrequent. The benefits of a primary-replica architecture include fault tolerance, improved read speeds, and relative simplicity. If the primary database fails, the replica can take over. Read requests can be distributed among multiple primary and replica nodes. However, there are drawbacks, including potential consistency issues if there are in-flight messages, and slightly slower write speeds due to the added complexity. Fault tolerance comes at the cost of increased latency. In summary, data replication is an essential practice in maintaining data availability and fault tolerance in software engineering, particularly in mature companies or those using cloud solutions. While it offers benefits, it also introduces challenges related to data consistency and latency.

  • WAL and Change Data Capture

    We explain the process of data replication, focusing on the interaction between the primary and replica databases. Initially, both databases are in the same state, and complex operations are copied onto both of them. When a new write operation occurs on the primary, it is sent to the replica for synchronization. The replica processes the write operation and ensures consistency by applying all previous operations in a sequential manner. This process is described as a "write ahead log" where operations are sequentially executed and can be rolled back, similar to transaction logs. Benefits of this approach include the efficient transfer of only the necessary operations, but there are also potential drawbacks. If the replica cannot catch up with the primary, consistency issues can arise. Timestamps and contextual values can also lead to confusion when the replica and primary show different data values. The solution for addressing these issues is introduced as "change data capture" (CDC). CDC involves sending events that indicate data changes and allow subscribers to process these events and make necessary transformations. This is particularly useful when you have different types of databases, some optimized for writes and others for reads. CDC can transform data and has built-in libraries for connecting to various databases, simplifying data replication. In summary, the transcript explains the process of data replication, the challenges associated with it, and how change data capture can be used to address these challenges and efficiently replicate data across different types of databases.

  • Write Amplification and Split Brain

    We discuss various challenges and solutions related to data replication. It begins by explaining why having two primary databases can be problematic, leading to inconsistencies and difficulties in reconciliation. This situation is referred to as "split brain." Manual reconciliation and techniques like "last writer wins" are mentioned, but they come with their own risks and complexities. To avoid split brain scenarios, the speaker suggests having an odd number of primary nodes to ensure a majority consensus in case of network partitions. The use of an even number of nodes can lead to reconciliation issues and require manual intervention. We also touch on the challenges of replication when dealing with multiple read replicas, such as the issue of "write amplification" and the potential for increased latency. I recommend using a consensus mechanism like Paxos or Raft to ensure that the entire cluster agrees on a value during data replication. This ensures strong consistency across the system. In summary, we need consensus mechanisms and odd numbers of primary nodes to maintain data consistency and avoid split-brain scenarios.

  • Database Migrations

    So you found a flashy new database technology. You run some tests, check the performance, and are convinced! Time to speak to your manager! But wait… How do you move millions of records, each entry representing a user's aspirations, into a new database? Without any downtime? AND without any mistakes… After sweating profusely, you come up with a plan: 1. Stop the existing DB. 2. This fails all incoming requests to the DB (You should probably have some email communication setup beforehand to notify users of a planned outage). 3. Take a Data Dump of the current DB into your shiny new DB. 4. When the process completes, point the existing servers to the new DB. This will take a deployment or a server restart What did you say? You can't afford downtime at all? No worries, we have your back! Watch the video for the full explanation!!

  • Migrating the Database

    This is going to be more complex engineering, but worth it the savings in downtime. 1. Set up a change data capture solution (Similar to a SQL trigger). 2. Take a database copy of all older records. 3. Setup a view or proxy to serve existing clients through the old database. 4. Set the clients to point to the proxy. 5. Now point the proxy to the new database. 6. Point clients to the new database directly (this requires a deployment or server restart). 7. Delete the proxy. Yey! Remember, despite our best efforts, there is always a chance of a data miss. It's best to check for correctness after deployment, and keep your clients in the loop! Watch the full video to know how the optimised approach for data migration works!

  • Migration Across Regions

    Cloud Solution Provider DB Migration Let's say you are moving your database from Microsft Azure to AWS or vice versa. How to handle it? This has to work for NoSQL databases and heterogenous sets of databases as well. DB COPY: DB copy to a new DB. DB-copy with a timestamp>x Add an insert and update trigger at an old-DB MIgrate all missing records from old to new. Timestamp>y and Timestamp<z.

  • Database Migration SummaryPDF

Security in distributed systems

  • Security aspects in distributed systemsFree

    This chapter introduces security protocols and techniques used in distributed systems. We dive into three parts of security Authentication SAML-based SSO (Single Sign On) User Tokens OAuth (Authorization abused for authentication!) Authorization Access Control Lists Rule Engines Secret Keys Attack Vectors Hackers Developers Malicious Code Things we do not discuss here are DRM Solutions Verification methods Rate Limiting

  • Token Based Auth

    In this video by Gaurav Sen, the topic of authentication is discussed, with a focus on using tokens. Tokens are a basic and commonly used method of authentication. The process involves a user sending their username and password to a server for verification. If the credentials match, the server generates a token, which is a piece of text signed by the server. The token represents the user's authorization and specifies the actions they are allowed to perform. The user can present this token to the server as proof of authentication and authorization without needing to provide their password again. The server can verify the token by using its private key to decrypt the signed text and confirm its authenticity. The private key is a secret number known only to the server, while the public key is shared with others. The public key can be used to decrypt the signed text and retrieve the authorization details. The token, being a signed piece of text, prevents unauthorized modifications. However, the token-based authentication method has some limitations. It does not protect against replay attacks or token theft. If someone intercepts the token, they can potentially use it to impersonate the user and perform actions on their behalf. Additionally, if the token is stored in an insecure location, it can be accessed by unauthorized individuals. Despite these limitations, token-based authentication is widely used due to its simplicity and efficiency. The permissions associated with the token can be quickly checked, and logging out invalidates the token. By including a version number or timestamp in the token, its validity can be controlled, and logging out will render it unusable. Overall, token-based authentication offers a practical solution for balancing security and convenience in many authentication scenarios.

  • SSO and OAuth

    The video discusses three different authentication mechanisms: tokens, Single Sign-On (SSO), and OAuth. Tokens: This is the most basic form of authentication. The user sends their username and password to the server, which checks if they match in its database. If successful, the server generates a token, which is a piece of text signed by the server. The token contains information about the user's permissions. The user can then use this token for future requests without sending their password again. Tokens are simple but not extremely secure. Single Sign-On (SSO): In SSO, the authentication process is delegated to an external service, such as Google or Uber. The user's credentials are checked by the external service, which then sends a token to the server. The server can decrypt the token and verify the user's permissions. SSO is useful for companies that want to have control over their users within their system, even if external services handle authentication. OAuth: OAuth is primarily an authorization system but is commonly used for authentication as well. It involves integrating external services, such as Google or GitHub, with a server. The user is prompted to give permissions to the server, allowing it to access certain information from the external service, such as the user's name or profile photo. OAuth tokens are generated, and the server can use them for authentication purposes. OAuth is widely used for authentication, even though it is originally designed for authorization. Overall, these authentication mechanisms involve the generation and use of tokens but differ in terms of where the token generation and authentication process takes place (server-side or external service).

  • Access Control Lists and Rule Engines

    The video discusses various authorization mechanisms used in systems: Access Control Lists (ACLs): ACLs are lists that define the actions that can be performed on objects. They can be user-based (a specific user has certain permissions), role-based (users with a particular role have certain permissions), or group-based (a group of users or resources have certain permissions). ACLs are grouped together to form the access control matrix of the system. Rule Engines: Rule engines involve using a set of rules, often implemented as if statements, to determine whether a user or resource has permission to perform an action. Rule engines are useful for handling complex authorization requirements and can handle multiple objects and rules efficiently. They are also helpful in centralizing common rules and allowing easy rule modifications without modifying ACLs. Secret Keys: Secret keys or client keys are an additional layer of security for authentication and authorization. They involve authenticating using a key or token rather than a username and password. This mechanism provides an extra level of security but is not considered ideal on its own. These authorization mechanisms differ in their approach and can be used depending on the specific requirements of the system. They help define and enforce permissions for users and resources, ensuring that only authorized actions are performed.

  • Attack Vectors - Hackers, Developers and Code

    The video discusses different types of attackers and how to protect resources from them: Hackers: Hackers often launch Distributed Denial of Service (DDoS) attacks to flood the system with requests. To mitigate this, techniques like distributed rate limiting and web application firewalls can be used to verify the authenticity of users and prevent flooding. Employees: Employees, either willingly or unwillingly, can pose a threat to the system. Access control lists are useful in preventing unauthorized actions by employees. Restricting resource access and minimizing attack surfaces, such as opening only necessary entry points and allowing database access only to relevant servers, helps reduce the scope of attack. Malicious Code: It is challenging to completely prevent the execution of malicious code. However, access to resources should be restricted, and rules in a rule engine can be used to prevent illegal modifications of code. Integrity checks and thorough code reviews can also help in identifying potential vulnerabilities. Additionally, the video briefly mentions the importance of virtual private networks (VPNs) or virtual private clouds for secure remote access and end-to-end encryption while communicating with the system. Taking regular backups of the database is recommended as a preventive measure against the impact of malicious code.

  • How videos are protected inside CDNs

    In the video by Gaurav Sen, he discusses a design pattern for protecting videos on a CDN (Content Delivery Network) server from unauthorized access. He presents three approaches: token-based authentication, domain-based authentication, and server-side authentication. Token-based authentication: The CDN server sends each user request to the main server to authenticate whether the user can view the video. If authorized, the CDN serves the video. This method provides strong authentication but is slower and less secure since the CDN can access and misuse the user's token. Domain-based authentication: The CDN server restricts video access based on the domain from which the request originates. If the request comes from an allowed domain, the CDN serves the video; otherwise, access is denied. This approach is simple and fast but vulnerable to domain spoofing and lacks user-specific authorization. Server-side authentication: The user's request directly communicates with the server for authentication. The server generates a token signed by the user's private key, indicating permission to access the video. This token is then forwarded to the CDN, which verifies its authenticity using the public key and grants access accordingly. The token can have a refresh interval to mitigate unauthorized sharing. This method offers a mix of security and efficiency. Gaurav Sen also mentions that for enhanced security, digital rights management (DRM) solutions can be used to prevent unauthorized copying or multiple IP addresses accessing the content. Additionally, access control lists and rule engines can provide resilience against attacks. The video concludes by emphasizing the importance of authentication, authorization, and security mechanisms tailored to the system's specific needs.

Observability in Distributed Systems

  • Introduction to ObservabilityFree

    Observability allows engineers and operations personnel to assess how a system is functioning and identify any issues. The topics covered include logging, monitoring, alerts, anomaly detection, root cause analysis, and incident reporting. By the end of the chapter, viewers should be able to identify areas for improvement in their organization's observability. The goal is to achieve a robust and fault-tolerant system through excellent observability capabilities.

  • Logging - Best Practices

    In this video by Gaurav Sen, the speaker discusses the concept of logging in the context of software systems. Logging is a simple way of noting events, similar to using print statements in programming languages like Java or C++. Log lines typically include a timestamp, an event level (INFO, WARN, ERROR), details about the event, and the line number in the code file where the log originated. Logging helps debugging, allowing engineers to trace and analyze the flow of requests in a system. Examples of log lines and their components, including the masking of sensitive information for privacy and security reasons, are discussed.

  • Monitoring Metrics

    Monitoring involves manually or automatically tracking metrics, which are numerical values representing various aspects of a system's performance. Examples of metrics include the amount of memory used, the number of IO operations, sales made by the system, or the number of visitors on a webpage. These metrics are tracked over time, presented in graphs or time series, and are monitored by operations teams. The dashboard allows for manual decision-making based on the observed data, such as adjusting the number of load balancers in response to changes in request volume. For organizations using cloud solutions like AWS, tools like CloudWatch are available for monitoring and alerting.

  • Anomaly Detection: Holt-Winters Algorithm

    Anomaly detection helps find problems before it's too late. This capability is often offloaded to cloud solution providers. The goal of anomaly detection is to quickly identify unusual patterns in metrics, such as unexpected spikes or collapses in metrics. We explain confidence intervals, and note that this method may not adapt well to seasonal metrics. A more reliable approach involves repeated differentiation of metrics. After differentiating multiple times, anomalies, especially spikes, become apparent and can be flagged as issues. The Holt-Winters algorithm is capable of detecting anomalies even in metrics with seasonal variations.

  • Root Cause Analysis

    Root cause analysis involves understanding what went wrong when a system experiences an anomaly, and aims to identify the underlying cause of the issue. There are two approaches to root cause analysis: manual and automated. The manual approach involves investigating log lines, metrics, and possible causes to build an incident report. The "Five Whys" technique, a series of iterative questions aimed at getting to the root cause of the problem, is a manual method that provides in-depth insights. But it becomes challenging as you delve deeper into the causes. Automated approaches focus on quickly identifying the cause of an anomaly. They involve looking at the factors affecting a metric and determining the root factor contributing to the anomaly. Mathematical algorithms such as principal component analysis (PCA) and Spearman coefficient, which can be used for automated root cause analysis. You can also outsource or using machine learning for root cause analysis. Services like Amazon SageMaker can use algorithms like isolation trees to detect anomalies in metrics.

  • Conclusion to Observability

    This video concludes our learnings on observability in distributed systems.

Distributed Consensus with Paxos

  • The Problem of Distributed Consensus

    In this video by Gaurav Sen, he introduces the concept of Paxos, a distributed consensus protocol. Paxos is an algorithm designed to achieve consensus in distributed systems where individual components can fail, and messages may be dropped during communication. Key points covered in the video: Consensus : Paxos aims to achieve consensus in a distributed system, specifically a majority consensus. This means that more than half of the nodes in the system must agree on a particular value. Distributed Nature : Paxos operates in a distributed environment where components can fail, and messages may not always be delivered. Use Cases : Paxos is used in various popular systems, including Apache Zookeeper and Google Chubby, mainly for distributed locking and maintaining order in distributed logs. Importance : Distributed consensus, especially in the presence of failures and timeouts, is a challenging problem. Paxos offers a solution where consensus is achieved with the agreement of more than half of the nodes. Prerequisites : Gaurav suggests reviewing topics like data consistency, isolation levels, two-phase commits, and quorum to better understand the challenges of distributed consensus. The video sets the stage for a deeper exploration of the Paxos algorithm, its practical use cases, and its applications in systems like distributed locking and log management. Paxos is highlighted as a valuable tool for maintaining consistency and order in distributed systems, even in the face of failures and unpredictable message delivery.

  • Basic Algorithm

    In this video by InterviewReady, the concept of distributed consensus in a scenario involving three application servers (S1, S2, S3) and three database servers (D1, D2, D3) is explained. The key points covered are: Consensus Requirement : Consensus in this context means that for a written request to be considered complete and durable, at least two out of the three databases must accept it. Broadcasting Write Requests : When a write request is received by an application server, it is broadcast to all three databases. Agreeing on a Log Line : Before writing the data, the servers must agree on which log line (position) the data should be written to. This prevents potential overwrites. Log Line Synchronization : Servers ask the databases for their current log line. If more than half of the databases agree on a particular log line, it is selected for writing. This process ensures that the writes are consistent. Locking the Log Line : To prevent concurrent writes and overwrites, a server locks the chosen log line by becoming its owner. This is necessary to maintain consistency and avoid conflicts. Write Operation : Once a server owns the log line, it can proceed with the write operation. The write is accepted or rejected based on the owner of the log line. Two-Phase Protocol : The process involves a two-phase protocol. In the first phase, servers request the logline and lock it, while the second phase involves executing the actual write operation. The video explains how this two-phase protocol ensures that only one server can write to a particular log line at a time, preventing inconsistencies and conflicts. It also emphasizes that the write operation and log line ownership must be accepted by a majority of nodes to maintain consensus. This approach helps achieve distributed consensus in a distributed database configuration, ensuring that data remains consistent even in the presence of failures and concurrent write requests.

  • Improving the algorithm

    To improve the algorithm, servers send timestamps instead of server identities when requesting locks. Timestamps indicate the recency of the request, allowing the system to prioritize the most recent operation rather than solely relying on server proximity. Servers store both the value of the log line and the timestamp. To ensure uniqueness, they assign timestamps with remainders when divided by three, making it easier to compare and identify the most recent requests. This approach separates the concerns of writing and acquiring locks. Servers can write to the database even without initially claiming the lock, and they can make multiple requests with different timestamps during a fresh election.

  • Termination Conditions

    In this video by InterviewReady, the Paxos algorithm is explained in more detail, particularly in the context of a two-phase protocol for writing data to logs. Here are the key points: Roles of Servers (Proposers) and Databases (Acceptors) : In Paxos, servers and databases play distinct roles. Servers do not communicate with each other, and databases only respond to queries without coordinating with each other. Timestamps and Loglines : The video uses the example of a simple write operation (e.g., x += 10 ) to illustrate the process. Timestamps are associated with each operation, and a logline represents the position in the log. Logline Locking : The system avoids overwrites by locking the logline according to the received timestamp. This means a server must receive the logline value from most databases before proceeding with the write operation. Majority Consensus : Paxos requires a majority consensus for operations to be accepted. If a database misses an operation, it cannot provide a consistent logline value, and consensus may not be reached. Termination : The algorithm terminates when more than half of the databases (a majority) agree on an Accept operation. This ensures that the operation is durable and can proceed. Algorithm Safety : Paxos combines elements of two-phase locking and majority voting to ensure the safety of operations. Locking is crucial for maintaining data consistency. In summary, Paxos is a robust algorithm used for achieving consensus in distributed systems, especially in scenarios where write operations to logs need to be coordinated among multiple servers and databases. It ensures that a majority agreement is reached before proceeding with operations, ensuring data consistency and reliability.

  • Practical Considerations

    Several important points about distributed consensus and the Paxos algorithm are discussed: Timing Problems : When multiple servers (S1, S2, S3) send log line requests, there can be timing variations or "jitter." To mitigate this, servers can stagger their requests at different time intervals to avoid simultaneous requests. Algorithm Efficiency : The Basic Paxos algorithm is slow, especially if used for every read or write operation. Instead of selecting a single log line, it's more efficient to choose a leader among the servers. This leader can handle all operations, reducing the algorithm's bandwidth usage. Catch-Up Mechanism : When some servers fall behind in terms of their log lines, there's a need for a mechanism to catch them up. While Paxos doesn't define this explicitly, the intent is clear: servers that lag behind need to be brought up to speed. In practice, optimizations like Multi-Paxos are used for multiple servers to streamline the process. Use Cases : Paxos is utilized in real-world systems, such as Google Chubby and Apache Zookeeper. These systems use Paxos for distributed locking by employing a virtual lock represented as a file in a file system. Implementation and Testing : Paxos is a complex algorithm and typically already implemented in distributed systems. Attempting to implement it from scratch can be costly and challenging due to extensive testing requirements. Engineers should leverage existing solutions unless there's a compelling reason for a custom consensus algorithm. We encourage viewers to understand how these consensus algorithms work underneath, even if they rely on existing solutions, as it can be beneficial for engineers to have a deeper knowledge of distributed systems.

Distributed Rate Limiting

  • The Oracle and the Timer Wheel

    The ways in which rate limiting can be implemented in a distributed environment. The Oracle and the Timer Wheel If our services are overloaded with requests then the time required to send a response will increase resulting in a bad user experience. If the number of requests keeps on increasing then our servers might crash due to Out of Memory exceptions. We want to avoid cascading failure. A cascading failure is a failure that grows over time due to positive feedback. It can occur when one service of a distributed system fails, increasing the probability that other services fail. Let's understand this with an example. We have two servers A and B. Both can handle at most 500 requests. Currently, A is serving 400 requests and B is serving 300 requests. Now suddenly B is getting 600 requests. It cannot handle that many requests and crashes. Now we need to route all the requests to the other server i.e., A. But A cannot handle 400 + 600 requests so it also crashes. This is known as cascading failure. Timer Wheel In a timer wheel, Size of the wheel (Number of buckets) = Timeout of the incoming request Each bucket is numbered from 0 to Timeout -1 Each bucket can store a limited number of requests Every time we add a request to the bucket number Time % Number of buckets Our system can pull requests from this queue one by one. Before inserting a new request into a bucket it deletes all the existing requests (These requests have not been processed by the system). Note: If there are a lot of slots then the distance between any two requests can be very large. In such cases, we will be wasting a lot of system resources and time. To tackle this issue we can use Hierarchal Timer Wheel. Well now we have rate-limited external requests but our internal services are also talking to each other. One service sending a lot of requests to another service might cause the other service to fail. So let's discuss how we are going to solve this issue in the course video!!

  • Partitioning and Real-life Optimisations

    Partitioning and Real-life Optimisations Rate Limiting Internal Requests If we are optimistic we can have a Service Level Agreement which states the number of requests one service can handle. But a better approach is to have rate limiters in the service itself. But how does a service know that it cannot handle any more requests? 1. Average Response Time If the average response time is increasing then either the service is not functioning or it is under too much load and cannot process requests fast enough. 2. Age of requests Let's say the service takes at most 1s to process each request but many requests are present in the queue for more than 2s. That means the number of incoming requests >> the Number of requests the service can process. 3. Dead letter queue Every time our service is unable to process a request it sends the request to the dead letter queue. If the number of requests keeps on increasing in the dead letter queue then we know that our system is failing to process requests. Watch the course video to get in-depth knowledge!!

  • Architecture DiagramPDF
  • Capacity EstimationsPDF
  • Summary PDF - Rate LimitingPDF
  • API ContractsPDF
  • FAQsPDF

System Design Tradeoffs

  • Introduction to design tradeoffsFree

    The video discusses design trade-offs in engineering, specifically focusing on non-functional requirements in distributed systems using the example of Instagram. The speaker highlights three non-functional requirements for uploading pictures on Instagram: recovery from upload failure, safety or persistence of data, and maintaining image quality. Meeting customer expectations and business considerations drive these requirements. Examining one feature or functional condition at a time is essential to identifying the necessary actions to make it successful. Trade-offs are common in engineering systems. This chapter will help you decide when and what you should choose!

  • Memory vs. Latency

    The trade-off between memory and latency is discussed, with a focus on the example of a cache. For popular profiles or highly active users on Instagram, their profiles and follower counts are cached in memory on the server. This allows for faster access and reduces latency to around 10 milliseconds. By keeping frequently accessed data in memory, latency is minimized but it requires more memory and additional costs. On the other hand, reducing the cache size may lead to increased latency for certain requests. A medium approach suggests keeping hot entries in the cache using a suitable replacement algorithm and write policy.

  • Latency vs. Accuracy

    Reducing latency involves avoiding blocking or waiting for resources. Sacrificing accuracy for faster responses involves making guesses on past data. Sacrificing accuracy for latency can be acceptable in certain systems.

  • SQL vs. NoSQL databases

    NoSQL and SQL databases are compared on factors like consistency and fault tolerance. NoSQL databases typically have more relaxed consistency, employ a cluster architecture, and distribute read and write requests among nodes. NoSQL databases are designed for high scalability, allowing fast read and write operations. SQL databases can also have fast read operations, but write operations may be slower due to potential bottlenecks. NoSQL databases often have built-in sharding, which involves distributing data across nodes based on predefined ranges or shards. Sharding is not inherent in SQL databases. NoSQL databases typically lack transaction support, meaning they don't provide ACID guarantees. The choice between NoSQL and SQL databases depends on specific requirements.

  • Relations Between Tradeoffs

    In this video, Gaurav Sen discusses the relationships between different parameters in design trade-offs. Memory: Increasing memory decreases latency because more data can be cached. However, it also increases cost. The impact on consistency and availability depends on the policies used. Throughput: Throughput is inversely proportional to latency. Reducing latency usually means reducing throughput. There is no significant relation between throughput and cost or availability. Cost: Increasing cost may reduce latency because more resources can be allocated. However, there is no direct relation between cost and throughput, availability, or memory. Latency: Increasing cost may reduce latency, but it is not guaranteed. Latency is inversely related to consistency. If the system can tolerate inconsistency, latency can be low. Consistency: Consistency and availability are inversely proportional. Increasing consistency reduces availability. The guarantees of consistency and availability depend more on algorithms than on cost. Availability: Availability is inversely related to consistency. Partition tolerance is an important factor, and a system doesn't need to be perfectly consistent or available, but good enough based on SLAs and user requirements.

Interview Questions

Design an Emailing service like Gmail

Chess Design: Building a highly scalable turn-based gaming website

  • Requirements of a chess website

    We set the following requirements to design a chess website: Users can declare challenges, and accept game challenges. A user should be able to play a game with another user. Users can analyze their games later. Optional (Users can spectate games) Optional (The system should identify cheaters and ban them) Non-functional requirements: The game-play should be smooth, and have low latency. The system must be fault tolerant.

  • Handling connections at scale

    Reducing latency is critical to the success of any gaming platform. "Lag" is the layman term used to describe this frustrating gap in time between the user doing an action and the server acknowledging that action. In this video, we see how latency is taken into consideration while users make moves.

  • Consistent Hashing vs. Sharding

    Rebalancing load, after a server is added or removed from a topology, is an important problem. Have a look at the consistent hashing video in the "Calling App" design for more details. Consistent hashing and its variants are popular solutions to rebalancing load on a server. In this video, we'll see some problems with rebalancing server load.

  • Connection related thundering herds

    Creating connections to users on the server is an expensive process. What happens when we restart this server? To reduce latency, and avoid thundering herds, this video talks about how we can continue maintaining connections in a changing code environment.

  • Request Batching and Conclusion

    This video concludes the chess system design and talks about one last trick to save bandwidth (and compute for parsing responses) on the server-side.

  • Architecture DiagramPDF
  • Capacity PlanningPDF
  • API ContractsPDF
  • Summary - PDFPDF

Hotstar: Live Video Streaming Design

  • Live Streaming Requirement BreakdownFree

    We define the product and engineering requirements of a live streaming system. Primarily, they are: Ingest live HD video. Transform video for different users. Transport videos to end user.

  • Video Ingestion and TransformationFree

    Video Ingestion Using RTMP. Transforming videos into different resolutions and codecs Assigning workers Publishing and subscribing video events

  • Transferring processed videos to end users

    Topics discussed: Content Delivery Networks Adaptive bitrate streaming protocols Network protocols for uploading video Caching video endpoint System interaction and behavior

  • FAQsPDF
  • Capacity PlanningPDF
  • Architecture DiagramPDF
  • Virtual Advertising for live sports eventsFree

    Context based advertising during sports events is a big challenge, especially when serving to millions of users with diverse interests. In this video, we dive into what works behind the scenes like recording calibrations, object recognition, ad ranking, user segmentation and finally video rendering. The question answered here is: How do you serve dynamic ads to different user cohorts in real-time?

Google Docs: Collaborative editor design

  • Google Docs Requirements

    High-level design requirements of Google Docs: Storing, retrieving, and editing documents Sharing permissions: READ, EDIT, LINK-SHARING Notifications on email Version History Spell Checking* Comments Support Offline and Collaborative Editing

  • Document Schema

    Document Database Schema What would be the database schema of a stored document? How would we store the metadata of the document? What technology solution should we use to store documents? File Storage, RDBMS or NoSQL?

  • Storing Documents

    How large can a document be? How do we store a document of any size? Where should we store this document? What type of database is suitable for this?

  • Version History

    Maintaining version history in a document that is updated often is a challenge. With what frequency should we store versions of this document? Where do we store these versions? How do we store changes?

  • Avoiding Thundering Herds in Crons

    Problem : How do you take a single job of size X , to be done in Y time, and have a constant flow of X/Y operations in your system? Challenge : This is a cron job, so the operation is called from time to time. Solution : Hash the job candidates into Y pieces, and run cron jobs at much smaller intervals.

  • Compression and Caching

    We use two ideas here to save on costs, while speeding up queries: Compression Caching the most recent version

  • Concurrent Writes with Locks

    Can we wing it? Are concurrent writes with multiple users in different regions possible with locking (optimistic or pessimistic)? We deep dive to see what works, and what fails.

  • Operational Transform Overview

    What is the heart of the concurrent document update algorithm? We need a method by which multiple users can make changes (possibly offline), on the same document, without completely overriding each other's work. This calls for some reconciliation algorithm. We talk about two major classes (Centralised and Decentralised). Since Google Docs would prefer a centralized algorithm with a single source of truth and low compute requirements on client devices, we choose the operational transform.

  • Permission Management

    How do you manage permissions to update, view or share a document? We use a permission manager service to keep track off all permissions in a document, stored in a NoSQL store.

Design a Cab Aggregator App like Uber

  • Requirements of a cab aggregatorFree

    The video by Gaurav Sen explains the high-level design of a cab aggregator system, using Uber as an example. The system design is centered around a user who books a cab from point A to point B, and the use case can be extended to other apps such as food delivery and pick-up/drop-off services. The video explains the challenges involved in pricing, search pricing, and matching, which are all important components of a successful cab aggregator system. The system relies on real-time calculations to determine the price, dynamic search pricing, and efficient matching of drivers with riders.

  • Static Pricing of Rides

    The video by Gaurav Sen discusses the design of a pricing feature for a ride-hailing app. The pricing feature must be both accurate and fast to ensure that users do not drop off from the app before making a booking. To achieve this, the app uses a cache to store information about regions and their pricing, but a rule engine is a better alternative as it allows for more flexibility and is easier to maintain. Pricing is determined based on the user's location, the time of day, and supply and demand in the area. The video also discusses how the rule engine can be implemented and how to choose the right granularity of regions for pricing purposes.

  • Ride Matching

    The video explains the process of ride-matching algorithms for cab aggregators. The most critical factor is how many riders are available nearby. The algorithm works by selecting a small region where a person wants to book a cab and finding all the available drivers in that area. If their estimated time of arrival is within a threshold, they will be sent a ride request. If none of the drivers accepts the request, the size of the region will be increased. Frequent location updates are crucial to finding the right driver and predicting the estimated time of arrival accurately. The final optimization required is for drivers' paths. The algorithm needs to consider the drivers who will be close enough to accept a ride in the next ten minutes and who will be closest to the person's destination when the ride completes. The estimated time of arrival from the current location to the destination and from the destination to the pickup point should be less than or equal to the threshold time. All the potential pickups need to be run through this equation to select the best match.

  • Calculating Estimated Time of Arrival

    This video explains how to find the estimated time of arrival (ETA) of a person from A to B. Traditional solutions like quad trees or Hilbert curves fail because the representation in terms of a continuous plane or blocks of squares is not helpful. We store information as a graph with nodes as points and edges for rules. To minimize time, an algorithm like A* search is required, which is more tuned towards costs. The video explains how A* search works and how to store the points using a database.

  • Matching Drivers with Live Rides

    The video discusses the process of ride-matching for cab aggregators. Location updates and the estimated time of arrival are critical in finding the right driver for a passenger. The ride-matching algorithm works by finding nearby riders and sending them a request to accept the ride. If none accept, the proximity radius is increased, and the request is sent to more drivers. The video also mentions that optimization for drivers is necessary, considering their current location and the next ride they'll take. Finally, the video introduces an equation to determine the right driver for a passenger, taking into account the estimated time of arrival from the driver's current location to the destination and then to the pick-up point.

Design a Location Based Service like Google Maps

  • Requirements of a Map applicationFree

    The video discusses the system design of a map application like Google Maps, which is used for estimated times of arrival and providing the best possible routes. The primary responsibility of the application is to track the location of the user and update them with the best possible routes. It also discusses the feature of adding locations and finding nearby stores, like restaurants and gyms. Searching for nearby locations is filtered using a Cartesian distance and ETA is an important feature. Traffic updates and estimated time of arrival is also important. Here is what we will be covering: Finding the best route from a source to destination Adding Locations in the map Finding Nearby Stores Notifying users with traffic updates Calculating the estimated time of arrival

  • Routing Challenges in Maps

    Finding the optimal route amongst an infinite number of intermediate points, with potentially infinite number of paths, is hard. We tackle the following two problems in this video: How do we reduce the number of points to contend with in a map? How do we efficiently find a solution route amongst a large number of potential solutions?

  • Partitioning Algorithm: Splitting the graph into regions

    The video explains how Google Maps deals with the problem of too many possible routes between two points. One solution is to reduce the number of routes by finding popular hubs or common routes that most people use, even though this may not be the optimal route. The other solution is to compute the result faster, which may require approximation. Google uses graph partitioning algorithms and artificial intelligence to split roads or entire regions in an optimal way to find these hubs. This technique can also be applied to public transport like bus and metro stations.

  • Finding the shortest path - A* search

    In the video, we talk about how Google Maps uses an algorithm called A-star search to find the best route between two locations. This algorithm takes into account the time taken to move from one position to another as well as the estimated time to reach the final destination. To get accurate times for the edges, Google Maps relies on actual, real-time data from people moving from one position to another. The best way to store this information is to use a graph database that can store points and edges. We recommend using Neo4j, which provides built-in graph algorithms.

  • GeoHash and other proximity filters

    We discuss the technology used in Google Maps, specifically the location tracking, adding new locations, and finding nearby stores. Gaurav Sen explains that location tracking involves sending a request to the server every five seconds to notify the server of your location. The speaker then suggests using a graph database to add new locations, and also suggests using open-source libraries to render locations on the client. To find nearby stores, we suggest using a keyword search and filtering out stores that are too far away. Finally, we discusses the use of Geo hash to build a hash of a location by splitting the 2-D space recursively.

  • Detecting traffic jams and broadcasting updates

    We discuss how Google Maps traffic updates work. The traffic updates are relevant when a user is traveling from a particular position to a destination, and there is an update on an intermediate road that the user is supposed to travel on. There are three factors that Google Maps uses to update the time it takes to finish a stretch of road: user updates, average time spent in the last five minutes, and the number of users in a particular stretch. To update a user's ETA, Google Maps needs to store the active routes of each user and inform them of any updates to their route. The server is responsible for giving the user the most efficient route, even if it means going through the same road.

  • Calculating an accurate and efficient estimated time of arrival

    We discuss the factors involved in estimating the time of arrival (ETA) using Google Maps. The ETA algorithm depends on routing and is a function of the route chosen. The video explains how Google Maps breaks down regions into subregions or hubs, making the routing process more efficient and accurate. The goal is to provide a highly predictable ETA so users can plan their trips and arrive on time. Google Maps limits the number of potential routes between regions to improve accuracy and reduce computation time.

WhatsApp Calling App

  • Requirements ListPDF
  • Calling App Design

    We design the calling app. This video highlights a good attempt at the design, typical of someone with less than 4-5 years of experience in the industry. Things to notice: The candidate does a very good job at describing what the system will do. However, they do not talk about the deep problems, since they haven't spent time to expand on the requirements.

  • Concept #2: The state machine

    We model the call's lifecycle as a state machine, and move both state storage and transition logic to the Call State Manager.

  • Concept #3: Charging Users

    We decide how the billing service should handle call invoicing, decoupling services and polling for call termination.

  • Concept #4: Consistent Hashing for caching call state

    Load Balancing is a key concept to system design. One of the popular ways to balance load in a system is to use the concept of consistent hashing. Consistent Hashing allows requests to be mapped into hash buckets while allowing the system to add and remove nodes flexibly so as to maintain a good load factor on each machine. The standard way to hash objects is to map them to a search space, and then transfer the load to the mapped computer. A system using this policy is likely to suffer when new nodes are added or removed from it. Consistent Hashing maps servers to the key space and assigns requests(mapped to relevant buckets, called load) to the next clockwise server. Servers can then store relevant request data in them while allowing the system flexibility and scalability. Some terms you would here in system design interviews are Fault Tolerance, in which case a machine crashes. And Scalability, in which case machines need to be added to process more requests. These two principles are allowed by Consistent Hashing, and hence it is an important building block to a system design architect's toolbox. Another term used often is request allocation. This means assigning a request to a server. Consistent hashing assigns requests to the servers in a way that the load is balanced are remains close to equal. Server architecture is a subjective concept, and there are outliers for many cases. Don't think of Consistent Hashing as a silver bullet for fault tolerance and scalability, but a useful concept for request allocation. Use it to solve software questions in interviews and real life. Best of luck!

  • Calling App ArchitecturePDF
  • API ContractsPDF
  • FAQsPDF
  • Capacity EstimationPDF
  • FAQs 2PDF
  • Summary - PDFPDF

Recommendation Engine Design

Case Studies

Scaling Memcached at Facebook

Google's Authorization Service

Timeseries Databases: Gorilla and Monarch

TAO - An in-memory graph database by Meta

Google Dremel - Interactive Analytics

Meeting Recordings

Zoom Meets

  • October - Kafka Internals

    What was cooking up interesting in the zoom class?? What's life without exploring? In this zoom class let us explore Kafka! How many of us know about the Kafka internals? Kafka is ubiquitous. The questions that come up to our mind when we want to dig deep into Kafka are: How is Kafka Internally built? It's a log. A log is a file where u write one line and when there is a new event you write another line. Every line will have a delimiter. Every line is appended after the delimiter. Why Kafka doesn't work on queues(priority)? Queues are implemented using linked lists or an array. The mechanism of Kafka doesn't give any priority. Kafka is fast, scalable, and distributed in nature by its design, partitioned and replicated commit log service. So there is no priority on topic or message. All events belonging to a topic need to be ordered. For example, Let's say we have topics A B C and P Q R. For instance, if all this is jumbled like P Q A B C R Rules: P comes before Q And Q comes before R The original order need not be maintained. We need to care about ordering between inside the topics not in between topics. Wherein C will be consumed by the consumer only after B gets consumed. What is Causal ordering? As long as relevant messages are sent to the same partition, message (causal) ordering is still guaranteed. In Kafka, this is achieved using keyed messages. Kafka partitioners would hash keys and ensure a key always goes to the same partition given the number of partitions is not changed. How do u search in a log? A log has a key and value. The log is already sorted by key. How we usually query is based on timestamp or the machine id. Instead of taking any random key, we take the 128 bits key and use 64 bits to store the timestamp. For Example: 18th August_123 Where ‘_’ is used as a delimiter. Elastic search The basic idea is it is a no-sequel database. It has key-value pairs. It is good at text search and aggregations. You can add tags to it. An Elasticsearch index is a collection of documents that are related to each other. Benefits of using Kafka Limited persistence Fast in terms of writing(if u want to read fast use another database). Used for debugging Finally, If a very popular search text exists how do u find the document that it exists in and store it efficiently? I leave it here for you to discover. Let me give out a little secret! This zoom class has the answer to this question!

  • October 2022 - Implementing a File System

    In this zoom video, get ready to start with interview-ready tips. Let's say there are 5-6 servers and there is a load balancer... How are the requests monitored? If we have authentication in the load balancer it becomes a gateway that checks whether a request is valid or not. If we need extra authentication needed( for example, only admins can send requests to a particular server), the question arises of where we should have the authentication logic. If you are up for scaling then you must be opting Load Balancer Authentication. The drawback of this method is that you will have to deploy a load balancer every time a request is sent. Another drawback is that it is a single point of failure. Who will handle authentication when the Load balancer is gotten rid of? AuthService We have an Auth service which takes in request from load balance them validates and send them back. What if I tell you we have a video to explain this? Go ahead and watch the video on Authentication and Global caching under the High-Level Design. Service Registry If we are going to have a sidecar to take care of invalidation, how do you maintain the states? The consistency becomes a question. The state here is the cache that we use. What is a sticky session? Dive into the video on time stamp 23: 50 to know more. Authentication for multiple systems works in different ways. The architecture-wise where you give the authentication matters in terms of performance. You have a bunch of externally exposed services one of the things you can do is if it is coming from the service which is sending the request can say whether it needs to be authenticated or not. If internal services are not talking to each other and the sidecar feels that the request need not be authenticated, the authentication can be bypassed. Benefits Of Compression: Lower bandwidth usage Faster to send Lesser memory being used Why shouldn't we compress messages? Computational overhead Time-taking latency. To know how to implement a file system in an interview dive into the time stamp- 1:19:00 and keep your knowledge caps open. For this problem, we need to be creating : 1. Create(file_id, content) 2. Read(file_id) 3. Update(file_id, content) 4. Delete(file_id) It's time for discovery. Go ahead, explorer. The zoom class answers this question!

  • Pessimistic and Optimistic Locking

    How to count the frequency of words from 1TB of file in 2 GB RAM consistency having only a single machine? There are 2 things that could be done in this case. To split the whole file into contents of 2 GB. The other approach would be to stream the file while tokenizing it and counting the number of words until it reaches the required size. Data Structures that can be used: Hash Maps The entire words of the file will be converted into a key set of 1 Terabyte long if there are distinct words in the file. The edge would be storing 1 next to all the hashed words of the file which will be altered with changing frequencies of the words. We might need concurrent hashmaps to deal with this problem The scatter-and-gather approach can be used to sort out the data. Have a look at the external Merge Sort algorithm (Which can also be used to solve this problem) Splitting the file into pieces is the complexity. Multi-threading can be used to solve this problem. Let's say we have a fileSplitter(n), There are 1,2,3 threads that are being passed. The splitter will split the file into pieces based on the number of words such that the limit memory is reached. We will have to stream through these chunks of threads. We will need to have seek() and parse() functions to go ahead with this method. Finally, we store all of this in hashmaps. Quick Fact!! BookMyShow makes use of Optimistic locks !! There are 2 copies to make sure your booking doesn't clash with already booked seats. While pessimistic locks won't create copies. The timestamps will help you to reach the point in the video where you can see a full video explanation of the same. Why not solve your doubts? When you are connecting to a server you are connecting to a port where multiple people connect. The server can identify who is connected to it. Let's go back to TCP connections, We take a message and route it through multiple routers and send it to a person. It is more of a logical concept rather than a physical one. Let's Move on to a system design problem: There is an e-commerce app which has a lot of price fluctuations for its products. They have a catalogue of schools that is already in the ERP systems. They want to add a new component to the e-commerce platform. Get to know the tips on how to solve this!! Come lets clear a user's doubt on over-utilizations of resources.

  • Producers and Consumers: November Meet

    This zoom meeting can never get interesting enough!! Get to know the resources on where to apply for SD roles from this zoom meet intro part.. Some websites that were discussed: Instahyre Cutshort Geektrust Angel.co Along with that why not hustle with a few more resolving doubts and general discussions? How do we handle 5000 transactions(writes) per second on a relational database with a max latency of 50 milliseconds? Moving ahead with the learning spree.. The first thing to note here is what type are consistency are looking at and what type of translation isolation level is needed? Choosing the translation isolation level in a wise manner will help increase concurrency. The requirements for the questions are below: The user should get appropriate tasks to work on based on skill language. Example: A user who only knows Hindi should be assigned Hindi tasks where these tasks should be served in order of priority and first come first serve basis. By default, users should get English tasks as well. The system management should also be able to view the life cycle of a job and we expect the people to be working in parallel. The system must be able to respond to the task within 15 milliseconds if the user is waiting for the job or not( If he is free). Let's take a queue and enter the tasks. This means that we have a priority queue which is implemented by a heap-based on which the tasks are served. Let's say we have multiple queues. The problem to address here is the priority and skill level of the worker and the language might be different. We will have a head-of-line blocking problem when we have the tasks blocking the assigning of work to the users based on priority when using only one queue. We can have a list of queues for the languages task and the queue is following FIFO(as we know) if the task on the front doesn't get assigned to any worker we upload that task to a separate log queue. We use the approach of using a sorted tree instead of a priority queue which makes the problem complexity O(log(n)). This is very similar to the Operating System Resources problem. Now, what happens when we use a hash map(brute force algorithm)? Depending on the language we have a worker. Then we see the skill level of the worker. To remove the problem where 2 workers pull the same task we shall assign locks. Acknowledgements for the acceptance or assignments of the tasks need to be mentioned. Some sort of expiry algorithm would do. Question: There are 2 services A and B deployed using EKS. What is the best way to communicate with each service? Also, where do we need to keep the URL? Go ahead and watch the video to know what is Gaurav’s take on this!!

  • IRCTC Algorithm: December

    We discuss how the IRCTC algorithm possibly works. It's a complex algorithm, but worth understanding! Irctc Algorithm Train ticketing system How do we efficiently book tickets in the ticket booking system so the seats don't remain vacant? It's a High-level design question. Example: There are 10 stations and some person X books seats from 1 to 5 and some other person wants to book from 6 to 8. How do we handle that? How do we know who booked prior? Here the exciting part is that we not only have to consider the allocations, but also we will have to search for trains. We draw the UML diagram by considering the following: •There is a Client who wants to book a seat through their web browser. •The clients say that they want to book a journey with the specifications • This is the case when the client wants to book a specific train. • The other option is when the client searches for all trains available to reach the particular destination with the number of seats, and start date. • The train should have information about the estimated time of arrival, estimated departure time, source, destinations, category and seats. Multiple perspectives and multiple calculations..let's head to see the exact answer. How do you book the seats? The other interesting fact is that the centre of gravity has to be taken into account or the train will be heavier on one side alone. If a single woman is travelling she needs to be placed with other women travelling on the same train. All these are the considerations that need to be considered. However, the brute force approach would be: Book from A - B ( 2 seats) The algorithm is like a rule engine. It goes through the preferences and chooses the best option for the client. The time the rule engine runs to check all the preferences will become a matter of consideration here. The duplication is good because it will make the search faster. Instead of a beat matrix, we can flatten it by making it an array. So how do we approach this? Let's head to the video!

  • Google Docs and Chess Design Queries

    Google Docs and Chess Design Queries When we use Google docs when we are sharing a document with someone else. Both will be able to edit the content of the docs. This makes it difficult to design the algorithm as we will have to lock based on the ranges. A specific algorithm is used for these kinds of test collaborations. The Operational Transform The algorithm is very complex. There is a video available in the System Design course for your reference. Let's move to the Chess design: The Chess design is going to have a gateway. The gateway speaks to another service that stores session data which does a protocol translation. In a chess design, the app should make regular changes. Different clients speak to the app to hit the services. The problem: If the API has a breaking change then all the clients need to be updated and restarted. When changing a version of the library you need to update and restart the services. The key here is to keep the gateway dumb as possible so that you do need not to restart the gateway. •Minor changes are not relevant • Services are not relevant • Breaking changes are not relevant The service registry needs to be dumb as possible. Apart from a reverse proxy it also acts as a parser. Let's say you are in version 3 of Google Docs and google docs says that a person deleted a character from the doc that will be stored in a new version. What happens is that one keeps drifting away from the lattice point. Eventually, the events are mapped to the current state. The algorithm will tell you how to do that. Each operation is mapped which will end up in the final result. Further, let's see what exciting discussions we have and how to approach them. Head to watch the full video to learn more!

  • Jan 2023: DNS, CDNs, and everything in between

    The process of navigating to a website, specifically mentioning the role of DNS servers and Content Delivery Networks (CDNs) in the process. When a user types in a website address, such as amazon.com, the request first goes to a DNS server which then determines which CDN to direct the request. The CDN acts as a first line of defense against Distributed Denial of Service (DDOS) attacks, and then the request goes to the real server (content server), which delivers the content. Protection from malicious attacks can be done at the CDN or Gateway level. He gives the example of a website hosted on S3, with the CDN being backed by S3 and replicating static content all around the world for low latency. The DNS server can determine the closest location for the user to connect and provide the appropriate IP address.

  • Jan 2023: API design and other discussions

    Gaurav discusses using consistent hashing in sharding and how requests are mapped to specific buckets in some distributed systems like Redis. There is a potential for one bucket to become overloaded due to a lack of uniform distribution. What is the potential latency when consumers pull during a load rebalancing event? Graph API vs. Open API. There are various ways to overcome the problem of an increase in logs, including compressing and archiving logs, scaling the storage for logs, and identifying which service is causing the increase in logs. It is impossible to stop an increase in logs as user activity increases, but steps can be taken to manage them more efficiently. Additionally, discussing with teams why they are logging so much is a practical consideration. The use of APIs and understanding how they work are essential. GraphQL may not be as helpful as REST in certain situations and clients must understand the objects beforehand. When there is a large request load on the API gateway, it can be scaled horizontally to handle the load, but that code can also be made more efficient. YouTube handles a similar scenario, but when there is too much load, the cached content can be used. There is a need to change and degrade the service in the API gateway when there is too much load. Gaurav shares his experience preparing for technical interviews, specifically for DSA. It's difficult to find an equilibrium between system design and DSA prep. It's essential to speak with the interviewer about their thought process and approach to a problem. For DSA, execution is very important, and being able to understand and catch hints is also essential. In an interview, it's best to be prepared and not overlook the importance of paying attention to the current needs. One complex DB query as an API versus using three different simple queries as three APIs: depends on the requirement, but usually, three simple queries are better, but sometimes it may be necessary to use one complex query for an atomic operation. We found a solution for handling large requests and responses for an API. The idea of breaking down a large request into smaller fragments and returning responses is applicable here. The input and decision on how to break down the request can depend on the business use case and that there is a separation of ingestion rate and persistence rate. We may also use an asynchronous approach, where data is uploaded and inserted into a database after all files have been processed and cross-checks can be made on the records.

  • Feb 2023: Permission Management System

    The video by Gaurav Sen covers a range of topics. The discussion starts with Kafka and SCS, which are message queueing systems. Gaurav suggests some resources for learning about Kafka, including Confluent's documentation. The conversation then moves on to managing permissions and naming conventions for meetings. The video ends with a design problem that involves designing an authentication system for microservices with complex rules for user actions.

  • March 2023: Eventual Consistency Levels

    In the video, Gaurav Sen talks about the importance of consistency in distributed systems. He gives an example of how likes on a post in a social media platform may not be immediately consistent across all nodes in the system, but it is still acceptable because it does not affect the overall experience. He also explains how consistency can be ensured in financial systems through settlements and reconciliations, and how a sorted string table can be used to maintain a record of changes in the system. Sen emphasizes that the level of consistency required depends on the application and the user's needs.

  • March 2023: Fetch top K hits in a distributed system

    The video by Gaurav Sen shows a live session in which he discusses system design, specifically the difference between machine coding and low level design. He also discusses the types of questions asked during SD interviews and gives an example of a recommendation system. In the latter part of the video, he demonstrates how to solve a system design problem using Redis and hash maps to return the most popular blog posts in the last five and 15 minutes.

  • April 2023: Distributed Consensus and Load Balancing

    The video discusses the issue of ordering events in a message queue or event bus with multiple topics, using examples of services posting response times and events like create, update, and delete. Enforcing order by timestamps or using distributed consensus algorithms like Paxos or Raft is mentioned as a solution. The use of logical clocks like hybrid logical clocks or Lamport clocks is also discussed as an alternative to the Spanner implementation, which requires Google-level hardware. We also discuss the importance of enforcing order in the queue to ensure the sequential execution of instructions in a multi-threaded environment. It explains how a thread pool can mimic a queue and how events can be sent to different threads depending on whether they need to be sequentially ordered. The video also addresses questions related to hashing, partitioning, and handling multiple events across different services. Finally, it suggests that ordering across services can be complex and require engineering coordination.

  • September - Design Judge Launch

    We attempt multiple questions of the system design judge!

  • October - Live Watch Party Design

    We discuss the architecture and tradeoffs that come up when designing a live watch party system.

Additional Resources

Basics

  • Horizontal vs Vertical ScalingFree

    Systems design a procedure by which we define the architecture of a system to satisfy given requirements. It is a technique by which the required amounts of scalability, reliability, performance and consistency are satisfied in real world systems. We discuss how to start with system design and why it is required. The first concept in designing a system is scalability. We discuss the two main approaches to solve this problem: Horizontal scaling and vertical scaling. Horizontal scaling is adding more machines to deal with increasing requirements. These machines handle requests in parallel to improve user experience. Vertical scaling is replacing the current machines with more advanced machines to improve throughput and hence response time. The techniques are used in conjunction in real world systems.

  • Monoliths vs MicroservicesFree

    Microservices are a hot topic in system design interviews. It is important to know why we use them instead of monolithic systems. The short answer is: Scalability. The detailed one would be: Advantages: 1) The microservice architecture is easier to reason about/design for a complicated system. 2) They allow new members to train for shorter periods and have less context before touching a system. 3) Deployments are fluid and continuous for each service. 4) They allow decoupling service logic on the basis of business responsibility 5) They are more available as a single service having a bug does not bring down the entire system. This is called a single point of failure. 6) Individual services can be written in different languages. 7) The developer teams can talk to each other through API sheets instead of working on the same repository, which requires conflict resolution. 8) New services can be tested easily and individually. The testing structure is close to unit testing compared to a monolith. Microservices are at a disadvantage to Monoliths in some cases. Monoliths are favorable when: 1) The technical/developer team is very small 2) The service is simple to think of as a whole. 3) The service requires very high efficiency, where network calls are avoided as much as possible. 4) All developers must have context of all services.

  • Load Balancing

    Load balancing is distributing work evenly onto a set of resources. In this video, we explain the main reason for load balancing, along with some common algorithms: Round robin Geo-distributed Least Connection

  • Load Balancer Summary PDFPDF
  • ShardingFree

    Sharding a database is a common scalability strategy used when designing server-side systems. The server-side system architecture uses concepts like sharding to make systems more scalable, reliable and performant. Sharding is the horizontal partitioning of data according to a shard key. This shard key determines which database the entry to be persisted is sent to. Some common strategies for this are reverse proxies. Database interviews ask for concepts like sharding to make databases more performant and available. This makes horizontal partitioning a logical choice.

  • Summary - Database ShardingPDFFree
  • Single Point of FailureFree

    A single point of failure(SPOF) in computing is a critical point in the system whose failure can take down the entire system. A lot of resources and time is spent on removing single points of failure in an architecture/design. Single points of failure often pop up when setting up coordinators and proxies. These services help distribute load and discover services as they come and leave the system. Because of the critical centralized tasks of these services, they are more prone to being SPOFs. One way to mitigate the problem is to use multiple instances of every component in the service. The graph of dependencies then becomes more flexible, allowing the system to resiliently switch to another service instead of failing requests. Another approach is to have backups that allow a quick switch over on failure. The backups are useful in components dealing with data, like databases. Allocating more resources, distributing the system and replication are some ways of mitigating the problem of SPOF. Hence designs include horizontal scaling capabilities and partitioning. It is important to note that the CAP theorem does not allow removing SPOFs if perfect consistency is required.

  • Service discovery and HeartbeatsFree

    Servers crash due to various reasons like hardware faults and software bugs. Service Discovery and Health Checks are essential for maintaining a service ecosystem's availability and reliability. We talk about how a heartbeat service can be used to maintain system state and help the load balancer decide where to direct requests. Now when a server crashes, the heartbeat service and identify and restart the service immediately on the server. Service Discovery is another important part of deploying and maintaining systems. The load balancer is able to adapt request routing. Both features allow the system to report and heal issues efficiently.

  • Capacity Planning and Estimation: How much data does YouTube store daily?Free

    Back-of-the-envelope calculations are often expected in system design questions. They help logically state the parameters influencing a result, and estimating the capacity requires multiple estimations on the way. Also lets us individually state our assumptions. Eg: Estimate the hardware requirements to set up a system like YouTube. Eg: Estimate the number of petrol pumps in the city of Mumbai. Chapters 00:06 Storage Requirements 01:20 Supplementary storage requirements 03:54 Back of Envelope calculations 05:38 Youtube caching estimation 08:58 Youtube video processing estimation 12:14 Conclusion ------STORAGE Let's start with storage requirements: About 1 billion active users. I assume 1/1000 produces a video a day. Which means 1 million new videos a day. What's the size of each video? Assume the average length of a video to be 10 minutes. Assume a 10 minute video to be of size 1 GB. Or... A video is a bunch of images. 10 minutes is 600 seconds. Each second has 24 frames. So a video has 25*600 = 150,000 frames. Each frame is of size 1 MB. Which means (1.5 10^5) (10^6) bytes = 150 GB. This estimate is very inaccurate, and hence we must either revise our estimate or hope the interviewer corrects us. Normal video of 10 minutes is about 700 MB. As each video is of about 1GB, we assume the storage requirement per day is 1GB * 1 million = 1 PB. This is the bare minimum storage requirement to store the original videos. If we want to have redundancy for fault tolerance and performance, we have to store copies. I'll choose 3 copies. That's 3 petabytes of raw data storage. What about video formats and encoding? Let's assume a single type of encoding, mp4, and the formats will take a 720p video and store it in 480, 360, 240 and 144p respectively. That means approximately half the video size per codec. If X is the original storage requirement = 1 PB, We have X + X/2 + X/4 + X/8 == 2*X. With redundancy, that's 2X 3 = 6 X. That's 6 PB(processed) + 3PB (raw) == 10 PB of data. About 100 hard drives. The cost of this system is about 1 million per day. For a 3 year plan, we can expect a 1 billion dollar storage price. Now let's look at the real numbers: Video upload speed = 3 * 10^4 minutes per minute. That's 3 * 10^4 * 1440 video footage per day = 4.5 * 10^7 minutes. Video encoding can reduce a 1-hour film to 1 GB. So 1 million GB is the requirement. That's 1 PB. So the original cost is similar to what the real numbers say. If we are off by order of magnitude, it's good. However, being off by 3 or more orders of magnitude is too much. We can then highlight the following: Where our assumption was wrong, or Which factor we didn't take into account.

  • Content Delivery Networks

    Content Delivery Networks are a bunch of servers spread across the globe to serve information. These networks are available on rent to deliver static content quickly to nearby users. Some examples of CDNs are Amazon CloudFront and the Akamai CDN. They are (relatively) cheap to rent and have high availability. They also provide pluggable algorithms to invalidate and fetch data. We discuss why content delivery networks are useful, and how they serve data with examples.

Design an Audio Search Engine like Shazam

  • Design an algorithm for an Audio Search Engine

    This video talks about Designing an audio search engine like Shazam. We shall try to find the "closest" sound match of a given audio clipping using audio spectrum analysis and hashing. Algorithm: Graph the frequency and amplitude of each point in time in the song. Find the points having high amplitude and frequency variation in this graph. Take all k consecutive points to get a set of chunks. Take the combinatorial hash of each chunk. Store these hashes in the DB.

  • Capacity EstimationPDF
  • FAQsPDF
  • Summary - PDFPDF