System Design of Twitter.
System Design of Twitter.
Requirement Gathering
- Functional Requirement:
Follow Others
Create Tweets
View Feed
- Non Functional Requirements
Highly available
Highly responsive
Pretty UI/UX
Let’s analyze the functional requirements of twitter.
- Create Tweets:
In twitter there are about 500M users and 200M daily active users. To start with we are assuming that each tweet has a character limit of 140 so it would be 140 Bytes. Since we have not considered the avatar, the username, the identifier hashtag, it could be around 1KB. Let’s Consider the Tweet also contains the Video or Image within itself so the estimated data for each tweet is about 1MB.
Out of the 200M users, let’s assume that about 50M users create tweets on a daily basis.
So, the storage estimate here is:
50M*1MB = 50MB*1000*1000 = 50B *1000 = 50T
- View Feed:
Let’s Assume that each active user reads about 100 Tweets per say, the overall tweets read is 20B. I.e for
200M daily active users view 100 tweets of 1 MB each. 200M*100*1MB = 20B*1MB.
Now lets calculate the overall throughput required here.
20M*1MB= 20B*1000*1000 = 20T *1000 = 20PB bytes
- Follow Others:
A user must be able to follow other users. Let’s assume that out of 200M active users, 10M users have followed each other then the relation here is 10*10M =100M connections. In this case however there can be users who could have over 100M connections/followers. This is a case to be considered.
High Level System Design
We will be using the C4 model to make the high level design. In the C4 Model for documentation, for the given software system we use containers and components to describe the software system.
The C4 stands for
System Context
Containers
Components
Code
The System Context Diagram.
We are using the C4 model for creating the context level diagram of the system.
For the Level 1 diagram we point out what basics needs to be done by the system.
Fig: Level 1 Context Diagram
For Level 2 diagram we will be expanding the twitter system and trying to figure out what the system would look like. The Level 2 diagram is subject to change as we might find new findings in Level 3.
In the Below diagram we have identified that the user will post tweet, view feed and follow users. This all is handled by the api service which in turn acts as a load balancing mechanism. The load balancing mechanism is added in order to distribute the load amongst multiple application servers.
The application servers handles saving of tweet, fetch feed and establish follower, followee relationship.
We have tried to separate the database into multiple master slave configurations in order to handle throughput. Some of the slaves acts as read only replicas which reduces the load of write which can occur during table lock.
Separate Document storage or Object holders mechanisms are in place in order to handle file delivery. Since twitter post might contain videos, photos, etc. It is wise to save it separately and distribute it through CDN. The CDN could be inhouse or a SAAS. The requests for file can independently go through CDN without ever touching twitter system which impacts the throughput of twitter.
Fig: Level 2 Diagram
Before Jumping into conclusion for the level 2 diagram we would like to revisit how the actual system handles such high throughput.
Before adding the master slave configuration of single or multiple read write replicas we need to find out how we can have max throughput with single DB.
Fig: Level 3 Diagram
For max throughput through single DB, we have planned multiple caches with different configurations, pub/sub model for async operations database partitioning to aggregate similar things altogether.
Our assumptions:
**Cache:
**There are 2 distributed caches in this configuration.
The first distributed cache holds the latest tweets along with its metadata and frequently requested tweets in picture. We might limit the number of tweets this cache can hold as this might not get called that frequently.
The 2nd distributed cache is created in such a way that it works in a fan out model. In the fan out model we need to be able to create a pipeline of feed timelines for each user. That is if a user posts something then for each follower the tweet content will be made available to the timeline. We know that this is a very write intensive task but this will help to fetch feed in an instant. Since new tweets written by a followee will get appended to the timeline of followers and we can notify them about new tweets via pub/sub.Pub/Sub:
As mentioned, it is used to notify users about new posts from followee, this is very important as it simulates a real time system and does not miss any tweets. For users to have new posts notification it needs to be first written into the fanout timeline and then needs to be notified. 2 separate topics could send out data, one sending that new tweets on feed which forces the frontend team to refetch if the user wishes to engage. On the other hand we might send notification regarding a followee posting new content with the gist of the tweet. If the user clicks on this notification then we load the tweet from redis cluster which holds new and frequent tweets.**Database Partitioning:
**We are partitioning the db in order to aggregate tweets from a user within a block so that we can fetch tweets faster, as it is going to be a sequential fetch.
Here we have changed certain aspects of the system. We would like to see how the internal of the system would work before we jump into conclusion on how level 2 would look like.
After including cache, pub/sub model and partitioning we feel that the database here has become a bottleneck. Let’s use a distributed db with sharded contents to reduce load on a single db. Taking this into consideration now we do not have to do partitioning but the shard key needs to be wisely created. We need to be able to retrieve the latest tweets of a user faster than past posts.
Fig: Actual Level 2 Diagram of twitter
Reference on how it's done. here