Manhattan — Twitter Distributed Database
How Twitter built Manhattan, its distributed storage system — sharding, replication, replica repair, and why existing databases did not meet Twitter's latency and operational requirements.
Manhattan - Twitter Distributed Database
Manhattan was built by Twitter’s Core Storage team. Boaz Avital was one of its original engineers.
UserId: 42 tweetId: 900 Text: "Hellow world"
UserId: 42 commentId: 900 Text: "Hellow world"
UserId: 42 tweetId: 902 Text: "Hellow world 2"
Comment table
for some cases: 1. Availabity > consistency - tweet 2. Consisteny > Availability - Username
Functional requirements Store and read key-value data. Divide data across many shards. Replicate data across servers and datacenters. Repair missing or outdated replicas. Support strong operations such as compare-and-set. Non-functional requirements Low latency. High availability. Horizontal scalability. Reliable performance during failures. Isolation between different Twitter applications.
Solution : 1 Sharded MySQL
Good Familiar database Transactions Strong consistency Problem Applications must understand shards. Moving data between shards is difficult. Operating many clusters becomes difficult.
Solution : 2 Distributed Database
Use something like Cassandra or another Dynamo-style database.
Already supports horizontal scaling
Already supports replication
High availability
Twitter said existing products did not fully meet its latency and operational requirements. Twitter also wanted more control over storage engines, multi-tenancy and internal tooling
Solution : 3 Different database for every workload
Tweets -> Database A
Profiles -> Database B
Recommendations -> Database C
Counters -> Database D
Each workload gets the best database.
More systems to operate
More monitoring tools