Confluent Developer ft. Tim Berglund, Adi Polak & Viktor Gamov

KIP-500: Apache Kafka Without ZooKeeper ft. Colin McCabe and Jason Gustafson

Confluent, original creators of Apache Kafka® Season 1 Episode 55

Use Left/Right to seek, Home/End to jump to start or end. Hold shift to jump forward or backward.

0:00 | 43:46

Tim Berglund sits down with Colin McCabe and Jason Gustafson to talk about KIP-500. The pair, who work on the Kafka Core Engineering Team, discuss the history of Kafka, the creation of KIP-500, and what it will do for the community as a whole. They break down ZooKeeper's role in Kafka, the implications of removing ZooKeeper dependency, replacing it with a self-managed metadata quorum, and how they've been combatting security, stability, and compatibility issues. With pending improvements towards scalability and inter-broker communication, and now that KIP-500 has been adopted within the community—there's a lot covered in this episode that you won't want to miss!

EPISODE LINKS

SEASON 2
Hosted by Tim Berglund, Adi Polak and Viktor Gamov
Produced and Edited by Noelle Gallagher, Peter Furia and Nurie Mohamed
Music by Coastal Kites 
Artwork by Phil Vo 

  •  🎧 Subscribe to Confluent Developer wherever you listen to podcasts. 
  • ▶️ Subscribe on YouTube, and hit the 🔔 to catch new episodes.
  • 👍 If you enjoyed this, please leave us a rating. 
  • 🎧 Confluent also has a podcast for tech leaders: "Life Is But A Stream" hosted by our friend, Joseph Morais.
SPEAKER_02

It's no secret that there's been some interest in the Kafka community in removing Zookeeper. The good news is this is finally happening with the kip that will go down in history with a nice round number of 500. Today I'm joined by two engineers who are working on this, Jason Guffseson and Colin McCabe. They're going to tell us about their progress on building Kafka's own metadata quorum implementation and all the other things that attend that change. We'll dig into it in this episode of Streaming Audio, a podcast about Kafka, Confluent, and the cloud. Hello and welcome back to another episode of Streaming Audio. I'm your host, Tim Berglund, and I'm joined today by two guests that I'm delighted to have on the show, Jason Guftison and Colin McCabe. Jason and Colin, welcome. Thanks, Tim. It's great to be here. What do you guys do? You're uh I can say you're coworkers of mine. Tell us uh what you do.

SPEAKER_00

Sure. Yeah, uh Colin and I, we're both part of um what we call here uh the Kafka Core engineering team. So, you know, we work we work on um you know open source Kafka, like the the the broker and uh the standard clients. So it's kind of like a focus on um you know the internal, you know, the replication, for example, of uh the the logs. Um we focus on um you know storage system, um, how that's working. Uh so basically all parts of this and you know, like the protocols as well.

unknown

Cool.

SPEAKER_02

So it's core Apache Kafka. You're both engineers in that. Yeah, that's right. So uh there's big news, and I know this is uh announced several weeks before we're we're recording this, and this will be released uh you know, probably a few weeks from now. But Kip 500 has been published and announced and even blogged about. Uh pretty big deal. What is Kip 500?

SPEAKER_01

So in a nutshell, KIP500 talks about removing the Zookeeper dependency in Kafka and replacing that dependency with a self-managed quorum for storing the Kafka metadata.

SPEAKER_02

There you go. So if you know Kafka well uh and if you've been working with Kafka for a few years, uh this is like, I don't know, spring rain. People ask me about this all the time. Um, you know, hey, when are you gonna get Razookeeper? That that's a question that comes up a lot in community interactions from experienced Kafka operators. And I always feel like Zookeeper gets it's it's just like the easy target and it gets beaten on a lot, and it shouldn't, you know, it's this awesome piece of infrastructure that grew out of Hadoop and it pops up in all kinds of distributed systems. You know, there's always this little Zookeeper quorum over in the corner that nobody talks about except to complain about it. And it it we should love it more than we do, but apparently we unlove it enough that we're planning on replacing it. And I want to talk through why and what we're replacing it with, and kind of what that vision is and how long this is gonna take and all of that. But I think to put it in context, it would help if you could just give us a little bit of a history of Kafka through this lens, right? Like, what was Zookeeper doing? What does it do now? Kind of work us through that history.

SPEAKER_00

