A Detailed Five Step Twitter Scaling Plan
whydoeseverythingsuck.com
whydoeseverythingsuck.com
Your intuition is excellent. These are all valid techniques to alleviate disk seek limitations. They are likely to be successful. Unfortunately, they're expensive options.
You can partially overcome disk seek limitations by performing database batch inserts. So, rather than inserting messages indiviually and updating indexes for each message, you insert 100 messages and perform less disk seeks updating the indexes for each batch. Unfortunately, it can be hard to reliably merge data into suitable batches.
So, you've got two non-linear curves, cost and benefit. If you've got an exceptional team, a great architecture and good finances then getting 99.999% uptime should be easy. Of course, this assumes that you don't have unforeseen circumstances, glaring omissions in redundancy - or exponential growth and utilisation without a revenue stream.
Even with generous financing, exponential growth can be deadly. A small but sustained increment in load can be enough to cause a backlog in requests that never recovers. If you're growing fast then it can occur anywhere in your system at any time. I've been a DBA for a renderfarm with thousands of cores and this situation is stressing.
We've also had this situation while working on search. Thankfully, it occurred on a smaller scale and most thankfully it occurred before launch. We had two database servers, two app servers, no redundancy, on unreliable hardware. We didn't make much headway with efficiency improvements because long tests were likely to encounter hardware failure. It was demoralising.
This slowed the development and testing of a database failover feature in the database wrapper. Now this has been implemented, progress is much better. A test that previously took three hours now takes five minutes and we've had much more time to migrate to reliable hardware and introduce redundancy. We've also got a setup which can be used to test resilience to real modes of failure.
1) It's not just a DB problem.
So not all scaling problems are database problems.
This is an argument in semantics.
You suggest storing all users physical addresses in memory across multiple machines. You also advocate "shard splitting" a la B-tree indexes which split when the node's population reaches 2n. OK, so when a split occurs, how do you maintain that massive table in memory. Do you lock it? Do you examine your existing queue of requests (which may soon point to the wrong node)? How do you handle new queue requests during "shard split maintenance"?
Calling these "database problems" is waving your hands. Why not just call them "computer problems"?
Regarding when shard split how you maintain the table, the easy answer is Terracotta. It is amazing for this and handles all that kind of stuff totally in the background. If you didnt have terracotta it would be a harder problem indeed. It handles locking and all of that magic for you. I highly recommend it. Regarding how you handle queue requests while the split is going on, what happens is that while the split is happening, the original shard continues to respond, but all new updates are sent to the old shard and the new one that is being created in the split, so during a split there is addition messaging going on. Once its done, the original shards range of users is cut back, but actually it is still storing all the old data. At that point you can delete the data that has been moved to the new shard - though this is not necessary.
Right. Just because we didn't write it doesn't mean it's magic.
There is still a lot of work to be done. It can get very complicated for high velocity apps. And it takes time.
So you're right. Bigger minds than the two of us are probably struggling with this (and lots of other stuff) as we write.
Now as you scale to 50 or so servers they will start to fail with some regularity and you can't reboot and or rebuild in five seconds so we need to have redundant systems ready. Splitting a shard probably takes longer than five seconds so you need to be proactive about such things. So we need an active rebalancing between shard so they can quickly pass around their load and fast failover incase things fail. These are DB problems but they are only problems because they system need to keep responding quickly as parts of the DB fail.
But this is still DB centric how do you load balance incoming and outgoing requests? (Routing each existing user to the same system every time and letting each system lock data when they want to change their password etc.) How do we recover from data corruption? How do we notice an outgoing or incoming mail server is down? You also have n^2 number of connections between systems as you grow can we keep that many sockets open? (More of an issue is we are segregating the network.) Etc.
PS: These are not hard problems by themselves but they are linked. And there are a lot more where they came from.
For example, I load tested my last rails app (load testing is my background-I have consulted for major pharma, tax, and other companies) and on a wimpy inherited server, processor maxed out first for the rails processes. I popped in a second CPU, then was able to get to the point where we could support 200 concurrent users, which was all we needed at the time.
I have seen limitations on bandwidth, firewall problems, load balancer misconfiguration, poorly written queries, memory usage,processor usage, thread deadlocks, and licensing limits (max database connections on a trial license) all be bottlenecks. At the larger scale, yes, database is more often the problem, but it's not clear twitter should be using a database in the way they are, anyway. Your statement goes too far.
Now consider a voting system you have 1,000 items and 1,000 users and you want to show how much people like them like each item? How well does that work when there are 100,000 users and 100,000 items?
If the average number of recipients grows along with the number of users, then there's a problem, because average number of messages should be growing linearly with userbase, so the total hardware required is quadratic. No way that would work out economically...
Basically you end up with this equation
cost = ((NumOfUsers * avgMsgsRevceivedPerUser)/maxWritesPerCPU) * costOfCPU
There is nothing quadratic here. Your are able to scale linearly with the number of messages received in the system. There is a fixed cost per message received which is exactly what you want. By the way, this is exactly how the internet mail system works and it is why it is so scalable.
E-mail's actually a really great example. It worked great as long as it was used for person-to-person communication, which held both the number of messages/user and the number of recipients/message constant. Then spammers found out that the more people who had e-mail, the more they could send to at once, and sent recipients/message through the roof. In 1997/98, there were actually a number of articles about how the e-mail architecture was in serious trouble, and unless some way could be found to limit spam, it would collapse within a couple years. Of course, they dealt with it through spam filters and blacklists on the server, so now e-mail works fine.
Twitter's problem is that both Messages/User and Recipients/Message average significantly higher than e-mail, and yet the resources they have available are significantly lower than the full resources of the Internet. So even if they don't have an architectural scalability problem, they could very well have an economics problem.
That doesn't make your opinion of what's wrong with Twitter any more valid than the crap TC keeps printing unless you happen to also have lots of internal details of the Twitter architecture.
This doesn't even get into the complexity of scaling your model from 1,000 or 10,000 users where your brute-force methods seem like they will work and applying it to 10,000,000 users where your model will die a horrible death, kinda like Twitter is today.
As I read your design, one of the biggest problems is that to build up the infrastructure I think (and infer based of my experience with designing these kinds of systems) you are describing is that it's a massive engineering effort in and of itself. You basically end up rewriting a realtime-optimized BigTable or HBase, and that's much to much to expect of Twitter's small engineering team.
People seem to forget that scalability is a tradeoff, like any other aspect of a service. You devote an amount of resources to the problem that makes sense to your business. Despite twitter's downtimes, people don't seem to be leaving the service, despite several competitors. The service is down, but it isn't down so much that it's unusable.
For a site that's grown at such a rapid pace and has such a small engineering team, I think Twitter's scaling plan is probably very reasonable. They're growing their capacity at a rate they can afford, and a rate that they know will retain most of their userbase.
It is certainly reasonable to suggest that Twitter is good enough as it is because people still use it, though the grumbling in the market suggests that that may not be something they can bank on forever.
And have you checked latency serving out of EC2? Our experience is that the connections are fairly fast but the latency establishing the connection is significant, especially from the west coast where twitter is based.
Regarding EC2, latency isnt really an issue because Twitter is far from real-time. That said, haven't really noticed any problems with EC2 connection latency, but we are not live yet so I would defer to you on that.
Better safe the messages in a central store by id and only write the ids to the different pages.
Why? They are proposals to solve a problem that has not yet been defined.
Fact is, you don't know what the problem is.
So proposing a database solution to what could just as easily be a memory management or throughput problem is premature.
I just don't like gloating on someone else's mistakes and failures (essentially kicking someone who's already on the ground), especially when I'll probably have my own set of spectacular f-ups to deal with related to my own startup
I am currently managing a system that serves over 100,000 DNS queries per second, each of which is dynamically evaluated as to how we respond to it. And each of which is logged, amounting to over 7,000,000,000 log lines a day, all which get processed, parsed and in many cases handed over to a database store. Yep, we're dealing with BILLIONS of rows of MySQL per day.
Hey thanks for OpenDNS by the way, I'm a long time user, and I love it.
also you forgot about messaging queues.
This isn't even about Ruby, it's about the Rails architecture.