Manhattan: Real-time, multi-tenant distributed database for Twitter scale
blog.twitter.com
blog.twitter.com
Summingbird (bear with me, I'll tie this in) is also twitter's answer for writing code once and seeing it run on a variety of execution platforms such as hadoop, storm, spark, akka, etc... Not all of these have been built out, but the platform was designed to be a generic framework to support write once execute everywhere.
Summingbird is written to support Manhattan's model as well. The high level idea is to use versioning to determine whether a request is precomputed (batch), computed (realtime) or a hybrid (precomputed + computed). These are expressed as monads with basic functionality present in algebird. One way to bring support to this model to the open source world would be to implement storehaus bindings for elephantdb and to resurrect elephantdb or build a similar service to provide storage similar to Manhattan.
Overall, very early yet promising work in the open source community.
[edit: book is not about elephantdb, but is a critical component. modified wording. Also added link]
Over time, Manhattan evolved into a fully fledged read/write database that is able to support batch + read/write in the same system. Batch is great for some use cases, but sometimes the cost is too expensive when you factor in how much processing power, storage, etc is needed for that model. We support both of course still we want developers to have the freedom to increase their productivity.
They have got the hadoop watcher right as well as the BTree/SSTable update mechanisms - heavy import, light update is a use-case that is not catered to usually.
The rest of it feels a lot like memcache (membase/zbase impls) when reading up on its architecture - plus SSDs, that's always going to beats the pants off anything with disks.
Looks good overall, this is probably not for you if want to do exact counters or match up different counters (i.e +1,+1,+1 for a counter funnel).
If you don't need updates to be consistent, like if you imported hadoop data in without modification, this starts to look like a really good model.
Not that I took it vary far, but my first stab to see what the scaling issues would be was actually plenty fast to run there feed process at the time.
PS: The 'trick' is to keep two lists one everything a user follows, and another is everything that follows each account. For showing the messages to someone when they log in you keep the last 10 message Id's with timestamp or sequence ID so if someone follows 5k accounts you can avoid looking at the vast majority of those messages then sign them up the device to listen to new messages. (Sure, sometimes you will need to look past that top 10 but it's rather effective.)
(And for the downvoters I can upload some code if you want to see it.)
PS: I did not keep old messages just there ID because that was not going to fit in RAM. My assumption was using Redis or other key value store would be fine what they needed was an internal index so you would only need to look up messages that would be displayed.
Note: there current setup once they worked the bugs out handled a peak of over 100,000 tweets per second in 2013. https://blog.twitter.com/2013/new-tweets-per-second-record-a... Which is well beyond the target I was shooting for.
That's the difference between "getting the basic process running" and operating at scale.
Back when they where growing from 500,000 users to 7 million total users they where having major issues and that's when I was looking into things.
Anyway, not I was suggesting 1GB would be fine today. Still, I was saturating a 1GB connection so 15M (over twice the total users back then) * ~200bytes * 8 / ~1000^3 = ~24 seconds did you have an extra 60x multiplier in there somewhere?
Things add up eventually. Suddenly your 6000 tweets/s has turned into a million op/s fanout.
Problems become a lot less straightforward over network links too.
By comparison, many complex machine-generated data sources (e.g. real-time entity tracking) that are sometimes fused with the Twitter firehose operate at millions of complex records every second (often tens of gigabytes per second) that need to be processed, indexed, and analyzed in real-time. You can't deal with this kind of data model using something like Twitter's current architecture because the several order of magnitude difference in velocity and volume exposes the limitations of most database platform designs people typically use.
A $DATABASE_COMPANY was recently looking to hire somebody to be a dedicated database breaker. The Jepsen series was mentioned in the job post.
Worshipping celebrity is unhealthy in what should be an engineering profession and is a sign of a deeply embedded pop culture.
Don't be lazy.
"When you open the Twitter app on your smartphone and all those tweets, links, icons, photos, and videos materialize in front of you, they’re not coming from one place. They’re coming from thousands of places."