Yeah, sure. I can do that. Um and and to answer one of one of your your points there, I think um, you know, Zookeeper actually is is you know quite a reliable system, and it gets kind of an unfair rap in some cases just because the problem it's solving is extremely difficult. Um so though you know, distributed systems in general are have so many corner cases. And you know, I think the people people see it, people see the problems there, but they they don't look at um you know the other system that might be there if Zookeeper wasn't there. Um so in any case, you know, Z ZK removal um in Kafka, um of course, like for in any kind of system, you know, the fewer dependencies you have, in general, people prefer that um because it kind of unifies all of the um you know the the concerns that you have. Um but this particular one, you know, we we've had this, this has kind of been in the ether for for quite a long time. So um, you know, back when I first started working on Kafka, uh you know, the first the first project that I really had a uh a big part in was um you know changing uh implementing the new uh Java consumer. And you know, one of the main things there, um there were actually like two ten kind of two key parts. So one was actually like just that it wasn't in Scala anymore. But the other big part of it, and I'd say this is actually the more substantial part, was that um, you know, the old Scala consumer uh depended on Zookeeper very directly. So all of the consumer offsets, you know, which is the way the consumer saves its position, um, they were maintained inside Zookeeper. Even the um the way that we rebalanced the consumer group, that was all done through through Zookeeper, and it had a lot of problems, um, you know, because uh, you know, like the the rate of offset commits were something that was diff in Zookeeper was having trouble keeping up with. So one of the big improvements that we made there was uh we implemented the um offset storage inside Kafka just as part of a a regular topic. Um and you know, we we put nice APIs uh behind it with you know uh authorization and such. Um so we've kind of been so that that was kind of like the one of the big first big steps has been uh you know getting rid of the dependence inside the consumer.

SPEAKER_02

So that so that clients never speak directly to Zookeeper.

SPEAKER_00

Yeah, that's right.

SPEAKER_02

And I think the second one is sorry, before you before you go onto that, I was gonna say um that was certainly a good change, right? That's a major architectural simplification. There's just a connection that clients don't have to have. Like when you think about putting Kafka in the cloud, you know, it's just so important that all that happened. But I can pretty much guarantee anybody who is writing their first Kafka, you know, like kind of counterfactual counterfactually, you just imagine you're the person who's building this for the first time. Again, especially 10 years ago, you're gonna use Zookeeper because it just does this consensus thing. It's this nice fast database that that is fault tolerant and and consistent. And you're gonna put consumer offsets there, right? I mean, that makes all the sense in the world. That's that's where you're going to start. Uh but it turned out that that was not a great idea, and we sort of evolved away from that. And so at the end of that first phase, like you said, rewriting the clients in Java and removing the Zookeeper connection, what was the next thing?

SPEAKER_00

Yeah, so this the next thing was um in addition to the consumers, we had a bunch of tooling, which was uh, you know, reaching into Zookeeper directly. So this this involves like like when you create a topic, you know, we didn't have an API for that. Um instead, what you did is you created a Z node inside Zookeeper. So all of the tools kind of you know assumed that you had this kind of uh API. So we had kind of like all these informal APIs, which were just really Z nodes in the end. Um so uh a next big step in this process was uh the addition of the admin client, which actually um you know Colin was uh responsible uh for the the initial kit which uh introduced that. Um so then we were able to take basically all of those tools, um all of those APIs, and we we put them behind um you know a Kafka API. So whether or not they were touching Zookeeper internally from the client perspective, you're just accessing some API. And so this was uh, you know, we had we did the consumer, then we did the admin client. And uh with with this kind of transformation, basically we we we we take the knowledge about Zookeeper, you know, this becomes more of an uh implementation detail that the the brokers are uh you know only aware of. Right, right.

SPEAKER_01

Yeah, absolutely. And I would add that um a big motivation for these changes was this concept of increasing encapsulation. So in the old world, when we had people talking directly to Zookeeper, clients talking directly to Zookeeper, there isn't really much in the way of security that you can have in that kind of world. So if people can sort of directly reach into your brain and just change stuff, then you know there's no security, right? And it's just, well, we trust this guy. I mean, and you know, that's that's suitable for like a project that's just starting out. But eventually people want like access control, they want to be able to set up ACLs, they want to actually be able to set up some kind of security on their cluster. And then you stop uh wanting to be in that kind of world where you can directly access uh the internal data structures. I another aspect of encapsulation is just sometimes the format of the metadata needs to change across Kafka versions. And in that case, having tools that are running around sort of directly modifying those data structures causes problems. Like what if an old tool overwrites data with an old format that was previously in a new format? Like information could be lost. So that's just in general, like having more encapsulation uh was a huge motivation for the admin client, a huge motivation for the client work, which Jason described in general.

SPEAKER_00

