Why are Facebook, Digg and Twitter So Hard To Scale?
highscalability.com
highscalability.com
Maybe things have changed since the last time I visited Digg over 2 years ago, but the social networking aspects are not very significant. The vast majority of their hits are practically fully page cacheable.
Twitter at least has an interesting scaling problem, but they don't have any features and they move at a glacial pace.
Facebook on the other hand has a graph that almost as nasty as Twitter's (minus the million followers thing), but they have 100 times the features, and they push new code every week.
But you're right with Digg, I don't see it having any excuses
The internal representation is quite elegant.
The killer app of Facebook - poking - does need to be reliable ;-)
Poking, while needing to be reliable is a 1:1 activity, so light on writes and easy to read, whereas broadcasting is difficult because of the wide variance in update frequency and number of followers from user to user.
I actually think having a bunch of large sites tackling similar issues will present a whole bunch of neat new technologies for addressing these sorts of issues in the near future.
All these sites have scaled to multiple servers, which means they've all had to address some of the hard problems, and no doubt that many of the issues are very similar. However even if Facebook's traffic were half of Twitter's it would still be a much more impressive architecture. They support a full-on integrated development platform for crying out loud. That is orders of magnitude more impressive than Twitter's lightweight API whether they're serving 100 or 100,000 requests per second.
I'm not dissing Twitter either, it's just that what Facebook has going is pretty incredible. Consider that failing to scale is basically what killed Friendster.
Huh?
(check out http://www.facebook.com/Engineering?v=wall )
The probability of (B fr C) is greater than the mean if (A fr B) and (A fr C). That's useful information.
Off the top of my head, I wonder how an algorithm like:
- I have N shards
- pick the top N most-connected users
- assign them each to a shard
- assign their immediate friends to the same shard
- randomly fill in other users
would work.
Possible refinements:
- if the %age of shared friends between two users in the top N is > X, put both users in the same shard + add the N+1th user at a new shard-seed
- chase more than one level of immediacy from the shard-seed users to fill the shard
- if a shard is full, don't drop to random allocation for 1st- or 2nd- level friends, but instead put them all onto shard+1
The idea here is that for pull or push you win if you need to contact fewer shards. i.e. the queries and updates needn't be per-user but per-shard. i.e. you can query/update for all users on a shard in one sql statement.
If you somehow manage to keep all of ashton's friends on 10 shards instead of 100, then that's a big win, surely?
e.g. If Oprah only publishes 3 tweets a day, but her one million followers each check 100 times a day (just to be on the bleeding edge of gossip), it's much less effort to push the change.
On the flip side, if you post status changes several times a day but your followers rarely check (daily/weekly), pull may make more sense.
Although, I suppose that a push inherently requires an update/write, while a pull is generally a read. Seems like this might need to be taken into consideration as well.
The problem in large scale information networks, more in the Facebook way than the Twitter way, is that you you potentially trigger a cascading effect in information updates if you go pub-sub. Specifically, applications that do interesting things with social graphs have to go beyond basically doing message passing.
I naturally look at things from a recommendations angle, but if you've got a new edge that enters the graph that may affect other edges that are connected to the end points. Those may in turn affect the edges that are connected to those nodes and so on. You want to avoid something that effectively becomes a breadth first traversal of the graph doing updates since that's well, slow, to put it mildly.
This is why large-scale graph algorithms like PageRank work on constantly regenerating static matrices rather than doing regeneration of the ranks dynamically, but that naturally is problematic when you're working on data sets where the most recent data is the most important and is being generated at very high rates.
I can't think of anything, really. If X and Y becomes friends, X's friends and Y's friends see it in their feed ("inbox"), but nobody else. Same goes for tagging someone in a photo, etc.
If accessed frequently, but changes less frequently, then push makes sense, if changes frequently, but accessed infrequently, then pull makes sense.
It also seems likely that the same piece of data may have different ratios from different perspectives.
If you think that's obvious, I think you're significantly beyond all of the people with all of the "hello world" twitter clones out there.
This is quite different from the problems of even large scale web apps where there's essentially a set of data that's pulled from a caching layer and assembled.
5k node updates per second might sound like a problem, but one core of one machine can easily keep up with that so you can have several copies and several views of the whole network graph. Public vs. private messages can be handled separately and then joined before presentation to the user. You can separate finding which message to display from the message data. You even get to display dirty reads as long as the data is <2 seconds old it's plenty good enough.
Ok, describing the solution based on the above insights takes some time and pictures but does it still sound horrible?
PS: Twitter was forced to morph an architecture built to solve a different problem into a working solution. That takes time and can be fairly difficult. But, starting from scratch it's not that bad.
The communication costs are significant, and probably similar to N-body simulations.
Hence the "eventually consistent" model.
BINGO! Twitter's initial implementation was a "my first blog" in Rails -- when you're starting from an impedance mismatch that massive, and you have to rearchitect it live, jesus.
However, Facebook rolled out a messaging system with little problem. I think the problem is guessing and simulating the load before people start messing with it. While not wasting millions building for load that never shows up.
It could very well be a RAM based approach might be viable once memristros (http://en.wikipedia.org/wiki/Memristor#Potential_application...) are used for memory.
As for the benefit... there was a picture some time back here on HN: http://imgur.com/X1Hi1.gif
They now use Scala/JVM rather than rails for the backend, Rails for the frontend.
http://highscalability.com/scaling-twitter-making-twitter-10...