KIP-500: Replace ZooKeeper with a Self-Managed Metadata Quorum
cwiki.apache.org
cwiki.apache.org
"Currently, Kafka uses ZooKeeper to store its metadata about partitions and brokers, and to elect a broker to be the Kafka Controller. We would like to remove this dependency on ZooKeeper. This will enable us to manage metadata in a more scalable and robust way, enabling support for more partitions. It will also simplify the deployment and configuration of Kafka"
Looking forward to seeing if this gains traction.
Zookeeper does expect a single “host” for every server but it runs perfectly fine with docker. But that’s no different from consul or etcd.
Dynamic reconfig in 3.5 addresses the "restarting every zookeeper instance" problem. [0] You stand up an initial quorum with seed config, then tie in new servers with "reconfig -add". Not sure how well it would tie into cloudy autoscaling stuff though. I wouldn't start there myself.
A much bigger pain IMO is the handling of DNS in the official Java ZK client earlier than 3.4.13/3.5.5 (and by association, Curator, ZkClient, etc.). [1] The former was released mid 2018 and the latter this year, so tons of stuff out there that just won't find a host if IPs change. If you "own" all the clients it's maybe not a problem, but if you've got a lot of services owned by a ton of teams it's ... challenging.
Even with the fix for ZOOKEEPER-2184 in place I'm pretty sure DNS lookups are only retried if a connect fails, so there's still the issue of IPs "swapping" unexpectedly at the wrong time in cloud environments which can lead to a ZK server in cluster A talking to a ZK server in cluster B (or worse: clients of cluster A talking to cluster B mistakenly thinking that they're talking to cluster A). I'm sure this problem's not unique to ZK though.
Authentication helps prevent the worst-case scenarios, but I'm not sure if it helps from an uptime perspective.
TL;DR: ZK in the cloud can get messy (even if you play it relatively "safe").
[0] https://zookeeper.apache.org/doc/r3.5.5/zookeeperReconfig.ht... [1] https://issues.apache.org/jira/browse/ZOOKEEPER-2184
Should it not be the case that one of the responsibilities of the quorum is to vote new members in or out? I mean, as a first-class feature of the system.
The main problem I see is that in most consensus systems, any member can be nominated as leader. But if the member was inducted during a partition event - which one could do if the partition were long duration - then nominations for new leaders will go out. What happens if the new member gets elected? How do the partitioned machines find that leader when they return?
So the process of joining the quorum would have to be incremental. Because demanding a unanimous vote to add a machine means you can never replace a dead one (except by impersonating it)
Indeed. What you're describing is one of the main motivations for Raft. Paxos showed that distributed consensus was mathematically sound, but did little to guide implementors in actually building such a system. Raft is not a fundamentally new consensus algorithm; just an incremental improvement that formalizes many of the improvements that you needed to make to Paxos anyway, like membership changes, log compaction, and multi-decree support from the get go.
If you're interested, the Raft paper is quite readable and goes into this in detail. [0]
> The main problem I see is that in most consensus systems, any member can be nominated as leader. But if the member was inducted during a partition event - which one could do if the partition were long duration - then nominations for new leaders will go out. What happens if the new member gets elected? How do the partitioned machines find that leader when they return? How do the partitioned machines find that leader when they return?
This isn't actually the tricky bit, as it turns out. Communication flows from the leader to the other nodes, so the leader, even if it's a newly-inducted node, will initiate the connections to the partitioned nodes when the partition resolves. (The leader necessarily knows the addresses of the partitioned nodes, because the leader knows about all committed entries, and the identities of all the nodes in the cluster are committed into the Raft log.)
Again, the Raft paper does a great job explaining cluster mebership changes—much better than I can!
Goes to show you should always go back to the source at some point, even if the third party descriptions have better facility.
Alternatively, you can create an ASG-per-node so that you get auto-replacement.
In my experience, ZK is one of the easiest and most reliable distributed systems to operate. I've only seen issues when it's used as a database instead of a distributed coordination service.
Even though Netflix hasn't updated it in a while, Exhibitor was helpful as well, in that it allowed ZK to bootstrap nodes off of state stored in S3. That did come at the cost of an extra 2-3 minutes per node on initial quorum startup.
I had a single Zookeeper ASG running under load for over three years without maintenance. I pinged my old co-founder to see if he's willing to open source the CloudFormation template.
edit: i guess i should RTFA -- looks like they are going to build raft into kafka -- at this point, I'd think raft libraries are fairly mature.
https://cwiki.apache.org/confluence/display/KAFKA/KIP-273+-+...
Considering most Kafkas will probably run in Kubernetes at some point, they could have shared the etcd used by Kubernetes.
To end users of whatever platform it is keeping together, Zookeeper is often seen as an unwanted dependency. The unwanted sentiment arises because maintaining zookeeper is not a gimme, and it takes additional knowledge and overhead to maintain. You need a quorum to be kept alive with an odd number (greater than one) of instances. So the question often arises, "how can we get rid of zookeeper?".
Edit: Oh right, the fact that etcd is golang might make that an issue for Kafka...
If you were building a kafka like system in house for private use only, then yes I agree with your sentiment. If you are building something to be used by hundreds or thousands of organizations then the cost benefit tradeoffs shift to where it probably makes sense to pull the consensus logic into the primary application itself.
As someone who enjoys sleeping at night, I wouldn't poke something like that with a stick before it has matured for a few years, and before many brave (?) operators have smoothed out most rough edges by landing on them repeatedly with their faces.
It does not simplify the system. It is simply replacing the ZK ensemble w/ its own raft impl.
Same number of JVMs, same operational complexity.
> although ZooKeeper is the store of record, the state in ZooKeeper often doesn't match the state that is held in memory in the controller
Maybe it would be better to put the effort in to changing ZooKeeper instead of writing their own consensus system in Kafka, but I would assume they know that tradeoff better than us since they work so closely with ZooKeeper.
Maybe the correct solution is to understand why the states don’t match in the first place?
But then we got the Raft paper, which has single-handedly essentially managed to commoditize distributed consensus. Raft is an elegant and simple protocol that has been implemented in many languages. To implement primitives like leader election and consistent, distributed reads/writes, you can grab an off-the-shelf library like Hashicorp's Raft library, which does all the heavy lifting as long as you implement the state storage yourself. It's absurdly trivial to turn a single-node app into a replicated, fully redundant one.
Of course, a project like Kafka might have particular requirements (I wouldn't know) that would mean Raft would not be a suitable solution for them, or would require modifications to the protocol. CockroachDB, for example, had to modify Raft ("MultiRaft") in order to achieve the scalability they needed within Raft's paradigm. Raft is a starting point more than a finished solution.
You can have a perfectly functioning consensus algorithm and still battle inconsistent reads across nodes, loss of data on node failure or network splits, write skew, and so on. For example, TiKV uses Raft, yet the problems uncovered in TiKV have nothing to do with Raft itself. Conversely, those found in Elasticsearch are absolutely caused by its consensus algorithm, Zen, which was invented in ignorance of modern consensus research (Paxos, etc.).
If you read my comment carefully, you'll note that it's above all an endorsement of Raft, which is a milestone in that it democratizes advanced consensus and makes it a tool anyone can use. My point isn't that anyone can trivially whip up an infinitely scalable, consistent database because consensus is a solved problem, but that we have solid foundational theory to build on, to the point where it is a well-understood problem (albeit with many complicated implementations) and not black magic.
Also, I did not use the term "easy". But I also think the OP's phrasing "extremely hard" is not true, unless your requirements are extremely hard – which, again, may be the case with Kafka; scalability tends to be diametrically opposed to consistency.
If that assertion is true, then the replacement will exist because something in Raft is difficult or impossible to do. In which case consensus is hard, for some quantity of hard.
The one that's most obvious to me is that Raft violates one of the 8 Fallacies of Network Computing: the network is homogeneous.
The star topology of messages means that some messages - identical messages with different destinations - will be competing for bandwidth. It also means electing the slowest machine on the longest network link is probably a bad idea. I've certainly heard of this sort of problem happening.
Zookeeper is great. Weird, for sure, and a bit of a pain in the arse. But it will neither lie to you nor lose data and this is much, much harder to achieve than you might think.
In my past I've seen this many times and each time people went back to Zookeeper after a while, because - as it turns out - consensus is hard; and Zookeeper is battle hardened.
Secondly, it's a pretty heavy dependency. The JVM is RAM-hungry and it's difficult to ensure that it always has enough RAM. Running multiple JVM apps on a single node must be done carefully to make sure each app has enough headroom. It consumes considerably more RAM than Etcd and Consul.
Thirdly, I think it's fair to say that ZK is showing its age. It's notorious for being hard to manage (see the other comments in this thread), with a fairly old design (based on the now-ancient Google Chubby paper) that, while resilient, is also less flexible than some other competing designs.
https://issues.apache.org/jira/browse/ZOOKEEPER-2164
https://issues.apache.org/jira/browse/ZOOKEEPER-2791
Requests (can't find them in JIRA at the moment, so I need to paraphrase) in the past to have a call to initiate a controlled leadership move to another node have been turned down as "you don't need this" yet leadership election fails in some circumstances! In addition there's no command or configuration to disable FastLeaderElection.
So the zookeeper maintainers keep operators limited to having to flip nodes off and on again, which is really a bad way to manage software because it impacts clients as well as leadership (and even if clients recover, most code that I've seen like to make some noise when zk connections flap). I would really like to eliminate all use cases for zookeeper where there is a chance that the zxid will exceed the size of its 32-bit counter component in the span of, say, a decade so that as an operator I don't have to set alerts on the zxid counter creeping up, and having to reset zookeeper and restart all of its clients (many versions of many zookeeper clients don't retry after connection loss, don't retry after a timeout, don't cope with the primary connection failing, will have totally given up after 15 minutes, etc.).
I think that the kafka maintainers have been doing a better job of actively maintaining their code and ensuring it works in adverse conditions, so I'm on board with this proposal.
Zookeeper isn't magic, it's just pretty good at most of what it does, and I think that projects that understand when they've pushed zookeeper into a bad corner may benefit from this kind of move, if they also have a good idea of how they can do better.
Previous HN discussion: https://news.ycombinator.com/item?id=13449728
> enabling support for more partitions
I don't know if anyone of you ever ran a high-throughput Kafka cluster with a large number of partitions (as in, thousands of them), but its not pretty. Rebalancing can easily take half an hour after a rollout, and throughput is degraded during that time. We recently had to move to shared topics because it became untenable.
This is a very welcome change!