Yeah, and and to take a little bit of a step back, um, I think um, you know, starting with Zookeeper is not really a you know a bad thing. Um it's it's a solid system for for what it does. And if you're coming up with a new system um and you you know you have to make a decision, well, am I going to implement some consensus system myself or am I going to use Zookeeper, which is already available, already has a lot of these cases figured out, it totally makes sense to uh use Zookeeper in these situations. I think as a project matures, though, um, you know, you do you do start to find some of the limitations of as Colin was saying, like, you know, have having having more things aware of sort of like the underlying representation and uh being able to access it, mutate it, all of this stuff becomes you know, becomes difficult to manage compatibility, becomes difficult to manage security. So all of this stuff, you know, we we've been thinking about this for for several years. So with the consumer, with um the admin client, and you know, kind of the next step that's um even before um you know Kip 500, um, you know, we we started looking at, you know, even in between the brokers themselves. So um, you know, in order to communicate between two brokers, often what would happen is, you know, we'd do it through Zookeeper. We'd create a Z node, which some other node uh broker would be, you know, we'd expect them to have some kind of a watch for it. So the other part of this, which um is it's a little bit more recent, has been you know, taking some of these inner broker communication mechanisms which are flowing through Zookeeper and then putting them behind an API as well. So basically brokers just communicate to each other the same way that a client communicates with it, just using some kind of uh request RPC.

SPEAKER_02

Okay. Okay. Is that that interbroker communication at present uh that link is broken? I mean, is that has that also sort of transitioned out of Zookeeperland?

SPEAKER_00

Um not yet. That's I think that's part of this work. Yeah. Yeah. So so there I can mention one specific um so you know, one one specific case that uh, you know, we have this notion of a controller. The controller is kind of like the um you know, the the manager of the metadata for the uh the Kafka cluster. And one of the things it's responsible for is kind of you know, when a when a broker shuts itself down, then it has to be taken out of you know uh uh in-sync replica sets and stuff like this. And the way we do that currently is um you know, the there's there's like a two-step process. The broker will send some a message to the controller, then the controller will update Z nodes, and basically like the the a lot of the communication here is actually indirect through the um you know Zookeeper nodes. So uh kip 497, which just got ahead of uh KIP500, was uh one of the changes that we need for this, which is uh you know replacing one this inner bro inner broker communication mechanism um with an actual API. So brokers rather than writing into Zookeeper are going to send a message to the controller saying, you know, I need you to do this, this or whatever. So um, but this this there there are still some remaining. So that I think is part of the the Kip 500 uh roadmap, which we'll probably talk about a little bit more.

SPEAKER_02

We will certainly get to. So in terms of Kafka history, does that bring us pretty much to the present? Are we missing anything?

SPEAKER_00

Yeah, I think that's pretty much it. So yeah, basically like we get getting um the ZK dependence out of the clients was was really the first step. And really there was no conversation uh about removal until we could have that. Um now that all that has has sort of happened, um, you know, I think there are still some straggling um APIs here and there which um we may need to fix. Um but for the most part, the producer has no dependence on Zookeeper, uh consumer has none, admin client also has none. I think there's you know there's just some straggling tools. Gotcha. And the interbroker stuff.

SPEAKER_02

And then there's just the raw function of uh there is metadata about which we must have consensus. It must be stored somewhere, and we have to be able to ask and get an answer we can count on.

SPEAKER_01

Yeah, I I think um the first phase of of ZK removal, which I guess is Jason was just describing like um, you know, adding those APIs, adding that uh those interfaces, was mostly about encapsulation. So it's mostly about having like a an API boundary that we can secure between us and the clients, potentially even between the brokers and other brokers or the rest of the system. As we move on to the next phases of ZK removal, I think there's there's two more big themes that I'd I'd like to comment on. So one of them, I think Tim, you touched on this earlier, is is this theme of deployability. Why do I have to manage this separate thing, right? Why do I have to set up a configuration which is in a totally different format? Maybe the security is configured in a totally different way. So why do I have to manage two services when I really want to be managing one? So that's that's deployability, right? And that's something that removing Zookeeper will get us. And I think upgradability is another issue which is sort of like related to this. Um we sort of hinted at this earlier, but in general, the idea that we want to decouple software from the representation of the metadata, right? So the more services you have directly writing metadata in a specific format, the harder upgrades will be. Because you can always have someone who overwrites your new stuff with old stuff. So getting rid of Zookeeper will actually make this a lot easier because we can store the metadata in our own format and have our own APIs for changing it. Uh, this might be a little hard to see without like a detailed explanation, but I'll touch on the stuff.

SPEAKER_02

