There's Postgres-X2 (formerly Postgres-XC), sponsored by NTT, which implements fully consistent multimaster replication and partitioning. After many, many years apparently still isn't production-ready nor particularly scalable, and the likelihood that it will ever be merged into mainline Postgres is zero (its planner changes alone are apparently several hundred thousand lines of code).
More interestingly, there's Postgres-XL, which is apparently a merger of Postgres-XC and a different implementation called StormDB that was bought by a company called TransLattice. (TransLattice also sponsored or acquired an earlier project that went nowhere, Postgres-R.) Unlike XC/X2, they say they aim to contribute changes back to the mainline, and they also claim their distributed query model is superior. With the commerical support behind it (it's used as the basis of a commercial product), it's possible that this is something that will be usable. Unfortunately, it seems very quiet and not very open-sourcy; most of the development seems to be by just one guy [3], and nobody seems to be using it in production at this point.
[1] https://github.com/postgres-x2/postgres-x2
[2] http://www.postgres-xl.org/
[3] http://git.postgresql.org/gitweb/?p=postgres-xl.git;a=summar...
Greenplum is an MPP database forked from the Postgres 8.x codebase. It'll be opensourced soon.
Is the recent Dell acquisition putting the open sourcing plans at risk?
Blah blah a lot of it I don't know rhubarb doctrine of confidence murmer cough professional ethics mumble mutter SEC regulations yadda yadda.
So far we have opened Apache HAWQ (incubating)[2], which is an SQL query engine for Hadoop based on the MPP query planner in Greenplum.
And Apache Geode (incubating)[3], which is based on Gemfire, which is the most successful in-memory grid database you've never heard of.
Edit: I suppose I should mention, entirely in my own self-interest, that we are hiring in our Big Data division across all our products. You can look at our webpage (https://pivotal.io/careers) or email me (jchester@pivotal.io).
[1] http://pivotal.io/big-data/press-release/pivotal-introduces-...
It's great that you're releasing stuff as open source. Perhaps you could persuade a colleague at your other division to talk about Greenplum in this thread?
Including: Cloud Foundry, Greenplum, Spring, Gemfire and of course the namesake -- Pivotal Labs.
Email me and I will connect you with our director of engineering for Big Data.
For software some things can be added on later, but some can't. Things like fault tolerance, distribution, security can't be easily bolted on after the product has already been developed (as those usually cut right through the whole stack).
People who like Postgres and make fun of NoSQL products for not having well ... SQL, usually forget that NoSQL isn't as much about not having SQL as about having a good distributed/scalable backend story.
There are also databases coughMongoDBcough which claim to have it all, but in reality they have neither[1][2].
[1] Performance of single instance: http://www.enterprisedb.com/postgres-plus-edb-blog/marc-lins...
[2] Scaling out: http://www.datastax.com/wp-content/themes/datastax-2014-08/f...
Any one of these database types can have a SQL interface (in various dialects and standards), ACID, different consistency models and scale-out options. It's just up to the specific database product in what it actually offers.
Perhaps it should be said that relational is DB that's for general purpose, while specialized (key/value, graph, etc) requires you to understand your data so you can make proper trade-offs when picking specialized db. You are not getting those features advertised by other databases for free.
Relational is specialized too - in dealing with relational data, foreign keys, normalization, joins, etc.
ACID is just a feature of a product, not tied to relational as there are relational databases that don't offer it or can have it disabled. Tradeoffs and understanding your data apply to every single type of database.
HBase is strongly consistent and given that it a core part of the Hadoop stack and used by Facebook for Messages would indicate it is hardly has systemic data loss issues. Cassandra likewise is used for PSN, Steam, EA Online, Spotify and is quite capable of running in strongly consistent mode. That's just two of the many NoSQL databases that are strongly consistent and being used for major systems where data loss would be unacceptable.
Also your comparisons between PostgreSQL and MongoDB are meaningless. The complete overhaul of the default engine in MongoDB 3 (WiredTiger) has increased performance massively with a notable reduction in disk space usage. It's basically a whole new database.
https://www.mongodb.com/blog/post/performance-testing-mongod...
Regarding the Mongo, the second link I provided was comparing WiredTiger. In link you provided they made an improvement, but 2.x was already slow. Initially MongoDB gained its popularity because it supposed to be fast. And it was when it did not worry about making sure the written data ends up on disk. Once people started losing data when system crashed they changed defaults, but then became much slower. And worst of all it still was losing data[1]. There's no good reason to use it, even if you would want to use it as a cache, there are far better solutions.
[1] https://aphyr.com/posts/322-call-me-maybe-mongodb-stale-read...
You're arguing the word "consistent" is used wrong, and then you go on to mixing ACID and CAP — two acronyms which use the word "consistent" with completely different meanings.
The "C" in ACID stands for data consistency, which is that you shouldn't be allowed to store any data which isn't internally consistent with the schema's consistency rules. For example, you're not allowed to store a "null" value in a "not null" column, and you can't delete a row that is pointed to by a different table's foreign key constraint.
The "C" in Brewer's original CAP theorem really represents atomic consistency. It's the idea that if you submit data, you should get the exact same data back, even if you ask a different node. It sort of represents both the "A" (atomicity, in the sense that either everything or nothing is updated) in ACID, as well as the "I" (isolation, in the sense that nodes don't see partial writes) and "D" (durability, in the sense that data written stays written).
In CAP terms, HBase is strongly consistent, I believe, as long as you query the primary region server and not a replica. But HBase is not transactional, and therefore it's pretty meaningless to talk about ACID.
By the ACID definition a database can either be consistent or not. It can't sometimes be one or the other (unless there is some bug. I was referring to the CAP definition which I assume is what you are referring to. And by that definition HBase is absolutely strong consistent.
And that second link was pretty bad (with all respect to the great guys over at Datastax). It is benchmarking a release candidate build of MongoDB (2.8.0-rc4) where quite a number of important fixes went in before GM. They should have just waited until 3.0 was released.
Also it would be good for that Call Me Maybe test to be redone considering that stale-read bug has been fixed:
Not entirely accurate. Relational database schemas are typically normalized into many tables and make heavy use of joins; it's fair to say that the concept of joins are integral the relational model. But joins, by definition, require fast access to data; it's extremely inefficient to join across network boundaries.
So-called NoSQL databases aren't just about ACID or CAP, but also about denormalization. In fact, I would go so far as to say that the single founding tenet of all of "NoSQL" was denormalization across multiple nodes. Google's Bigtable (on which HBase is based, and which arguably started the public's interest in "NoSQL") depends on extreme denormalization in its data model‚ which does away with tables and instead treats everything as columns, and also discards referential integrity (the "C" in ACID, which isn't the same as in CAP) and isolation (the "I"). Sharding flows naturally out of denormalization. Eventual consistency came later with Amazon's Dynamo paper; Bigtable was strongly consistent.
A consequence of denormalization is that you can hit any database node and get consistent performance. No query will ever get bogged down in complicated, long-running nested loops and hash joins because the mechanisms don't exist. It creates an even playing field for queries because they're mostly individual key lookups; the result is (in theory) predictable, stable performance across the cluster. The complexity of join logic that was in the relational database is moved into the client, which gets the burden of untangling data. Which is why Google also invented early solutions like MapReduce and Sawzall, because once you have a huge, denormalized, distributed cluster of data in this poor man's relational model, your individual clients can no longer access the data efficiently. So denormalization also leads to parallel computation, and especially precomputation, rather than the ad-hoc queries that relational databases are good at.
It's all about scale, and fitting the data model to work in your favour at that scale.
So yeah, I would definitely say that relational databases are designed in a way that make them hard to distribute. The above is just the tip of the iceberg. How do you enforce referential integrity (foreign keys, etc.) across nodes? How do you assign unique primary keys? You talk about making tradeoffs, but would you really want a distributed relational database without distributed transactions?
Can you implement a performant distributed database with joins and other relational features? Yes. VoltDB and Vertica come to mind. They are designed from the ground up to do this, and have complicated logic for partitioning queries. They make other tradeoffs: for example, VoltDB is in-memory and to perform well, all transactional access code has to be written as a "stored procedure", a Java JAR file that you upload to the server, and transactions are expected to be very short-lived. But they will, pretty much by definition, not perform at scale the way that "NoSQL" databases are designed to.
repmgr is an open-source tool suite to manage replication and failover in a cluster of PostgreSQL servers. It enhances PostgreSQL's built-in hot-standby capabilities with tools to set up standby servers, monitor replication, and perform administrative tasks such as failover or manual switchover operations.
repmgr has provided advanced support for PostgreSQL's built-in replication mechanisms since they were introduced in 9.0, and repmgr 2.0 supports all PostgreSQL versions from 9.0 to 9.4.
https://github.com/2ndQuadrant/repmgr src: http://instagram-engineering.tumblr.com/post/13649370142/wha...
For Sharding try:
http://instagram-engineering.tumblr.com/post/10853187575/sha... http://rob.conery.io/2014/05/29/a-better-id-generator-for-po...
For clustering try:
https://github.com/smbambling/pgsql_ha_cluster/wiki/Building... https://wiki.postgresql.org/wiki/Replication,_Clustering,_an...
I too wish there existed a nice dashboard db management tool like the ones Rethinkdb and Couchbase provide out of the box.
If you design your application to communicate to your database through a certain sharding solution, and the you find that sharding product becomes abandoned, you can be in a very difficult position.
That wiki lists at least one product as "stalled".
If I run into a Postgresql bug, am I going to get told "we have no way of debugging that with your third party clustering extension in place"?
These are all obstacles that can be overcome, but just jumping on a random clustering project should be done with caution.
As a related part of clustering, I'd love to be able to do downtime free version upgrades like in Oracle - it would remove one of the major reasons for clustering.
Which you can do, but it takes pgpool + VIP trickery, and is a lot of setup for a small project. Hosted from AWS or similar is a lot better as they do the crazy setup, but not for local setup. Nice Digital Ocean finally got floating IPs this week which should help there (have not looked into if it will work for PG failover).
Plan for uptime and you'll be more likely to get it, than planning around a buggy server process.
(Also, in this case, it would have violated SOX and PCI rules, as well as general sanity, to host a 512GB 64-core server in the cloud)
So last week I was setting up a sql server for a side project I was doing with a friend, and went with MySQL just because I knew it would take me ten minutes to setup async replication, whereas with PostgreSQL, I'd have to spend a bunch of time figuring out what third party thing to use now.
Even just having simple master/slave asynchronous replication functionality built in would help a lot.
Check out [0], or the docs on the same topic for the latest version at [1]. It looks like built-in async master/slave replication has been available since at least 2010-09-20, in Postgres 9.0.
[0] http://www.postgresql.org/docs/9.1/static/high-availability....
[1] http://www.postgresql.org/docs/9.4/static/high-availability....
There are some parts I'd like to see smoother like checking replication status without DB superuser creds, but the general built-in replication seems fine.