Postgres high-availability cluster with auto-failover and cluster recovery
github.com
github.com
Clustered relational DBs are really remarkably hard to get right. A shocking percentage of projects and vendors don't really even make a serious try (cf GaleraDB.)
This project may or may not have done a reasonable job (no clue,) but the lack of information about how they tested it is a bad sign.
Also PostgreSQL HA has two options, either you have 1 Master and 1 Sync or 1 Async, if you have a Async Secondary it will of course loose data. Also you will have a short downtime on switchover.
And arguably, active/standby with automatic failover is just a special case of active/active, with particar tradeoffs about throughput and failover latency.
"Call Me Maybe: simulating network partitions in DBs"
Running Jepsen against a log shipping automation script is, to my mind, largely redundant.
Running it against this system would be testing Postgres [1], not yoke.
two node systems can never work as a primary/secondary failover solution in the presence of a network partition.
disclosure: I helped write yoke.
It would be nice if the readme listed what failure scenarios are handled and what guarantees you intend the system to have.
EDIT: Shameless plug: I believe it was a shame that Sentinel was initially badly received because of the Jepsen test (now tons of users are using Sentinel with success, fortunately), since for what it does, it is very advanced compared to * SQL failover solutions, so a porting of Redis Sentinel to failover those systems could make sense IMHO. A few things Sentinel has that are desirable in * SQL failover systems:
1. It's distributed so Sentinel itself is not a single point of failure.
2. When the partition heals Sentinel is able to have a coherent view of the configuration.
3. It is able to automatically reconfigure the old master and the other slaves to replicate from the new master.
4. It is able to work as a configuration provider for clients, with a well defined handshake in order to trigger the reconfiguration of all the clients.
It's already possible to use more than one Postgres database while not maintaining consistency, but that isn't desirable for most users.
PostgreSQL is not clustered by default so it is only from the point of view of the theory a CP system. Most proper CP systems have the ability to provide a limited form of Availability, which is, they are available in the majority partition (this is not "A" of CAP but is a lot better than, if primary is not available, the whole cluster is not available). If you use PostreSQL with asynchronous replication and a best-effort failover solution to provide some form of Availability, you are doing something that makes sense, assuming your product can tolerate the resulting weak consistency properties.
This is a great example on how CAP does not capture certain important semantics of a system. Technically a single primary node is CP but the real world CP systems have better availability than that, or they are useless (as CP systems).
- clustering to increase performance - clustering for HA
Most people want at least the later one. The former one is a bigger challenge in case of relational databases so I'll ignore it here.
Regarding consistency vs availability typically is not all black and white, and people are willing to make tradeoff. For example losing few transactions as opposed to being completely down. This is not too bad as long as developers are aware of that special case and write code that can handle that.
Having said that, I'm wondering how a failover system based on BookKeeper[1] would work (a relational database is effectively a state machine, so seems like this could fit). The tradeoff here would be most likely reduced performance though.
Sentinel isn't really meant to solve what Aphyr/Jepsen tests.
Redis Sentinel's semantics are a terrible solution for things I use relational DBs for.
The fact that clustered relational systems accepts transactions in SQL and then just handle them wrong is completely unacceptable.
Redis has the advantage of making it pretty darn clear that it's best-effort, simple ops, non-transactional. If you spent all your time bragging about ACID semantics for Redis, we'd rake you over the coals at least as hard.
(And you don't. Thanks for that! Also, Redis is awesome.)
That is probably be true, but I'd like to point out that clustered non-relational DBs are also really remarkably hard to get right, it's just a hard problem to solve.
The tradeoffs you make with a non relational DB to make this easier are usually also available in relational DBs. It's not like Galera is your only choice.
Like, tradeoffs in consistency? No transactions? That requires not implementing some important bits of SQL.
Funnily enough I ran into this yesterday when Mist was posted and I looked at what other projects were under the umbrella.
Technically you don't need ZFS to use the Manatee state machine. You could swap out the ZFS snapshot transfer service to use pg_basebackup instead.
For anyone interested I have made the appropriate alterations to Manatee to get it to run on Ubuntu with the standard ubuntu-zfs packages.
The fork is here: https://github.com/josephglanville/manatee
Chef cookbook to help you get started along with the other dependencies: https://github.com/josephglanville/cookbook-manatee
Manatee didn't work for our use cases, which should not be interpreted as "manatee is bad". Here are a few of the reasons we forged a different solution:
- While node.js is a great language for many things, we didn't want the overhead of running the node.js/v8 runtime alongside postgres. 50M might seem negligible, but with thousands of postgres clusters for our clients every MB adds up.
- We didn't want the administrative overhead of managing a zookeeper cluster, and instead built the cluster management semantics directly into the yoke project.
- This may be different now, but at the time manatee was very heavily integrated into illumos/smartos. We really love smartos, but recognize that not all clients are able to use this tech.
For example, in this case it's not at all clear how the cluster management works, or even what the semantics are.
https://spilo.readthedocs.org/en/latest/DESIGN/
https://github.com/sorintlab/stolon
I wonder how they compare to yoke.
it's better to have a flexible 'correct' implementation of a DB instead of baking in features which make it more complex and less extensible.
Cf. their JSON support, which regularly outperforms dedicated NoSQL databases; they're also slowly integrating building blocks for multi-master replication, while still holding out on actually implementing it – seeing MariaDB's recent track record with their solution ( https://aphyr.com/posts/328-call-me-maybe-percona-xtradb-clu... ), probably a good decision.
If you read the MySQL documentation, it's full of "until version 5.7.7.3 this returned NULL if..." or "this will accept invalid dates unless a certain mode flag is enabled" (I'm paraphrasing, obviously). Postgres has none of this.
They do occasionally deprecate features, such as OIDs, which are supported but not recommended (or particularly useful), or they change some key behaviour (standard_conforming_strings).
Sometimes the development strategy leads to incomplete or overlapping functionality because they prefer discrete, additive changes over big, one-off overhauls. But in general, Postgres' conservative strategy has really paid off.
Another example, INSERT ... ON CONFLICT UPDATE (UPSERT/MERGE) available just now in the 9.5 beta.
- system stays online when any server goes down
- after confirming a write it will always be returned in corresponding reads following a failover
By skimming the installation instructions I deduce that every time a 'write' call hits the primary server, the primary server sends a message to 'monitor' server which sends a message to primary or secondary server(s) to run rsync?
The monitor node appears to be there to act as neutral witness node so that a secondary is only promoted leader if it can reach the monitor and the monitor cannot reach the current leader. Hopefully, the leader also commits suicide when it cannot reach the monitor.
MS SQL server implements this pattern and calls it "Database Mirroring Witness": https://msdn.microsoft.com/en-us/library/ms175191.aspx
Not saying it doesn't work for them or it couldn't work for me, but I do feel its justifiable to deduct a few "points" there.
Secondly, this is one possible design. The Postgres team prefers tools like these to complement the core so that people have a choice between different solutions depending on their use case, network topologies, availability tolerances, etc.
This just a helper tool that can control Postgres and ship logs around. A native Postgres solution could do it better.