I think it makes sense. You kind of, when you were talking about the admin client, yeah, uh, that's kind of the problem you're solving, where there are these external actors that have sort of claims on metadata, they want to write things, they want to read things, and that that API has to be exposed directly to them versus um have essentially the broker expose an API and then the broker manage its metadata in some way that we all assume is good and works, but it's encapsulated. We we don't need to know that. And when the broker is versioned, if it also needs to upgrade or version or change its metadata store, well, it'll know how to do that based on you know its own internal upgrade behavior.

SPEAKER_01

Yeah, absolutely. And you know, I mean the problems that we have with clients, we actually have inside Kafka too. Like uh when you're doing a rolling upgrade, you actually have two versions running in parallel, and they need to be able to work together. And having well-designed APIs helps a lot with that.

SPEAKER_02

Yeah, and of course, also of note, it's really hard to do that, right? That's a lot of special effort and a lot of engineering, yes, you know, entropy that is reduced to make that kind of thing work.

SPEAKER_01

Yeah.

SPEAKER_02

You don't want to push that sort of thing onto every client because number one, they'll be sad, and number two, they won't all do it properly. So you have one stable API, and that's easy to accomplish here when we're talking about, you know, essentially key value pairs that we're editing. That's a the science is is more or less settled on on how to build an API to do that. Uh so you expose that and that's stable, and then everything else on the inside you evolve uh like Kip 500 envisions.

SPEAKER_01

Yeah, um I mean APIs don't have to be key value pairs. They can they can be a little more complex and they can have types. Um I think config I mean configurations probably are key value pairs. But just to give an example, like the API for creating a topic, um, it can have like typed fields like this is the number of replicas or or stuff like that.

SPEAKER_00

Yeah, I think actually the benefit is really like that that semantic awareness, right? So when you're writing to Zookeeper, you do see just a key value API. But when you're writing, you know, you're creating a topic or you're altering a config, you know you actually do have a structured schema and the broker can do validation on it. So, you know, you you you have like basically like this gate. So clients are not able to just kind of blindly insert their their whatever they think should be the state uh into the database. They have to go through this gate. And the gate is where you can control authorization as well as um you know just schema enforcement.

SPEAKER_02

Aaron Powell Fair point. I didn't mean to overly trivialize it. I just mean that it's possible to construct a stable, relatively stable API for Kafka configuration metadata and have that not change even as the internals evolve.

SPEAKER_01

Yeah, I mean I I mean I think you have the rate high-level idea of like having a stable API. I think that when you get down to the details, Jason is absolutely right that you're you're gonna want to see types there. Um and another thing you're gonna want to see is like if clients request a feature that the broker doesn't have, we want to be able to tell them that we don't have that, right? If if they request a feature that the broker is too old to deliver, it's nice to get an error code and say, no, we just can't. With Zookeeper-based APIs, you don't really know, right? I mean, you write your data and you sort of toss it into the void, something happens. Um there's not there's not really a response when you're writing to Zookeeper.

SPEAKER_00

Yeah, I can give you kind of a concrete example of that. So um, you know, we have an API um historically which allowed you to change the replica assignment for a partition or set of partitions. And the way this was done was um by you know writing a Zenode, which says like basically it's like a just a JSON blob of all the reassignments that you want to do. And we expect that the controller is sitting on the other side watching that Zenode. And uh, you know, when whenever there's some um you know new new uh reassignment posted, the controller is responsible for picking it up and executing it. Um but now as far as like the management of the schema for that reassignment, if you want to add new metadata into it, well, how do you do that? How do you actually make kind of the broker aware? Because you know, you don't have that gate in place. So you write that data into Z node, and the controller may or may not understand it, may not may not understand like what your new intent is because you have no way to kind of control access to the you know the current version. Right.

SPEAKER_02

So that's encapsulation and deployability. What is are there any other concerns? Uh we're sort of we've kind of inadvertently stumbled into uh things we want to do better or you know, yeah, dissatisfaction, areas of dissatisfaction with Zookeeper, if I could put it so negatively, even though we have much respect for it as a long-trusted component. What else is in there that has kind of been pushing on the community, making people want to upgrade?

SPEAKER_01

So I think um one area that we can greatly improve on is metadata scalability. And the idea here is you know, Kafka needs to be able to scale to not just thousands of topics, but potentially like millions of topics, right? Down the road. Right.

SPEAKER_02

That's literally several orders of magnitude.

SPEAKER_01

