Master-Less Distributed Queue with PG Paxos
citusdata.com
citusdata.com
In particular, suppose I do the following sequence of operations on an empty table:
BEGIN
INSERT INTO table VALUES ('foobar', 123);
SELECT * FROM table;
ROLLBACK;
Would the INSERT be submitted to the Paxos log immediately, causing it to be applied on other replicas even though the transaction never committed? Or would it be deferred until commit time, causing the SELECT to return an empty result? Or is there something more sophisticated going on?Submitting INSERT into Paxos log immediately would allow the rest of the cluster to see those 2 logical states of the table. Thus one can imagine [ have no idea whether Citus implemented it or not] that if failover happens after INSERT before SELECT, the cluster would still be able to continue the transaction and executing the transaction's SELECT on another node would still produce 1 record while any other session would still see empty table.
Upon COMMIT or ROLLBACK, the same as in one node situation, one of those table states will become permanent and another will disappear. Through Paxos log whole cluster would be aware (and in consensus :) about it. If there were other transactions in progress, the relations between them would be subject to isolation levels configured as usual. For example if there were another transaction running with READ_COMMITTED isolation level on some another node and that transaction were wanting to do SELECT on that table, it would see the COMMIT and thus would be able to read that one record.
Paxos serialisability of the state of the cluster has nothing to do with transaction isolation levels. Paxos just provides that all nodes are in consensus about the state, for example the state of relationship between transactions in progress.
- Conflict-free. No need to apply conflict resolution.
- Synchronous. No eventual consistency. This is strong consistency (with C in the CAP meaning).
- Master-less. No need to elect a master. All nodes are treated equal.
Paxos introduces higher latency and lower throughput, so it is best used for low bandwith applications, like metadata management. But many applications may benefit from the very strong guarantees that Paxos in general and pg_paxos in particular, provide.
Is it really masterless or do the nodes elect a master using Paxos to hold an election?
"The drawbacks are high latency in both reads and writes and low throughput. Pg_paxos cannot be used for high performance transactional systems. But it can serve very well for low-bandwith, reliable replication use cases."
http://michael.otacoo.com/postgresql-2/postgres-9-6-feature-...
In any case, synchronous replication means that all of the participating servers have to participate in the replication process. If one of them slows down or hangs, replication (and your transaction) does not proceed.
Paxos, on the contrary, can proceed when N/2+1 of the nodes are available. That's a huge difference, and it's irrespective of the latency and performance. In other words: while 9.6's synchronous replication is a really welcomed addition, a single miss-behaving node will halt transactions on the cluster, while pg_paxos will continue operating without problems. Both are meant for different use cases.
The difference is that with Postgres' replication, when the master fails, write operations can't be executed until a new master is promoted. This has to be done carefully, because you want to make sure that no in-flight operations are still happening on the old master (aka STONITH), and that the most up-to-date slave becomes the new master.
Paxos avoids the need for manual (or very delicately-automated) failover, at the cost of extra network round-trips and disk syncs on every operation.
But there are more differences between both setups:
- Paxos is master-less, so you can write to any node (there's no need for a master).
- Failover is very tough to get it right. Indeed, other than consensus, there are no other bullet-proof solutions to achieve it under any circumstance, so relying on a master is a significant difference.
Regarding the extra round-trips and syncs, they can be pipelined if wanted too. I wouldn't conclude this is necessarily slower (it of course depends on the Paxos imlementation) until properly benchmarked.
- Compared to a single Postgres instance, you get better availability for reads and writes
- Compared to a master-slave failover process, you get automated failure recovery without the risk of split-brain or lost writes.
- Compared to a NoSQL database, you get strong consistency, and all the powerful SQL support that Postgres provides.
The cost is higher latency and lower throughput for both reads and writes.
> Does this only work for append only unique data structures (ex: immutable log style)?
Paxos is based on a technique called state machine replication. It replicates an append-only log of changes to an initial state, which allows you to replicate arbitrary data structures. For example, pg_paxos logs SQL commands on a table (e.g. UPDATE).