7 ms·
Lessons Learned While Building Reddit to 270 Million Page Views a Month
- arkitaip 14y agoI love reading about how companies scale their BigHuge data but it bothers me that we still haven't reached the point where scalability is a commodity instead of a patchwork of technology that everyone actor solves in their own way.
- joering2 14y agoThat's because 99% of internet websites will never see 200 million requests.
- jeremyjh 14y agoMuch fewer than 1% will.
- NelsonMinar 14y agoScalability can only be a commodity if all big sites were built the same way. But they never are. A scalable read-only system is entirely different from a system with a mix of read/write which is different again from a message routing system like Twitter or Facebook.
- batista 14y agoThat doesn't mean they cant take advantage of ready-made "functional architectural components" for each of their needs. So they can pick and match from the "scalable read only" commodity system, the scalable "message routing system" etc.
- nostrademons 14y agoIt's because "how to make a website scale" depends heavily upon which website, what it does, and how big it needs to be. Making a messaging queue like Twitter or G+ scale is very different from making Google Search scale. Hell, making the indexing system of Google search scale is a very different problem from making the serving system scale. You can't really avoid having a patchwork of technology, because it's a patchwork of problems. Instead, there're a bunch of tools at your disposal, a few "best practices" which are highly contextual, and you have to use your judgment and knowledge of the problem domain to put them together.
- batista 14y ago>It's because "how to make a website scale" depends heavily upon which website, what it does, and how big it needs to be. Making a messaging queue like Twitter or G+ scale is very different from making Google Search scale. Hell, making the indexing system of Google search scale is a very different problem from making the serving system scale. For most websites is not THAT different. Actually, most have pretty similar needs, and you can sum those up in 3-5 different website architectural styles anyway. There is far more duplication of work and ad-hoc solutions to the SAME problems than are "heavily different" needs.
- nostrademons 14y agoWhat would be those 3-5 different website architectural styles?
- batista 14y agoNews/Magazine/Portal like (read heavy), Game site (evented, concurrent users, game engine computations), Social Platform (read-write heavy), etc. Most needs are bog standard. If you really look at most successful sites they use might same-ish architectures, only with different components/languages/libs each. Basically all high volume sites use something like the notions behind the Google App Engine, and the services offered. The various AWS tools are also similar (S3, the table they offer, etc).
- nostrademons 14y agoI think you're missing a lot of complexity of the considerations that actually go into implementing any of the above. I can think of 3 subsystems within Reddit alone (reading, voting, and messages) that all have different usage patterns and (if they're doing it right) require different approaches to scaling. Where's e-commerce on your list? The approaches for scaling eBay are completely different than for scaling Reddit or YouTube, because eBay can't afford eventual consistency. You can't rely on caching to show a buyer a page whose price is an hour out-of-date. Here's something else to think about: why do (the now-defunct) Google real-time search and GMail Chat have completely different architectures, despite both of them having the same basic structure of "a message comes in, and gets displayed in a scrolling window on the screen"? The answer is latency. With real-time search, a latency of 30 seconds is acceptable, since you aren't going to know when the tweet was posted in the first place. With GChat, it has to be immediate, because it's frequently used in the context of someone verbally saying "I'll ping you this link" and it's kinda embarrassing if the link doesn't arrive for 30 seconds. Real-time search also has to run much more computationally-intensive algorithms to determine relevance & ranking than GChat does. I've personally worked on Google Search, Google+, and Google Fiber. I can tell you that they do not all use something like the notions behind Google AppEngine. There's no way you could build Google Search on AppEngine, and G+ would be a stretch.
- ksec 14y agoWell because every site are made different, and uses different tech, with different bottleneck. But with SSD, RAM becoming a commodity, ( both prices has been dropping shapely ) I/O are much easier to deal with. And with every release of PostgreSQL / MySQL, Easier Database Replication and more common practice of scaling we are much much better at it then say 2 - 3 years ago.
- ndemoor 14y ago"Instead, they keep a Thing Table and a Data Table. Everything in Reddit is a Thing: users, links, comments, subreddits, awards, etc. Things keep common attribute like up/down votes, a type, and creation date. The Data table has three columns: thing id, key, value." I hope they introduced some NoSQL sweetness by now.
- encoderer 14y agoExactly. They turned their RDBMS into a NoSQL database that's still slowed-down by all the relational machinery. Though I'm sure this was an informed choice at the time, I seriously hope nobody finds this advice actionable anymore. Use Cassandra. (Or your NoSQL DB of choice)
- rfurmani 14y agoThey in an ad hoc way turned postgres into a key-value store, but in reality they have Cassandra running as a "permacache" in front of it and, in practice, everything really just hits Cassandra. It looks like postgres could at some point be phased out. Source: ive hacked at the codebase to produce arxaliv
- jebblue 14y agoThe only problem with NoSQL is what happens when one day _you_need_to_relate_data?
- encoderer 14y agoIt's an all-of-the-above strategy. I remember when the Bigtable paper was released. It was very early in my career and I remember it sounding so alien to me. Sure, i had Memcached in my stack, but no SQL? It seemed like something they had to trade off to be able to build the kind of services they offer. I felt the same way after reading Dynamo. Sure, I thought a lot about data design. I thought about usage patterns to inform how we denormalize. And I grew into using, eg, Gearman, to pre-compute dozens of tables every night. I evolved, a bit. But a few years ago, a little before this OP was written, I had a great experience with some Facebook engineers and had an a-ha moment that has made me a much better software engineer. Basically, I realized that I needed to let my data be itself. If I have inherintly relational data, then it should be in a relational database. But I've built EAVs, queues, heaps, lists, all of these on top of MySQL and Postgres. Let that data be itself. We have more options now than ever before. K/V stores, Column stores, etc. I use a lot of Cassandra. A lot of Redis. Some Mongo. And a put a lot more in flat files than I ever thought I would. I know a lot of people that are smarter than me left the womb knowing these things. But for me it was transformational and has made me much happier. I realized how much energy I wasted fighting my own tools.
- citricsquid 14y agoArticle is from 2010, if I remember correctly their architecture has changed substantially since this article.
- chaz 14y agoTheir most recent infrastructure blog post was from January 2012, and shows that they're using Postgres 9, Cassandra 0.8, and local disk only (no more EBS). I'm curious if the recently-announced provisioned IOPS would enable them to go back to EBS. http://blog.reddit.com/2012/01/january-2012-state-of-servers.html http://blog.reddit.com/2012/01/january-2012-state-of-servers...
- mdellabitta 14y agoProbably the better move would be to go to SSD-backed High I/O instances. Netflix did: http://techblog.netflix.com/2012/07/benchmarking-high-performance-io-with.html http://techblog.netflix.com/2012/07/benchmarking-high-perfor...
- ecaron 14y agoOnly disagreement (although I feel like I'm arguing w/ Linus about git) is don't memcache session data (lesson 5.) Memcache's 1mb max-block (exceeding that removes too many performance perks to be considered viable) introduces a "I need to constantly worry about my sessions getting too big" mental overhead that isn't worth it. Go with Redis for storing session data.
- batista 14y agoWhy would you need 1mb session data? Or 200K for that matter?
- ecaron 14y agoBecause it happens. Typically in places where it shouldn't (e.g. bad programming) but sometimes it does (e.g. enormously complex active system state.) The point is that there are solutions offering equivalent speed without that barrier, so why select a technology that has limitations?
- batista 14y ago"Because it happens" is not a very good technical answer. Anything can happen, even 1GB session data. But even anything approaching even half a MB of session data I would take as an architectural failure, and work to fix it, instead of letting it dictate what tools I use (ie. redis vs memcached). >The point is that there are solutions offering equivalent speed without that barrier, so why select a technology that has limitations? Because all decisions have tradeoffs and it more wise to depend on your particular use case to choose than to depend on a scenario (> 1MB session data) that can only happen with a seriously fucked up application design.
- disordinary 14y agoYou can easily get to 1mb if you follow the cache everything mantra that reddit is using here.
- ecaron 14y ago> The Google crawler will hit you has fast as you let it, so when gets slow just crank up the rate limiter and it quiets the system down without hurting users. How does one tell Google crawler to slow down without potentially hurting SEO? (Source: http://googlewebmastercentral.blogspot.com/2010/04/using-site-speed-in-web-search-ranking.html http://googlewebmastercentral.blogspot.com/2010/04/using-sit...)
- surferbayarea 14y agothat's just ~12 queries/sec..not that huge a deal..
- jonknee 14y agoThere are 2,592,000 seconds in a 30 month day so actually that's more like 104 page views per second. But that also does not count peaks (peak is likely 250+) and that each page view requires multiple queries. Either way, no need to hate on the traffic figures of what everyone knows is a very large website.
- masklinn 14y agoAlso, it was back in 2009[0]. In December 2011, they were up to 2.07bn page views (~750 pages/sec on average) [0] according to http://blog.reddit.com/2012/01/january-2012-state-of-servers.html http://blog.reddit.com/2012/01/january-2012-state-of-servers..., December 2010 saw 829 million page views
- ralfd 14y agoInteresting is also traffic peaks. The IAMA from President Obama a few days ago really stress tested Reddit: http://blog.reddit.com/2012/08/potus-iama-stats.html http://blog.reddit.com/2012/08/potus-iama-stats.html > At the peak of the IAMA reddit was receiving over 100,000 pageviews per minute. > In preparation for the IAMA, we initially added 30 dedicated servers (20%~ increase) just for the comment thread. This turned out not to be enough, so we added another 30 dedicated servers to the mix. At peak, we were transferring 48 MB per second of reddit to the internet. This much traffic overwhelmed our load balancers which caused a lot of the slowness you probably experienced on reddit.
- deleted 14y ago[deleted]
- deleted 14y ago[deleted]
- deleted 14y ago[deleted]
- mthreat 14y agoI'd love to see a 2012 version of this article, with SSDs as a viable option.
- petercooper 14y agoIt's pretty cool how nowadays Redis can cover several of those concerns quickly and simply (open schema, caching, replication..)
- dotborg 14y agoHow are PHP, RoR, node.js, Perl/CGI etc. engineers supposed to build offline processing? Crontab?
- hythloday 14y agoThere are worker queues for at least Perl & Ruby. I have no idea about PHP and JS but one assumes solutions exist.
- CWIZO 14y agoWe use gearman (and PHP, but it works with many other languages). It's great.
- benologist 14y agoNodeJS is insanely easy to do that with, it's one of the best reasons to use it in my experience: var queue = []; // request handling module.exports = function(request, response) { queue.push("bla"); response.end("bye"); }; // out-of-request work setInterval(function() { // do something with your queue either locally // or moving it to an external job queue that can // leverage the same script on a different thread }, 1000);
- dotborg 14y agowhat about concurrency?