Yeah, exactly. So we want to be able to remove the limits to what we can do. And Zookeeper's design focus is really more on being a configuration store. So Zookeeper was never really designed to be like the metadata store for something like Kafka. So we would like to be able to view metadata as a stream of events that we can keep up with over time and we can cache locally. So I know that um, you know, we spend a lot of time talking about how great Kafka is and how Kafka actually makes it easy for you to follow changes to a topic, right? So you can materialize that view, right? And uh I know Ben Stopford talks a lot about the duality between tables and streams. So you can when you're following a stream, you can construct a table of those key value pairs that you're seeing uh in the stream, and you can sort of materialize that view, and you can materialize it in multiple places, right? Um and this is exactly the same idea that we would like to use inside Kafka. So uh when a broker needs to know about the changes that are happening to the cluster, it should only be able to it should only be able to learn about the deltas, just the changes, right? And not the full state. So currently when when the broker starts up, um the controller has to send it the full state of the cluster. It has to send it every topic that's in the cluster, every partition that's in all of those topics. And you know, we can filter that somewhat potentially, but still, like that's a lot of information and it's only gonna grow. So what we'd really like is for the brokers to mostly have that information and just request what's new, right? We want an event-based system rather than a system that's based on sending around snapshots all the time.

SPEAKER_02

I love that. I uh just as an aside there, I often, you know, when I'm explaining elements of the platform that are not the internal consensus mechanism, I'll talk about how Kafka builds Kafka out of Kafka, you know, like keeping consumer offsets in a topic and um you know, Kinect sort of reusing the consumer group balance protocol for balancing tasks among Connect workers, just things like that. There are these little things where where stuff gets reused. Very judiciously. Schema registry keeps schemas in a topic. You never make stuff up if you don't need it. And now we're doing that on an even more fundamental level. We're turning the metadata consistency protocol into something based on an event stream.

SPEAKER_01

Yeah, I think I think this is a very powerful idea. Like the idea of metadata as a stream of events. And it's something we'd really like to exploit. It had better be a powerful idea because I tell people that's really powerful, a really powerful way to build applications. It's a Kore Kafka idea, and we should we should be taking advantage of it. And you know, there's there's also sort of a further thing as well, which is that the controller has to keep in its memory most of the state of the cluster, if not all of the state of the cluster currently. And Zookeeper also needs to keep around that state as well. Now, Zookeeper, because of its implementation, will also keep that state in memory. Um so there's duplication there, right? And there's also kind of a kind of a split brain phenomenon in the sense that you know there's two states, right? There's the controller state and there's the zookeeper state. Absolutely. In theory, they're never different. In practice, well, they can be. Yeah. And that that can create some very difficult uh difficult to analyze bugs, difficult to analyze situations.

SPEAKER_00

Yeah, and I think it's a it's actually even even a little bit worse than that, because there's there's the state that is in Zookeeper, there's the state that the controller has, and then there's even more the state that all of the individual brokers have. So the controller is kind of responsible for uh you know pushing out metadata updates, but it doesn't necessarily, you know, if if we we get into some kind of uh you know networking hiccup or something, or an unexpected case, those states, you know, any that can diverge at any layer. So where where the broker may think that it has the latest metadata, but the controller you know hasn't sent it actually. Um so basically like there's like this three-level thing where they can all be inconsistent with each other.

SPEAKER_02

Yeah, absolutely. So that's where Kafka has been. That's you know, kind of the remainder of our dissatisfaction with the current approach to consensus, um, which honestly I think is you know another way to put that is that's why we're doing this. Uh we want better encapsulation of metadata inside the cluster with respect to any external concerns like clients. We want it just to be easier to run this dang thing because not everybody is going to use Confluent Cloud in the future forever. You know, there are some people who, for regulatory reasons, have to operate their own clusters. So we'd like it to be more, more operatable. And then metadata scalability. That's a key concern. And it's not, it's, I mean, it's key if it's key, because there's plenty of people who never run into that. You know, they've got their several thousand topics and on the order of 10 to the one partitions in each topic, and everything's kind of fine for their 20-node cluster, and nobody never think about it until you do, right? Until you're really trying to do something big at the level of operating a streaming service for a large enterprise or a cloud Kafka or something like that, where you you just need more metadata than is conventional. So that all makes sense. Is all of that, you know, I think that the title for KIP 500 is simply Zookeeper Removal. Um, is all of that tied up in the KIP? Or what what's the process like? And I don't I don't know if you even have a roadmap you could talk about, but um what should people expect to see from here?

SPEAKER_01

