That is, a table on the jepsen.io frontpage, or at least on each product's review page, with database products and configuration on rows and consistency properties on columns, and a nice "Yay!" or "Nope!" mark in the cell, plus links on how to achieve the database configurations in the table (esp. how to configure each database to have the most guarantees).
Also, ideally the analyses should be rerun automatically (or possibly after being paid, but making it easy for the company to do so) every time a new major release happens rather than being done once and then being stale.
Finally, there should be tests for the non-broken databases (PostgreSQL for instance, both in single-server mode, deployed with Stolon on Kubernetes and using the multimaster projects) as well to confirm they actually work.
This is a wonderful idea, and I've got no idea how to actually do it in a standardized, rigorous way. Vendor claims are often contradictory, it's hard to get a good idea of anomaly frequency, availability is... a rabbithole, and it's hard to come up with a standard taxonomy of anomalies--most of the analyses I do wind up finding something I've never really seen before, haha. With that in mind, I've wound up letting the reports speak for themselves.
Also, ideally the analyses should be rerun automatically (or possibly after being paid, but making it easy for the company to do so) every time a new major release happens rather than being done once and then being stale.
I don't know a good way to do this either. Each report is typically the product of months of experimental work; it's not like Jepsen is a pass-fail test suite that gives immediately accurate results. There is, unfortunately, a lot of subtle interpretive work that goes into figuring out if a test is doing something meaningful, and a lot of that work needs to be repeated on each test run. Think, like... staring at the logs and noticing that a certain class of exception is being caught more often than you might have expected, and realizing that a certain type of transaction now triggers a new conflict detection mechanism which causes higher probabilities of aborts; those aborts reduce the frequency with which you can observe database state, allowing a race condition to go un-noticed. That kinda thing.
If I'm lucky and the API/setup process haven't changed, I can re-run an analysis in about a week or so. If I'm unlucky, there's been drift in the OS, setup process, APIs, client libraries, error handling, etc. It's not uncommon for a repeat analysis to take months. :-(
It would be kinda like you including this sort of thing on your resume. Which would also be a bad idea.
https://web.hypothes.is/about/ or similar could be used to develop commentary overlays on top of marketing materials.
Plus maybe a column indicating what [the company behind the database] claims?
Postgres is widely understood to be a robust database with safe defaults. I, and perhaps others, would love to see you aim your array of weapons at Postgres. Do you have any plans to look at stock Postgres?
With tools like repmgr it is just a single command invoked on the standby.
If you absolutely don't want to lose any data, you should have two masters in close proximity (so the latency isn't high) set up with synchronous replication, then have one or two standbys with asynchronous replication. This will reduce throughout, but then you can be sure that the other machine has all the same transactions. If something happens to both you then can fallback to the asynchronous one which might be a bit behind.
Automatic failover for PostgreSQL works great and can be done safely if combined with synchronous replication.
Multiple tools will implement this correctly:
https://patroni.readthedocs.io/en/latest/replication_modes.h... https://github.com/sorintlab/stolon/blob/master/doc/syncrepl...
Quoting a former colleague here, but "if it hurts, do it more often". That is what you should do with your PostgreSQL failovers.
I have clusters running on timelines in the hundreds without a byte of data loss due to using synchronous replication, tools that help out with leader election, and just doing it often.
I would actually be interested if aphyr's analysis of Patroni and other distributed add-ons to PostgreSQL.
The only question is how soon are you going to page humans. After the automated mechanism flipped your master 2-3 times but the cluster still hasn't made progress [nothing coming out of the master; or it locks up after a few minutes again]), or right after some other automated mechanism detects that there's a problem.
Whatever automation you have in place, it has advantages and disadvantages. In the GitHub case - I suppose - they determined post-mortem that it would have been better to just let the master chug through the incoming onslaught of queries instead of failing over, and over, and over. (But of course this seems like a trivial problem in any auto failover setup, so I suspect there's more to the story.)
The postgres documentation will tell you that you'll need to set up your own mechanisms for this, and that they will need to integrate with OS facilities as appropriate. One-size-fits-all does not cut it. Not wrt. replication, not wrt. HA/failover.
No. But the contract Patroni has is this:
I only serve a master (primary) if I have the lock. If I do not have the lock I will demote.
This results in that there can be only 1 primary active at any given point in time, even if the network is partitioned.
This in and of itself does not guarantee no-split-brain situations, a split-brain can occur if writes were made on the former primary, but not yet on the future primary. This however can be mitigated with synchronous replication.
> Postgres has both asynchronous (the default) and synchronous replication options, neither of which offers automatic failure detection and failover [12]. The synchronous replication only waits for durability on one additional node, regardless of how many nodes exist [13]. Additionally, Postgres allows one to tune these durability behaviors at the user level. When reading from a node, there is no way to specify the durability or recency of the data read. A query may return data that is subsequently lost. Additionally, Postgres does not guarantee clients can read their own writes across nodes.
> > Postgres has both asynchronous (the default) and synchronous replication options, neither of which offers automatic failure detection and failover [12]. The synchronous replication only waits for durability on one additional node, regardless of how many nodes exist [13]. Additionally, Postgres allows one to tune these durability behaviors at the user level. When reading from a node, there is no way to specify the durability or recency of the data read. A query may return data that is subsequently lost. Additionally, Postgres does not guarantee clients can read their own writes across nodes.
> From http://www.vldb.org/pvldb/vol12/p2071-schultz.pdf
This is like those commonly seen tables comparing your product with others where your product had checkmarks in all categories, and of course competitors are missing a bunch of them. The problem is that the categories were picked by you, and are often irrelevant to the other product. This is the case here.
PostgreSQL is not a distributed database, the master is the one doing all writes. The replicas are read only. By default replicas are asynchronous which means they won't affect master performance, at the cost of having data there being late by few seconds. Since you can't write to replicas, this won't cause data corruption, only delay which often is acceptable. If you design your applications in such way that will have two database endpoints: one for writes and one just for reads, you can then decide based on context which endpoint you want to use. The read only is easy to scale, but as mentioned earlier it is read only, and might slight delay.
Now, for failover, you might also opt on using synchronous replicas this will add extra latency, but then you always have at least one machine that has the same data. They mentioned that if you have multiple synchronous standbys then it only one needs to write. Actually that's configurable, you can specify group of synchronous machines and how many and which need to be synchronized, the remaining ones are a backup in case those that you specified aren't available.
Besides, the writes don't work the same way as in mongo, when a standby node is in sync it isn't just in sync for that particular write, it is completely in sync, so their following argument about not being able to specify durability/recency of data on read is redundant. If you contact the master or synchronous replica, you will always get the most recent state. If you don't mind slight delay you should query asynchronous replicas (in fact you should prefer them whenever you can, since those are cheap to add)
> the master is the one doing all writes. The replicas are read only. By default replicas are asynchronous
The same is true with MongoDB's defaults in an unsharded cluster.
All the tooling that provides extra distributed functionality not present in postgres (auto failover, multi master replication, sharding etc) will surely have issues, but then you aren't testing the PostgreSQL itself, but the tooling, so to be fair, you the article should evaluate these tools, and any shortcomings shouldn't go to PostgreSQL (unless it really is a PostgreSQL issue).
And looking at this table, basically the future seems to be WAL shipping anyway ( https://www.postgresql.org/docs/13/different-replication-sol... )
I believe RDS Postgres is probably the right answer for lots of applications, especially for those that already depend on AWS for baseline availability. I'd love to see if that holds up against a rigorous analysis.
However, if I may suggest, Stolon, Patroni, Postgres XL or Citus Data might be interesting to you.
Even common highly available configurations take the route of no consistency guarantees by doing primitive async replication and primitive failover.
In a classic single node configuration, a confirmation that its transaction isolation behaviors exhibited the corresponding anomalies would be valuable.
So I think there’s value in this ask.
I think he did something similar for MySQL when evaluating the Galera cluster.
In a single write master configuration, Postgres runs transactions concurrently, so the consistency analysis is still quite relevant.
I don’t think it’s a stretch to say that everyone expects Postgres to get top marks in this configuration and it would be worth confirming that this is the case.
But it was long ago, and maybe needs to be redone?
Edit: after re-reading it he treats it as a distributed system because client and server is over network. And that is true, it can also be thought of as a distributed system because as you said transactions are concurrent and are running as separate processes. Although in these cases you can't have a partition (which aphyr uses to find weaknesses), or maybe there is something equivalent that happens?
Not in itself, but it does offer a PREPARE TRANSACTION - COMMIT PREPARED / ROLLBACK PREPARED extension that could be used to add such support in the future. This would not be unprecedented, as the simpler case of db sharding is already being supported via the PARTITION BY feature, combined with "FOREIGN" database access.
Very happy for (informed) speculation here, I recognise we'll probably never know for certain, but I'm interested to avoid making similar mistakes myself.
The middle part of the report talks about unexpected but (almost all) documented behavior around read and write concern for transactions. I don't want to conjecture too much about motivations here, but based on my professional experience with a few dozen databases, and surveys of colleagues, I termed it "surprising". The fact that there's explicit documentation for what I'd consider Counterintuitive API Design suggests that this is something MongoDB engineers considered, and possibly debated, internally.
The final part of the report talks about what I'm pretty sure are bugs. I'm strongly suspicious of the retry mechanism: it's possible that an idempotency token doesn't exist, isn't properly used, or that MongoDB's client or server layers are improperly interpreting an indeterminate failure as a determinate one. It seems possible that all 4 phenomena we observed stem from the retry mechanism, but as discussed in the report, it's not entirely clear that's the case.
I get the impression that MongoDB may have hyped themselves into a corner in the early days with poorly made (or misleading) benchmarks. Perhaps they have customers with a lot of influence determining how they think about performance vs consistency.
Maybe this combined with patching, re-patching, re-patching again their replication logic/consistency algorithm means that they'll be stuck in this sort of position for a long time.
This section seems to be the most worrying results in your report, Kyle, with no work around. Did I read that correctly?
That's not to say that workarounds don't exist, just that I didn't find any in the documentation or by twiddling config flags in the ~2 weeks I was working on this report. :)
So I'm curios how would you have described the ability of finding violations with Elle using read-write registers with unique values vs the append-only lists?
If you look at Elle's transaction generators, you can cap the size of any individual key, and use an uneven (e.g. exponential) distribution of key choices to get various frequencies. That way keys stay reasonably small (I use 1-10K writes/key), some keys are updated frequently to catch race conditions, and others last hundreds of seconds to catch long-lasting errors.
So I'm curios how would you have described the ability of finding violations with Elle using read-write registers with unique values vs the append-only lists?
RW registers are significantly weaker, though I don't know how to quantify the difference. I've still caught errors with registers, but the grounds for inferring anomalies are a.) less powerful and b.) can only be applied in certain circumstances--we talk about some of these details in the paper.
https://web.archive.org/web/20150312112556/http://blog.found...
Jepsen draws inspiration from a long line of work on property-based testing, especially Quickcheck & co. It also draws on roughly 10 years of experience building & running distributed systems in production. A lot of Jepsen I invented from whole cloth, but some of the checkers in Jepsen are derived from specific research papers, like work by Wing, Gong, and Howe on linearizability checking.
Then it's just... a lot of thinking, experimenting, and writing. Jepsen's the product of ~6 years of full-time work. Elle, the system which detected the anomalies in this report, was a research project I've been puzzling over for roughly two years.
I write the Jepsen series, and open-source all of the code for these tests, partly as a resource so that other people can learn to do this same kind of work. :-)
My bias: I like and heavily use ZooKeeper in production. HN seems not to like it as much.