Skip to main content

Command Palette

Search for a command to run...

Database at Scale

Updated
•11 min read•View as Markdown
B
I am a third-year undergrad student at Chandigarh University. I am into building backend and applied AI stuff using the modern technologies. I loves to document my learning through the blogs so that I can contribute to the community.

Stateless and Stateful Architecture

What is state ? State is the information about previous interactions that are made.

  • Stateless: These are systems where each request is completely independent and is unaware of the previous interactions. They are horizontally scalable. For ex - Request/Response cycle of http.
    REST APIs, Serverless Functions are stateless systems.

  • Stateful: These are the systems where the interaction are remembered between the requests. They are not horizontally scalable. For ex - A session is created in order to remember the user's preferences while making http request.
    Databases, web sockets are stateful systems

As we know we have mainly two types of scaling i.e. Vertical Scaling and Horizontal Scaling. So in this blog we are going to go deep into how databases are scaled, what are the different ways to do so and many more interesting topics.

Why we need Scaling ?

Databases itself are servers which communicate with our application servers using TCP over internet. So as our number of users increases our application layer starts getting more requests to handle. And if the application is getting more requests our database server will gonna be receiving more read and write queries. This increase in traffic can overwhelm our database which can cause database crash and failures. And if databases crashes then the whole application will be of no use and if by any chance the data in the databases gets corrupted then it will be a big loss for the whole company. So to handle this we need to scale databases and protect it from any sort or corruption.

We will be starting with Vertical Scaling at first then we will move to Horizontal(distributed) Scaling. We will come across different techniques in each type.

Vertical Scaling of Databases

So by this word ’Vertical’ what we can think? Vertical Scaling is the way where we increases the same server’s hardware specification like RAM, CPU and storage in order to handle more requests.
For example - We have our PostgreSQL server running on a machine whose specs 4 cores CPU, 8GB RAM and now our databases queries are failing, the CPU RAM usage is touching more than 90% mark. So now we have added more powerful CPU like 8 cores CPU and increase the RAM to 16GB. Now our usage metrics came down to 50% which is nice to handle the current increased number of user. This technique in terms of scaling is known as Vertical Scaling as the server specs are vertically increasing.

Lets first come across different techniques of doing Vertical Scaling:

  • Increasing Hardware: As we have discussed above using an example that base condition of vertical scaling is upgrading the hardware’s specs. Upgrading CPUs, RAMs, storage type(hard dish, SSD, NVME) and increased bandwidth.

  • Query Optimization: Optimizing the read/write queries for the database on the single machine is also the way to make one machine more performant. Using proper indexes, choosing views over joins and avoiding full scans.

  • Table Partitioning: This technique is little interesting. It is of two types row based and column based(it is rarely used). Lets say we have 10M rows in a table, so rather keeping all this data in a single file the database splits that file into smaller segments called partition. So there are mainly three types of partitioning we make:

    1. Range Partitioning: Lets sat we have a very large orders table containing the history from last 5 years. So in this type of partitioning we split the data in the form of ranges (in this case it is order_date). So the orders of 2020 lies in separate partition, 2021 in other and so on.

    2. List Partitioning: Here we split the tables on the basis of a fixed value like we separated the tables using the value of country means users with country ‘India‘ will be on separate partition and users of ‘US’ will be on other partition.

    3. Hash Partitioning: In this technique we don’t use some specific range of value rather it is used for even distribution using a hash function to go to their respective partition. Example: if user_id % 4 equals 0 then it will go in partition 1 otherwise in partition 2.

      NOTE: Internally when a a write query runs the databases checks the partitioning rule and accordingly routes the query to the required table where the data gets stores physically. Similarly the read query also ignores the whole partitioned tables which is not required.
      This process is also known as pruning where we have just optimized the query performance by ignoring the irrelevant data by routing it to the targeted table.

So all the above mentioned techniques are used in vertical scaling in order to scale our database to handle more request but there are several downsides of using this which makes it difficult to scale further. Some of them are:

  • High Infra Cost: As we are upgrading the hardware so the cost increases exponentially because machines with higher specs are more expensive.

  • Hardware Limitation: Every cloud provider has limitation of hardware upgradation. We cannot have unlimited number of cores and unlimited ram in a machine. So for very large companies it is not feasible.

  • Single point of Failure: If something goes wrong in the database the whole system will not work resulting in bad experience. Even maintenance or migration on this machine will also create a downtime in the application.

  • Bottlenecks: Even if the CPU and RAM are very high in the system then also the database instance is one due to which there will be only one WAL (write ahead logs) and disk.

Advantages: Higher Consistency, No distributed system maintenance, Only one server to monitor, Less engineering cost. Though for most early staged startups vertical scaling is very good option.

Horizontal Scaling of Databases

