An overview of distributed Postgres architectures
crunchydata.com
crunchydata.com
The cons of the mentioned distributed shared-nothing SQL databases are questionable:
- "Key-value store" is in fact an LSM-tree-based document store that supports column-level versioning (Postgres supports row-level versioning only).
- "Many internal operations incur high latency." - I guess this conclusion is based on the referenced Gigaom benchmark that was paid for by Microsoft to compare apples to oranges.
- "No local joins in current implementations." (YugabyteDB certainly has colocated tables that store a whole table on a single node. CockroachDB and Spanner might do this as well.)
- "Not actually PostgreSQL..." - There is only one 100% compatible database with Postgres...It's Postgres itself. Citus, CockroachDB, Aurora, Alloy, YugabyteDB, and others can be classified as "not actually Postgres."
- "And less mature and optimized." - Well, both CockroachDB and YugabyteDB are on Gartner's Magic Quadrant of the top 20 cloud databases. Toys don't get there.
It feels like the author joined Crunchy to work on their own distributed version of Postgres. Good move for Crunchy, good luck!
> Related tables and indexes are not necessarily stored together, meaning typical operations such as joins and evaluating foreign keys or even simple index lookups might incur an excessive number of internal network hops. The relatively strong transactional guarantees that involve additional locks and coordination can also become a drag on performance.
You handwaved this away saying you can just store an entire table on a single node, but that defeats many of the benefits of these sharded SQL databases.
Edit: Also, before attacking the author's biases, it seems fair to disclose you appear to work at Yugabyte
- true Index Only Scan. PostgreSQL doesn't store the MVCC visibility in indexes and have to look at the table even in case of Index Only Scan. YugabyteDB has a different implementation of MVCC with no bloat, no vacuum and true Index Only Scan. Here is an example: https://dev.to/yugabyte/boosts-secondary-index-queries-with-... This is also used for reference table (duplicate covering indexes in each regions)
- Batching reads and writes. It is not a problem to add 10ms because you join two tables or check a foreign key. What would be problematic is doing that for each rows. YugabyteDB batches the read/write operations as much as possible. Here are two examples: https://dev.to/franckpachot/series/25365
- Pushdowns to avoid sending rows that are discarded later. Each node can apply PostgreSQL expressions to offload filtering to the storage nodes. Examples: https://dev.to/yugabyte/yugabytedb-predicate-push-down-pbb
- Loose index scan. With YugabyteDB LSM-Tree indexes, one index scan can read multiple ranges, which avoids multiple roundtrips. An example: https://dev.to/yugabyte/select-distinct-pushdown-to-do-a-loo...
- Locality of transaction table. If a transaction touches to only one node, or zone, or region, a local transaction table is used, and is promoted to the right level depending on what the transaction reads and writes.
Most of the times when I've seen people asking to store tables together, it was premature optimization, based on opinions rather than facts. When they try (with the right indexes of course) they appreciate that the distribution is an implementation detail that the application doesn't have to know. Of course, there are more and more optimizations in each release. If you have a PostgreSQL application and see low performance, please open a git issue.
I'm also working for Yugabyte as Developer Advocate. I don't always feel the need to precise it as I'm writing about facts, not marketing opinions, and who pays my salary has no influence on the response time I see in execution plans ;)
I just clarified one-liners listed under the closing "Cons" section. My intention was not to say that the author is utterly wrong. Marco is a recognized expert in the Postgres community. It only feels like he was too opinionated about distributed SQL DBs while wearing his Citus hat.
> Also, before attacking the author's biases, it seems fair to disclose that you appear to work at Yugabyte.
I'm sorry if I sounded biased in my response. I'm with the YugabyteDB team right now, but that's not my first and I bet not the last database company. Thus, when I respond on my personal accounts, I try to be as objective as possible and don't bother mentioning my current employment.
Anyway, I'm very positive to see that this article got traction on HN. As a Postgres community member, I truly love what's happening with the database and its ecosystem. The healthy competition within the Postgres ecosystem is a big driver for the database growth that's becoming the Linux of databases.
That was the first thing come to my mind when I read the paper on Spanner and CockroachDB (haven't read the paper on YugabyteDB yet) though, and surely I'm not the only one.
We had to provide our own active / active storage backend, and fabric. If there were any hiccups, the entire system fell over. The horizontal scalability was nice, but if you caused I/O saturation on the backend, you'd end up knocking over the entire cluster due to shared components. Several times, the entire DB just "broke" (this was a couple scenarios, one where the DB was using 100% of CPU, one where the DB was frozen, and not allowing connections, and things like hung queries that couldn't be stopped), and it required restarting the whole cluster, for which there was minimal tooling.
Perhaps SQL server is better, but that comes with an entirely new ecosystem, and other problems.
My biggest thing is being able to crack open the DB, and look at the source code when it breaks. In this, the likes of Datastax Cassandra and Cockroachdb are great, but I wouldn't call them "proprietary" by any means.
Also, Cassandra is active-active as well.
I'm sure there are others, but I'm less familiar with them (Yugabyte, Couchbase, etc..)
So it's not consistent in the CAP Theorem sense. Comparing it with other databases that don't lose writes and don't have split brain issues makes no sense
Of course, the issue with all shared-storage systems is how much it costs to have a reliable and fast shared storage.
This seems highly dependent on how you define “critical”. I think most people’s definition allows for everything to be in the cloud.
Specifics are here: https://www.yugabyte.com/postgresql/postgresql-high-availabi...
The free alternative would be Mysql/Mariadb + Galera Cluster. Not as solid as proprietary ones, but far easier to use and less buggy than Postgres + tons of tools.
Until someone accidentally run an expensive DDL on your Galera Cluster: now your cluster is down for hours without anyway to cancel that query except nuking the entire database and restore from backup.
And the only fact about performance is a benchmark comparing elastic and resilient distributed SQL to non-HA Citus running on larger machines.
Good, fair and reproducible benchmarks are a rarity. Do you have any (independent) benchmarks that compare different distributed PostgreSQL-based solutions?
But I agree that this benchmarks seems questionable at best then. I usually only work with on-premise deployments, so I don't know the details of all the cloud offerings. People will probably judge products on what is available on the first few search results though.
Though I do agree calling distributed SQL as key value store inaccurate and maybe a bit inappropriate.
I think it often goes overlooked just how slow network attached block storage is though, and some organizations get very surprised when moving from an on-prem data center to cloud.
One thing I'd add is a sense of scale - are these architectures for 100 queries per second or 100,000 or 100,000,000 ?
Single node Postgres (with a beefy machine) can definitely manage in the 100k transactions per second. When you're pushing the high 100k into millions read replicas is a common approach.
When we're talking transactions, question of is it simply basic queries, bigger aggregations, and is it writes or reads. Writes if you can manage to do any form of multi-line insert or batching with copy you can push basic Postgres really far... From some benchmarks Citus as mentioned can hit millions of records per second safely with those approaches, and even without Citus can get pretty high write throughput.
cache is King
Some apps might do just 1000 ops/second but still run on a distributed database for high availability or data locality reasons. For instance, shared-nothing databases usually guarantee RPO=0 (no data loss, recovery point objective) with RTO (recovery time objective) measured in seconds for zone and region-level outages. As for data locality, think automatic data placement/pinning to regions/data centers for data regulatory and low latency reasons (serve read/write requests equally fast for folks living in NYC, London, Tokyo).
But once you outgrow the primary/standbys severs storage or compute capacity you would need to scale to larger machines that can incur downtimes. With distributed Postgres such as YugabyteDB this is not gonna happen because you can scale horizontally
It seems to me that you would need to run some sort of consensus algorithm to ensure the replication is consistent but that’s obviously very expensive in latency. Is it actually done this way?