Distributed consensus reloaded: ZooKeeper and replication in Kafka
confluent.io
confluent.io
So while providing a lot of flexibility and decreasing the communication overhead when using a lot of nodes, it may introduce performance degradation and latency spikes when nodes are slow. I'd probably prefer a design where some kind of reduced majority is required, but not a full list of nodes, so that outliers would be filtered out.
Are there good percentile numbers over there to check metadata writes over the, say, 99,99%-99,999% percentile, to check how bad this may become?
If I understand your comment correctly, it isn't the case that the ISR is fixed. The ISR has a minimum size, but it can change over time, so brokers can be removed from the ISR and they can rejoin later once they catch up.
One can drive latency almost arbitrarily low if one is willing to give up exact knowledge of a where the data is.
A simple example is N out of M writes (N and M can be arbitrary with N <= M). M is a set of machines, N is number of replicas we want to have. Now at any given time we write to the N machines that respond fastest. Data is now arbitrarily sprayed over the M machines, but as long as the data itself has ordering information the right state can be reassembled by talking to the M machines.
The "spray" can now be controlled by favouring some nodes from M up to a timeout (i.e. we put more uncertainty in time). We can reduce the reassembly work by using learners and hence increase the likelihood that one machine has all state.
http://research.microsoft.com/en-us/um/people/lamport/pubs/v...
That's not a criticism at all, though. There are a lot of good design and systems ideas there, and they are well explained in this post.
The subject is exciting and I wish to learn more on how to design an efficient and safe replication scheme on top of two coordination protocols, each with its core set of garanties and constraints.
The post gives a good overview of the trade-off between safety and performance made with in-sync replicas (ISR) compared to quorum acknowledgement.
But it remains very vague on how to deal with the problem stated in the "ZooKeeper and consensus" section: being consistent does not mean that the values read [by two workers] are the same necessarily but are only computed after consistent growing prefixes of the sequence of updates. I totally fail to understand the way proposed to break the tie.
I would expect the schema better shows late replica and even late views of the ISR. I would expect more evidences on how the system ensures a message produced to a consumer is never retracted.
What I find vague in "ZooKeeper and consensus" is the answer to "Why does this proposal work compared to the original one?":
>> Because each of the workers has “proposed” a single value
Does the processes read or propose the value ?
>> and no changes to those values are supposed to occur.
Not suppose to ? How do you ensure that ?
>> Over time, as the configurator changes the value,
>> they can agree on the different values by running independent instances
>> of this simple agreement recipe.
Sorry, but I fail to see what recipe you speak about.
It does both, it first proposes by writing a sequential znode and then reads the children (all proposals are written under some parent znode). This is certainly assuming some experience with ZK, and I wonder if that's the problem. It was not the goal to go into a discussion of the ZooKeeper API, but I'm happy to clarify if this is what is preventing you from getting the point.
>> Not suppose to ? How do you ensure that ?
It is not supposed to in the sense that if this is implemented right, then the proposed values for each client won't change. The recipe guarantees it because it assumes that each client writes just once.
>> Sorry, but I fail to see what recipe you speak about.
I'm referring to these three steps: creating a sequential znode under a known parent, reading the children, picking the value in the znode with smallest sequence number.
Indeed, I have no experience in programming zookeeper, even I use it a lot behind tools like kafka or storm.