Companies/Twitter/

Manhattan — Twitter Distributed Database

Lesson overview

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

Loading Manhattan — Twitter Distributed Database