Yeah, so that's a good question. Um so kip 500 is more of a high-level kip that lays out a vision, and we anticipate having follow-on kips, some of which are described in the kip itself. I could break this down into a few different stages, I guess. One stage is is this idea of removing Zookeeper dependencies from the clients. Um as Jason said, that's mostly done, but there's still a few stragglers here and there. There's a few cases where uh we still need to break that dependency from the from the admin client. Or not not the admin client per se, but an admin tool is relying on Zookeeper and needs to be changed to use the admin client. Um so once that's done, uh the second phase would be trying to centralize Zookeeper access in the controller. And what I mean by that is brokers other than the controller should not access Zookeeper. They should actually go through the controller, go through APIs that the controller exposes. Um when we've uh succeeded in uh accomplishing that, we'll have reached a state where we can create what I call the bridge release. And the bridge release is a release that can be upgraded to the zookeeper free uh future. So just to step back a little bit, um normally we support upgrading from any version to any other version. Uh but in this case we actually want to say that if people are running a version older than the bridge release, they have to upgrade to the bridge release before they upgrade to a post-Zookeeper release. Got it.

SPEAKER_02

That makes a lot of sense.

SPEAKER_01

Yeah, yeah. So there's uh there's a lot of reasons for that that are technical reasons. Um we realize there's some inconvenience there, but I think that being able to do the rolling upgrade is actually pretty powerful and and hopefully not too hard to support. Um we have we have some good plans for how to do that.

SPEAKER_00

Yeah, I think it's just like, you know, if you if you kind of imagine that um, you know, as Colin was talking about earlier, when you're in the middle of an upgrade, you've got some brokers which are on an old version, you've got some which are on the new version. Now your ability to like reason about that, some brokers are maybe reaching into Zookeeper, doing arbitrary modifications, um, while some of them are kind of expecting to go through controller, just like the explosion in the number of um, you know, cases to think about and way to protect things, you know, it's it's it's quite you know intractable. So I I think that the bridge release uh is is intended to avoid that. So we're we're only transitioning kind of like to safe points where you know we can kind of limit the number of um you know like that compatibility matrix.

SPEAKER_02

Absolutely. And I think that's fair given the scope of of this change. It's the sort of thing that one needs to plan for. And if that means one needs to go through two upgrade phases, then I don't I don't feel bad politely asking anyone to do that, just given the scope of of what's happening.

SPEAKER_01

Yeah. So once we have the uh the post-Zookeeper release, a few things will be different about that release. So we will have multiple controller nodes, and those controller nodes will form a quorum, which is similar in in many ways to the existing Zookeeper quorum. Um the existing Zookeeper quorum uses uh a variant of Paxos called Zab, Zookeeper Atomic Broadcast. Our quorum will use Raft. Um they're they're very, very similar. Um in both cases, you get uh you get pretty strong, you get strong consistency for writing to that quorum. Um now, because there's a quorum, you will also get the ability to have uh standby controllers. So if the active controller fails, these other nodes in the quorum can take over. And they sort of will function as hot spares in the sense that they'll have everything in memory before they fail over. So failover time should be very quick. And this eliminates one of the one of the metadata scalability problems we have today is that when the controller first starts up, it doesn't know very much. It has to load everything, it has to load its brain from Zookeeper. And of course, that's not a problem if the zookeeper state is pretty small, but as the zookeeper state grows bigger and bigger, um, that loading process can take longer and longer. So eliminating that loading process is actually one of the big scalability wins here. Um yeah, so so we'll have the quorum and uh we'll have the APIs which we added in the previous phase. Um we are we're also gonna have a lot of tools which will replace the Zookeeper modification tools that we have today. So we'll have ways of introspecting into um what's going on with the metadata, because that's very important. And that's a big component of this work as well. Um there's also some changes to uh the broker lifecycle that we'll do. Um basically uh we want to be very clear about when brokers are partitioned from the cluster. And this is another sort of existing problem. So we currently have this flaw where brokers can be partitioned from the controller and yet they're still connected to Zookeeper. And so they stay in the cluster. They're not removed from the cluster as they probably should be, if they can't get any messages from the controller. So that will be fixed because we won't have this dual concept of liveness. And there's some other technical details around the broker state machine, which you can check out in the KIP. Um another thing which I should mention here is as I said before, the brokers will receive metadata updates as deltas. So in the same way that a client fetches uh data from a topic, a Kafka topic, the brokers will fetch metadata from the metadata topic. And they'll be very similar code paths. So the brokers will basically uh, you know, they'll give an offset in the metadata topic, and they'll receive whatever's after that, if there is anything after that. And so this is how we can sort of make that concept of metadata as a stream of events real.

SPEAKER_02

And uh as a as a overly detailed implementation question there, will a broker that is bootstrapping um have to consume from the beginning? And is that the sort of thing where uh there's infinite retention on that stream of events and you just consume it all? Or how do you solve that problem?