This is part which we we will be discussing in detail because the real scaling and distributed systems are made in this area because of the serious limitations of the vertical scaling.

So what does 'horizontal' means here ? Unlike the vertical scaling here we add more machines in the infra rather upgrading the same machine. So here we will be using several machines with low specs, all working together and serving the users. In the horizontal scaling of the databases there is mainly two most prominent techniques which are used that are:

Sharding

It is the technique where the the data is splited across multiple database servers so that one server do not gets overwhelmed by handling all the requests. Each database servers(shards) gets their own CPU, RAM, Disk Storage and bandwidth.
If one has understood the Partitioning of the databases which happens in the same server then understanding this will not be big issue. So just like there we have partitions created according to the condition here we create shards on the basis of conditions.

Like the partitioning we have different methods of doing the sharding like:

  1. Range Based Sharding: Like previously mentioned splitting the data using the range of date. These values are called shard keys which are used to split the data.

  2. Hash Based Sharding: Using some function to evenly distributing the load among all the shards.

  3. Geo Based Sharding: On the basis of country databases are sharded.

  4. Directory Based Sharding: Here we create a lookup table which stores which range data is stored in which shard like:

User 1001 → Shard 3
User 1002 → Shard 1

Best Practices:

  • Choosing Stable Shard Key: So shard key play a important role because it makes sure none of the shard is getting overloaded with the queries otherwise the whole purpose of sharding gets destroyed.

  • Keeping related data in same shard: Keeping the related data in the same shard makes sure that cross joining of tables can be avoided.

  • Introducing Consistent Hashing: It is a advanced way to distribute the load among all the shards

Disadvantages: Cross shard queries can cause latencies, Application layer complexity is added for routing logic, Maintaining consistency in these systems is often very difficult, Monitoring becomes hard.

Some Scenarios:

  • Let's take scenario where we have a e-commerce application with three shards which is created using geo based sharding. Like in diagram below:

    So it might some sale happens in India, sometimes sale happens in US. So it might happen our shard during these sales start getting overwhelmed. Here we have the option to scale again vertically or horizontally. When these shards are horizontally scaled then it is done using replication which we are going to discuss next.

Replication

Now we will discuss the last portion of the horizontal scaling i.e. Replication. What is replication mean? copying(duplicating) data from one database server to one or more other servers. This topic is the trickiest part in the whole database infrastructure because of consistency issues. We will be discussing some the issue but we will not go deep into it.

So as we have discussed it might happen a particular shard can get overwhelmed due increase in traffic. So here we have two options - Either upgrade the shard machine which is expensive to do and on the other hand we have option to horizontally scale it which is done using replication.

Master Slave Architecture:

As in the above image we have seen there is a master node which will be responsible for all the types of writes which is going to happen in the application. This node is often known as Primary or Master node. On the other hand we have three slaves which will be used for read only purpose and will be in sync with the master node.
So in this architecture we can see we can make system resilient:

  • Against any type data loss which might happens in any of the nodes as we have many replicas.

  • Read load are well distributed among several replicas which protects the master from getting overwhelmed.

  • All this collectively makes the system highly available with less downtime.

How the sync happens under the hood ?

The affected rows' copy by the query is first copied from the main file to the memory(RAM) and there it is modified according to the query. Now WAL(Write Ahead Logs) is written and appended to the WAL file sequentially. After the actual database forces it change the main file. In this way the WAL helps is creating durable data so that if something goes while writing on the disk then the whole flow can be rolled back easily.

This logs are then streamed to the replicas where these replicas use this logs to write in its WAL file and replays it in order to make those changes in its main file. In all this a LSN(Log Sequence Number) is used in order to maintain the consistency of WAL between the master and slave nodes.

These Write Ahead Logs in MySql is called binary logs.

There are mainly three types replication:

  1. Synchronous Replication: Here the master node waits until the slaves node confirms that all the changes are successfully saved in them or not. This works synchronously so there is a latency while syncing but is highly consistent.
    Used mainly in banking system where no data loss can be entertained.

  2. Asynchronous Replication: This is most common way to replicate data where the data is first committed in the master node after that it is streamed to replica which increases the throughput of the system but can cause replica lags and in some cases data losses if the master node crashes before sync. So we have to make a trade off here in order to use it.

  3. Semi-Synchronous: Here the master node waits until anyone replica confirms that all the LSN are synced.

Advantages:

  • Reads are highly scalable in this architecture.

  • High availability because of more numbers of replicas.

  • Backup safety is there which reduces the chances of data loss.

Disadvantages:

  • Master node becomes the single point of failure for write operations. This can be solved by introducing other architectures which makes the overall infra very complex.

  • Consistency is the issue due to replica lag and master node failure.

  • Monitoring becomes difficult for such a complex system.

  • Lazy population and eager population

  • Different types of indexes