Scaling out Postgres, Part 1: Partitioning
petrohi.me
petrohi.me
We also initially wrote the app to use a mix of Postgres and MongoDB but ended up moving everything into Postgres.
Our experience thus far in Mongo vs Postgres: go with what you know. We probably could have made Mongo work for us (especially Mongo 2.0+), but it was easier for us to work out the performance issues with Postgres. Plus, things made more sense in SQL.
If you're using Mongo and need a cloud provider, I can recommend the people at MongoLab. They really went above and beyond helping us get up and running and are really cool people.
Didn't like parallel queries using pgpool2?
I guess we could use pgpool2 parallel queries, but this would likely require some adjustments in how they are implemented.
No wonder people are moving to products like Cassandra, HBase, MongoDB, CouchDB etc which automatically and transparently shard across as many nodes as you like.
My point is that if you want to spread out your writes on a database that supports rich queries and honors the "C" in CAP, you're going to have to do some work. The distributed DBs are nice and shiny if you can forego some of the nicer querying features of relational DBs (or MongoDB), but not everyone's requirements fit into document-based records, and to assume they do is just NoSQL fanboism.
Note: I am using MongoDB for my latest venture because we are realistic about it's abilities and also because our development time with MongoDB is no joke about 1.5x-2x faster than with MySQL (or equivalent). So MongoDB has its place and it definitely a good system (especially after v2) but NoSQL isn't for everybody, especially if you can't get past the hype and look at decisions very technically and objectively.
While automatic sharding definitely has its place, we wanted to be very explicit in our consistency tradeoffs.
But hey if you wanted to be explicit about consistency then I am sure you had your reasons. Just strikes me as all a bit odd. Then again I didn't realise PostgreSQL didn't have an automatic sharding feature.
Automatically and transparently sharding across an arbitrary number of nodes is a case that is actually required by a tiny fraction of the industry. And that fraction is already paying people to monitor and control that sharding because they're too large to trust the "automatic and transparent" algorithms.
The big users of Cassandra and Hbase use them for deeper architectural reasons than simply automatic sharding.
I'd say it is closer to Oracle RAC in the space of use cases, without some of the nutty (but very impressive, and expensive) engineering to get this "shared-something" tradeoff RAC has.
Seems odd for PostgreSQL not to have made cluster features more of a priority given the more towards cloud based deployments.
Sure, there are many project goals, but they're rather informed by what various correspondents of pgsql-hackers (this could be loosely coined "the community". Admission is free.) seem personally willing, interested and/or funded to do. You, too, can be a database internals engineer! And Postgres's code is still considered crisp enough to fork for Your Very Own Database Startup (which I think is pretty remarkable -- most such stand-alone programs are typically not worth reworking), so you can waltz right on in.
One piece of project gridlock is that the bar is very, very high for any new code to be committed (and almost certainly released), and that goes double if you are not someone who is proven to suffer through all the bugs one is going to find over the foreseeable future and no one else is excited about doing that for you (if you submitted some excellently written, badly needed functionality, I think it would be accepted under that rationale, though, in spite of most complexity concerns). Postgres has a long history, so "the foreseeable future" lasts an awfully long time.
Some old features like the way hash-group and hash indexes work I am reasonably confident could not be committed if submitted as-is today. (The former can't spill to disk and the latter has no crash recovery, which is why the documentation shoos you away...but some people use it?)
Besides bugs, there also needs to be a lot of convincing that the feature is worth whatever complexity it brings, in implementation and in user-interface.
In spite of that, a lot of code has gotten committed, and not without some trepidation as to the sheer amount of growth. I used to hack a fork of Postgres a few years ago, and the project is probably about 10-20% bigger now:
http://www.ohloh.net/p/postgres/analyses/latest/languages_su...
The problem is hard, there is not even a single solution everyone agrees is the best, reliability and maintenance requirements are stringent, and even in spite of all that the project's code is growing at a rather scary clip.
All in all though, I think a lot of the pre-requisites for what you want are locked up in the attack on logical replication, which is ongoing in the 9.3 release cycle.
Also, the number of tables has nothing to with the problem of sharding, it's the amount of data you're trying to store.
Good thing about HBase is that it's designed to hold a massive amount of data, I've seen people with hundreds of thousands of columns on a single row, it's pretty remarkable.
And the number of tables DOES matter for sharding if you are implementing it yourself.
Never with the sharding part though.