SPEAKER_01

Yeah, so that's a very good question. So one thing that I think is clear is that you know you don't want to obviously if if you change something back and forth, back and forth, you don't want to have to consume all of those intermediate states in order to reach the final state. So there's a few different methods that have been proposed for achieving this. One of them is just a compacted topic. Um the other is is sort of a deltas and snapshots approach, more similar to HDFS. So this is an aspect of the design which we haven't completely finalized right now. Uh so I don't want to say too much. But suffice it to say that you know there will be a cutoff uh of the data. So you're not gonna have to go back uh you know, if you have a coffee cluster around for five years, you're not gonna be consuming five-year-old events to reach your latest metadata state. So we'll make sure that there is a way to compact that state somehow so that uh brokers can reach a uh brokers can avoid that. Another thing which I which I should mention here is uh we do want to store this metadata on disk, I think, on each broker, because I think that will actually greatly improve the startup times of brokers. And exactly how we'll stage that remains to be seen. So we may not have this metadata stored on disk in the first version, or we may. Um, but I think that in the long term, as the metadata size expands, this is something that will make total sense to do is to store a local copy of that metadata.

SPEAKER_02

Yeah, and that actually does make a lot of sense. That's by the way, it's it's good to know that. And I think as we move forward and this stuff starts to actually uh get merged and and released, um, that will be a good thing to talk more about because on the one hand, it's a very deep internal implementation. You know, the way the way Kafka achieves distributed consensus is uh the internals of the internals. But that problem is a problem application developers solve all the time. So it is whatever solution you arrive at is going to carry with it a certain authority. Uh, and that'll be a good thing to talk about how you did, because there are people building applications who want to know how to do that, who are persuaded that you know, uh keeping track of of some body of state as a stream of events is the right thing to do. But then how? You know, that that that question comes up. So it's good to know that you're working on that. What else uh you know, what uh what what should we expect to see? Uh, we talked about roadmap a little bit and and sort of the steps, but uh what would be the next thing we would look at?

SPEAKER_00

Um yeah, I can I can mention um you know, as part of this, of course, you know, we're we're you know, the whole point is we're placing Zookaper with a self-managed quorum. So I think um, you know, well the question people are always asking us, you know, are we going to use you know some other library that provides the a similar service or um or something homegrown? And something I actually like uh I like talking about it a lot um because I've been involved in the um the Kafka replication um like the pro the protocol for quite some time now. So and the in the in it's it's had like an interesting evolution over over the years, um, where you know we we had kind of you know Kafka 0.8, the first first release of uh you know log replication, there were some problems, kind of like iterated on it a few times. And I think the interesting thing that came out of it, you know, after after kind of you know seeing a few bugs um that and even after we kind of thought we fixed them and looking at the end result of um how this log replication protocol looks like, it's it's actually you know very very close to to Raft. Um so you know, let me give a I can give a couple specific examples of that. So one of the big changes that we did um in you know, like KIP 101 or something like this was um adding the addition of the leader epoch into the message format. So a leader epoch is basically just an sequential number, which every time there's a new leader, we we bump it. And this gives us a way to detect, you know, the ordering of events. Um and you know, this is also a concept which is in Raft. They call it a term, but it's the same thing. Basically, every time there's a new leader elected, you bump up the term. And all of this data gets into the log, into the individual messages. And what you do in Raft and what you do in Kafka log replication is you use this to reconcile. So in after a new leader election, it's typically the case that the replicas are kind of like behind or are inconsistent with each other. So you use these epochs in order to reconcile. So I think the interesting thing is that um, you know, sort of somewhat independently, we've kind of been converging on something like the Raft protocol already. You know, the Raft protocol, what is it really? It's a protocol for log replication. And what does Kafka do? Well, really, it's log replication. Um so I think the you know, the real, the real kind of like what I think of the potential um here is that um, you know, we can we can build on top of our log replication. You know, the the the main thing that we're missing is is really uh uh the leader election. So I think um, you know, the next thing I think I would I would expect that we'll see um that uh you know I've I've been working on a little bit is uh a proposal um you know for a raft-like protocol inside Kafka. So I think this is going to be something that is just gonna make total sense for the Kafka project itself to maintain. So we don't we don't really want to depend on yet another library because that kind of puts us back in a similar state where we had with Zookeeper. It may not be a separate service, but it has some of the same limitations of you know having this other component that we have to uh depend on. You know, for example, we try to get bug fixes uh in in this other project. I think log replication is just so inherent in what Kafka does that it makes sense for uh Kafka to take this piece. And I think even this is a little bit more longer term, but I think we could even get to the point where you know we're using Raft at the partition level too for uh log replication. Um but I think the main thing is that we really want um Kafka to kind of control this piece.

