5 ms·
When I've planned out sharded infrastructures, the database usually wasn't that big of a concern. The web framework or system architecture are usually the pain
by subwindow 18y ago
When I've planned out sharded infrastructures, the database usually wasn't that big of a concern. The web framework or system architecture are usually the pain point. And with Rails, you have to abuse establish_connection if you're going to have each web head read from multiple shards.
The easiest route I've gone when setting up a sharded infrastructure is to use subdomains and a 1-1 Web:DB setup. Have each subdomain go (either thru a reverse proxy or hardware load balancer) to a different (sharded) webhead. Each web head talks to two databases- the common database and its sharded DB. With this you'll probably want a "common" web head to handle home page traffic and authentication (after they are authenticated they'd get dished off to their shard).
Salesforce.com was my inspiration for this method, and it is probably reasonably common. It probably also has a name, but I do not know what it is.
- sanj 18y agoI've seen folks approach things this way -- Dr. Nic's done some of it in an exploratory manner: http://tinyurl.com/36twmo http://tinyurl.com/36twmo But I don't think it is the best approach. I'd much rather push ALL of the magic down into the database layer and not have the app worry about it at all. I want ONE layer of magic and I'd prefer it to be in the DB, where it seems it "should" be.
- cstejerean 18y agoI'd rather have the layer of magic be in the part of the code I understand best so I can fix it when things go wrong.
- sanj 18y agoI'm pretty confident I can learn what I need about DB configs, so I'm more concerned about it being in the "right" place from a complexity and efficiency standpoint.
- j2d2 18y agoWhy doesn't anyone ever talk about distributed caches? Is it really necessary to hit the database for all data? Perhaps I'm confused...
- mihasya 18y agoI disagree with you on many points there. I'm not sure what you work on, but the bottleneck in my experience is almost always the database. Doing 1:1 web to db setup is not a good idea because more than likely you'll end up under utilizing the web servers as they will outpace the db (unless your db has one three-column table in it of course). Even if it's the other way around, you're still handcuffing whichever server is faster/under less load. Every time, say, a DB server is overloaded, instead of just adding another DB server, you have to add TWO servers. Also, this ends up being sort of sloppy b/c the subdomain for your website changes by user once they've logged in. What happens when one user pastes a link to something he saw to the other user? You have to figure out how to handle that. To the OP: I have recently been looking around for similar things and, as expected, this is a very specific problem that hasn't been tackled (or certainly tackled well) by that many websites. In any case, there is no cookbook solution for it. It is probably a good thing, since your needs in this case will likely be very specific. You are most likely going to have to go it alone. I have yet to implement this in practice, but my approach would be to have a single DB master that just has the shard allocating table (make sure to cache the shit out of that, hopefully using your app server's local cache, though I don't know if RoR has something equivalent to PHP's APC, so you don't have to do a lookup on every call) followed by the additional masters. Figure out what you're going to shard by (user id etc). Then you will want to write some sort of algorithm that distributes new content between the existing servers (this will depend on how much content there is per-user; if you are sharing something that is database intense and each user has many of it, tying users to a single shard could lead to that shard getting screwed by a few super active users, so you may want to shard on a more finegrained level than that). At first, you will probably just have to assign new shards based on how many shards are on the servers you currently have (as you add an empty server, all new shards go to that), but I would recommend overtime implementing something that uses actual db server load or some other statistic that is actually more telling of how much work each db server is doing to distribute. As your DB grows (in terms of adding tables), I'd make sure to keep a script around that makes moving a shard easy. Basically something where you can put up a "Your account is under maintenance" for one user and then just run a script that updates the central allocation database and moves all the data associated with that shard accordingly. I think this goes without saying, but replicate every master both for reads and for failures. You also want memcached in there somewhere, I would assume, but that's another story. As I said, I have yet to actually start implementing this, but this has several advantages. Database servers are easy to add without manipulating your app (especially if your sharding algorithm is good), you can shard according to load, you aren't handcuffing your app servers to a single database server, etc. Hope this helps. Cheers. Edit: you might find some things that will help here: http://en.oreilly.com/mysql2008/public/schedule/proceedings http://en.oreilly.com/mysql2008/public/schedule/proceedings I can't remember which talks they were, but I saw quite a few people show layouts of approximately how their architecture worked. Not all of them were great, but it might give you some ideas. Look through the memcached slides for sure.