SPEAKER_02

Yeah, and have it have it really own it as a part of its own code base. Um plus then you get to write a raft implementation, and otherwise you don't. And uh actually I I often make that joke that uh what Zookeeper removal is really about is everybody wants to write their own raft implement. Come on. I mean, you know you do. Nobody's ever argued with me about that. So um but uh that makes all the sense in the world, and there's uh an architectural symmetry there that I think is even aesthetically pleasing to just push the concept of log replication into these other little corners where that hasn't been the solution in the past, but you know, getting it into more and more of the system. That's tantalizing that you suggest uh partition replication being done that way. I mean, that clearly would open up some possibilities in terms of uh multi-data center deployments and things like that, that the behavior there would suddenly be very different.

SPEAKER_01

Yeah, I think Raft is it sort of sits at a very interesting point in the trade-offs between latency and consistency and so on. And I think down the road it it makes total sense to think about it as an option for regular partitions. Absolutely. And and yes, everyone want everyone wants to write their own raft implementation. I won't deny it. Who doesn't?

SPEAKER_00

But uh yeah, I do think um in this case, um, you know, we we do kind of have a head start. Um in look I was like I was kind of alluding to. I think you know some of the work that we've already done. I think there there are some subtle differences between the way that Kafka does log reconciliation versus Raft. Um, but they're really sort of, in my opinion, kind of paper-thin. Um and so there's there's some work that we've already done basically in in the log layer that um can potentially be reused for this implementation. So and then that kind of gets us closer to you know potentially using Raft down the road for partition, but basically having kind of like a unified um layer where you know uh that's you know basically uh all parts are kind of implemented.

SPEAKER_02

Yeah, I have to have state in multiple places. Here's how I have confidence that that state rapidly rapidly converges on uh a single story.

SPEAKER_01

Yeah. Yeah, and I I wanted to say also, I think you commented about this earlier, but but I have a lot of respect for the Zookeeper guys. Like um, I think they did a pretty good job with the project, and I know some of those guys as well. Um, so for us, like as Jason said, it's not really about what's bad about Zookeeper, but it's about what what Kafka should be doing. Like Kafka should be managing logs, right? Kafka should be the best way to manage logs. And Kafka should be, you know, replicating logs in the best way and storing metadata in the best way. So for us, like that's kind of what it's about. It's not um like moving to something that competes with Zookeeper is not like that's not a good trade-off for us. Because Zookeeper is good at what it does, but we feel that uh by integrating this into Kafka, Kafka can be much better at what it does.

SPEAKER_02

Absolutely. And uh uh it's good to make that point. It's like I said in the very beginning, uh, it's been sort of a uh target of lots of abuse, kind of an easy, easy target to pick on. But um, for anyone who has written anything and contributed to open source, you probably aspire to create a project that has been as useful and as uh widely adopted uh and has had such a good run. And to the people who have worked sort of doing what I do in evangelizing and explaining Zookeeper in the community. I remember my introduction to Zookeeper was a talk by uh Camille Fournier, um, who did a great job at that. You know, so there's this really robust history there, and we should not be kicking it and saying good riddance. We should be saying thank you for making it possible for there to have been a Kafka at reasonable cost, and like now it's time to move on to something else. But you know, you don't want to have have to have written that initially, and it's great that it existed.

SPEAKER_01

Yeah, absolutely.

SPEAKER_02

My guests today have been Jason Gustafson and Colin McCabe. Jason and Colin, thanks for being a part of streaming audio.

SPEAKER_00

Thank you, Tim. Thanks.

SPEAKER_02

Hey, you know what you get for listening to the end? A Kafka Summit discount code. Kafka Summit is coming up on September 30th and October 1st in downtown San Francisco, and you can get 30% off if you go to Kafka-summit.org and use the discount code AUDIONET during checkout. Just enter Audio 19 while registering at Kafka-summit.org, and that 30% off is all yours. I'd love to see you there. But hey, I hope this podcast was helpful to you. If you want to discuss it or ask a question, you can always reach out to me at TL Bergland on Twitter. That's T-L-B-E-R-G-L-U-N-D. Or you can leave a comment on a YouTube video or reach out in our community Slack. There's a Slack signup link in the show notes if you want to register there. And while you're at it, please subscribe to our YouTube channel and to this podcast wherever fine podcasts are sold. And if you subscribe through iTunes, be sure to leave us a review there. That helps other people discover the podcast, which is a good thing. Thanks for your support, and we'll see you next time.