📱

Get Our Mobile App

Take your business learning on the go!

Download on the App StoreGet it on Google Play

Real Time with Redpanda Episode 8: How consumer groups work in Redpanda

Redpanda Data45:16

Transcription

All right, hello and welcome to Real Time with Red Panda, the show about real-time streaming that we're streaming in real time, so you can ask us questions if you need to. This is episode eight. We've kind of got a little bit of a last-minute switch-up here. We're going to be talking about Materialize. We're going to be talking about timely data flow, differential data flow, a lot of interesting stuff. We're still going to have that episode. We're going to push it back a couple weeks and reschedule that, so, so make sure you catch up on that one.

But today, we're going to be talking about consumer groups and how consumer groups work in Red Panda. And to help me out here, I've got a group of great engineers as well. We've got Alex G, a regular host on Real Time with Red Panda. We've also got Noah Watkins, a principal engineer here at Red Panda. So, Alexi, Noah, thanks for joining me today.

Thanks for having me. Alex G, just a fun, fun fact. Um, I met Noah on Twitter. Uh, he actually, maybe it's on Twitter or GitHub, I don't know, something like that. But basically, Noah submitted a thousand-line CMake patch to this, uh, open-source project that I had called Smurfs, this like kernel bypass RPC. And then one day, I see this one-thousand-line patch. I was like, who is this guy? And so it turns out we were friends on Twitter for a while, and then I, you know, I called them up and it's basically love at first sight. And Noah joined us, I think, number four or five. So super excited to have Noah just talk about the, the details of the storage engine and, and, um, our thread record model and our consumer groups. Um, yeah, so it's exciting. Noah, anything else that you wanted to add?

Um, I don't, I don't think so. Uh, that was a, that was a great intro. Thanks.

[Laughter] Awesome. Twitter and GitHub, where all the best development relationships are made. So yeah, I think for C++, also, Twitter is pretty good. So I think that's, that's where we found each other. Yeah. Cool. So, um, yeah, thanks for having us on today. I guess, you know, we're going to talk a lot about consumer groups today, get deep into a lot of different things, including architecture. But maybe at a high level, let's just start off with, what are consumer groups in, in Kafka, in Red Panda? How do they work? Why do we use them? Um, I don't know, what do you want to start on top of that?

Uh, sure. Yeah. So, um, so consumer groups are, are really part of the Kafka protocol. And Red Panda, of course, maintains Kafka compatibility. So when you have clients that are using consumer groups, there's, there's no difference there. So just, that's just to set the stage, right? This isn't like a feature specific to Red Panda. This applies to anybody that's using Kafka. And what consumer groups are, are a way for, one, consumers to deal with parallelism, right? So it's like, if you have a bunch of consumers that are consuming from topics, how do they organize in order to consume from a lot of partitions? Right? One way would be for all of those consumers to, I don't know, like talk to each other and figure out what they're going to do. But that would mean that basically every application in the world would have to invent their own coordination protocol. So what consumer groups do, it's a service provided by Kafka, the protocol, and as well as Red Panda, that allows these consumers to all coordinate and decide in a structured and standard way, which partitions each consumer is going to consume from. So just briefly, like say you have 100 consumers and you have a topic with a thousand partitions, these consumers can now easily connect to a group in the cluster and divvy up this work. So it would be like, say, each consumer is going to choose 10 partitions each to consume from, and you get this data back. So that's the first thing, parallelism and coordination.

The second thing is that it tracks progress, right? So if you have consumers, say you have a 100 consumers, a thousand consumers, I don't know, um, the mean time to failure is, you know, something where you have like consumers failing every day, or I don't know, maybe they're, maybe they're buggy and they're failing all the time. But the, the point is that consumers need to be able to restart their progress, right? So they're consuming from a partition or multiple partitions and they're making progress. And if they crash, where do they resume? We don't want them to have to resume from the very beginning of the partition. So consumer groups also allow, are, so like a, um, sort of like a, a place for consumers to store their progress. So as consumers are reading from partitions, they're periodically committing the offset progress that they've made. And if they happen to crash, they can, they can, uh, come back and then query where they, they've gotten to to resume. So consumer groups are about this parallelism and coordination, and also about tracking progress, uh, to deal with failures. So I think that's like, sort of the high-level nutshell of consumer groups in Kafka.

Yeah, let me, let's, let's add actually a couple more details. Actually, you want to switch to the drawing board real quick? And for those of you in the audience, remember that you can ask, uh, questions on on Twitch and we can talk about the things that are fun and interesting here. The last thing I'll mention, um, um, is we're testing a new, a new drawing app that is collaborative. So please give us feedback. Is the, you know, Keynote on iPad better for, for the drawings? Well, we'll try this and, and let's see how it goes. So remember that, uh, topics are ordered collections. And so, excuse me, on, on ordered collections. And so a topic will typically have, let's say, two, um, let's say in this case, say, uh, topic foo has two partitions. Um, and a partition is an ordered collection. And so a topic is really a named, uh, for, for these two partitions. And so let's say that you have, so let's actually rebrand this piece as producer. So you have two producers, um, producing to, to partition, let's say partition zero and partition one. And so this is, this is on the producer side, okay? So let's talk about the concept that Noah mentioned, uh, in just basically this. Let me give you a visual representation of why this is so, so important. Like, basically, why should this be actually part of the Kafka protocol? And therefore, like, you know, why do people want this, right? And so the, the, the multiple possible universes in, in, in which we need this. As, as part of this collection, well, some, if you're producing, well, let's assume that, that you're also interested in consuming, right? Because otherwise, then you just pipe it to death. Now, and so you have consumers on the other end, um, and let's actually start with one consumer. So say you have consumers on the, on this end, and, um, right now you have one consumer. And so one, one, one thread, the one Python process, it doesn't really matter, right? Like, one consumer will consume these two partitions, uh, in the beginning. And so it was like, okay, let me consume from partition zero, and then let me consume from partition one. Okay, that's great. And then you're making progress. And now, um, uh, this, this, this application, this consumer crashes. And so the checkpointing that Noah happened, it's like, well, actually, what, so what happens? So this thing crashes and, and what do you do? So let's say for simplicity, you put it on a Kubernetes pod and then you say, uh, pod auto-restart policy. And so therefore, every time this thing crashes, it'll move the computer around. Um, but it don't know which offset it actually consumed. That's called the checkpoint. And I was going to talk a great deal about that in the mechanics. I'm just kind of giving you a visual representation of why this is useful. Now, let's talk about the concept of parallelism. So right now, this single consumer is consuming in the union of partition zero and partition one, right? It's basically, uh, and so if you have one partition that, that is is full, then it's probably going to dominate the amount of data that is going to consume in the consumer. And so at some point, the consumer is simply not going to be able to keep up. And so you're going to want, you're going to want to add a new consumer. This is where the complexity starts to happen. So when the world is simple, things are very easy. But what actually happens is when you start to scale? Remember that, that, uh, the concept of Red Panda is actually rather trivial if you have five things, right? So when you show up to a restaurant and you're queuing up in line, and it's, you know, first in, first out order, you can actually write these things in pen and paper, or you could use your tablet. But when you're dealing with hundreds of thousands of events per second, or millions of events per second, like some of our customers, then things get complicated at that particular scale. It's the same thing that happens with consumer groups. So with one consumer, it's easy. Like, okay, well, you consume this state of the world, you know, maybe your checkpoint manually. Well, what happens is this thing can't keep up. So what do I do? I still want the same mental model, but I need to add a new thing. That's when, uh, sort of the notion of consumer group really kicks in. Because what it's saying is that, oh, instead of having, um, let's say this, uh, you know, consumer zero, oops, uh, consumer zero, let's just give this name, consumer A, and consumer B, uh, uh, consumer A consumes from both partitions. You have consumer A, consumer from partition zero, and then consumer B, consumer from partition B. This is the parallelization potential that Noah talked about, right? This clearly is, is like, it's its own thing. It's its own basically failure mode. And then you're going to start to make progress. And hopefully, uh, you can't keep up. Now, I'll mention the last, I think, super, super interesting point, which is availability of your consumers. Let's say for whatever reason that you really care about, uh, the uptime of your consumer applications, maybe you have some, some tight bound SLAs. As part of the consumer groups, you can have this notion of a standby consumer. In the consumer groups, called consumer C. Because there are no more partitions to hand out to consumers, well, consumer C is going to do nothing, but it's there on standby, such that if this thing crashes, there's a protocol built into the, to this consumer group protocol that Noah will talk about, that will basically end up and, and pick up where, where consumer B left off, right? So all of these problems of, uh, coordination, availability, parallelism, concurrency for the consuming applications, basically the things that are actually interested in processing the data, that's what consumer groups is all about. So it, it, it's, it's very challenging to do it at the application. You could do it, but it's very challenging to do it at the application level. Um, uh, and now here's, here's another like, while dimension that no one's going to talk about, which is, okay, let's assume that you have that in two partitions. This thing is still even imaginably easy to understand. What happens is, what if instead of two partitions, uh, you know, one of our customers is trying to launch a hundred thousand partitions, and like a consumer group with 20,000 consumer groups per, per consumer group. So when you get to this kind of a scale, this protocol, it starts to have a, an and non, um, non-trivial overhead. So even this coordination protocol is non-trivial, and things have to change. Anyways, this is the setup, and, and now it's going to start to talk about all of these issues. I'm just trying to give you the context of like, well, what is the problem that this thing is actually solving? Like, why do we have consumer groups in the first place? And why does it have to be embedded into Red Panda proper, uh, to, to give you the primitive so you don't have to solve these things manually? Um, any, any questions?

I mean, I love that. So, you know, I, I basically want to have these stateless consumers. Not, but there's a lot of state associated with them. But now, I, I, because that's sort of implemented in, in the Kafka protocol, I create a consumer, I just do a pull while loop and let it go. Um, and now I don't have to worry about state as much. I just worry about processing my messages. Yeah. Exactly. And this whole latency and rebound and stuff. So now, I don't know if you want to, should we start kind of diving? Why don't you give us a background? I guess, like, what's the challenge then of implementing this protocol on a threadcore model and architecture?

Um, let's see. Well, I mean, this, the challenges that exist implementing this protocol and thread for core are really no different than, uh, the challenges that exist for probably any of the other topics you've discussed on the, the show about thread record, right? So clients are connecting to, um, Red Panda, and they get assigned a particular core. And that core is the core that's handling messages from that consumer. It turns out that, um, in, in a threadbar core model, though, we pin state to different cores. So the core that your consumer is connected to might not be the core that is actually hosting the state for this consumer group. So when that core is handling a consumer connection, um, there's messages that are sent over a cross-core bus, message bus. And this message bus has a capacity. So it can, um, you know, if it fills up, there may be some issues. And we can talk about what happens there and what we've done to scale, uh, some of this stuff. Um, but really, that's, that's the crux of it. Is getting the messages from the core handling the consumer over to the core that actually has the state. Once you're on the core that has the state, we can do all kinds of super efficient batching because we don't have any locks or any, uh, contention for, um, for manipulating the state. And the throughput can be extremely high. Um, so, you know, I think that that's, uh, I think that kind of summarizes the, the challenges related to thread for core.

Yeah. Dear Galaxy, do you want to pipe just the, the drawing real quick so we can, we can see if we can keep up with Noah? So, okay, let me move this thing. I need to give himself a tablet. I think this drawing is very effective. Yeah. So assume that you have, um, I'm sure everyone has seen like the, the meme on, um, on like, basically the Linux kernel where it's like one construction worker who has a shovel and all of the other cores are like, you know, just looking at the one person doing work. It's, it's basically, you know, the consumer group, you know, inward Panda system. You have core zero, core one, two, three. So you have a four-core computer. And, and what, what's happening is that, that let's say that you have a, a cons, actually, let's say, um, a consumer, and then you have a bunch of partition, let's say a bunch of consumers, right? And so some of these consumers, let's give these numbers, consumer zero, one, two, three, consumer four, oops, um, yeah, and then let's say even consumer five. So, wait, hold on. I think I could do this. Let me do this better. Okay. So, are you saying these are all consumers in the same group? Or these are distinct? Yes. Yeah, in the same group. Okay. Yeah, let's talk about this. So, actually, let me make this even more challenging. Let's say that you only have two cores, because and that also highlights the problem. Um, this, this gets, uh, exacerbated the more cores that you have, but the problem is, so you have five consumers in one group. And so right, so this is a single group. Let's call this group A. So you have group A, and, um, let's say half of them land on core one, and then the other will land on core zero. And so what happens is that the consumers that land in core one, it has to make one hop, right? One cross-core, one cross-core, communication hop from core one into core zero. Well, now, what I was saying is that once you land in core zero, because we don't have, you know, we use a thread record architecture, we could do all sorts of efficient things like, uh, you know, debouncing and grouping and batching and, and a bunch of other things like that. And so as the consumer is making progress, so if you scroll up, Alex D, on the, on the partition, on the order partition, remember that each consumer is keeping track of like, okay, which offset, right? I consumed one, two, three, four, five, six. Um, I checkpoint, and then I crash. And so next time I come back up, I want to pick up from from from from offset. So, so it's continuously checkpointing. And so part of the consumer group protocol is basically to save data into into core zero and, and kind of make that happen. So the challenge that Noah is highlighting is this cross-core communication. It's like, how do you debounce? How do you make sure that, you know, you protect the cores from spamming each other? Um, and that's, you know, I think that's basically the fundamental challenge with, with the, with the core architecture. Um, now, do we want to talk about kind of what was the progression once you land into core zero into, uh, the, the, you know, basically the, the serialization format, what ended up happening with underscore underscore consumer offsets, and some of the optimizations that we ended up doing because of a thread core architecture?

Sure. Um, so you're, you're talking about once we have made the cross-core hop and we're on the core that is handling that group? Is that what you're referring to?

Sure. So, um, you know, basically this, you know, the model is very similar to Kafka as well. So we are going to, you know, execute the protocol according to the spec, uh, on that core. Um, and then at some point, we're going to reach a state where we want to checkpoint that data. We're going to do that directly down into a raft group. Um, I think the last stream you had Mihao Masri Lanka on, and he discussed the transition from that raw raft group over into a consumer group, underscore underscore consumer groups. So that's the state. So the state that's maintained on core zero that's updated by the consumers as they execute. So stuff like who are the members, who are the pending members, what is their offsets, all of that data gets, um, serialized using a super efficient, like native format that we have. And then we, we put that into the raft group. And that raft group is distributed across the cluster and protects that data and allows us to recover if one of the brokers that is acting as the consumer group coordinator, if it crashes, we can recover very quickly on a, on another broker. And clients will then automatically reconnect to that broker as part of the, the normal reconfiguration protocol and fault recovery protocol in Kafka.

Wait, so, so then how does, um, a user configure Red Panda to, for two things? One, to make sure that the consumer group, um, has enough, um, because, so if it's a raft group, then it has to have partitions, right? Like, what is the replication factor in that group? And how many partitions should we talk a little bit about how do they think about configuring the cluster? Like, okay, if you lose, then a rep on the broker, you don't lose the state and progress. How do we handle that failure?

I think there's two ways to, there's two dimensions of configuration, right? So the first dimension is how many partitions are you going to put inside of, uh, this consumer groups topic, right? So you may have, um, uh, you know, let's say, as an extreme example, you have one partition in the consumer groups topic. A little background, all of, uh, consumer group coordinators, they follow around or they have affinity for the leader of that raft group of the partition, right? So if you have one partition, you only have one leader, you only have one consumer group coordinator location that's possible. This means that no matter how many raft groups, or sorry, how many, uh, consumer groups you have, they're all going to be co-located on one core in the entire cluster. So you want to choose to have lots of partitions in this topic so that you can spread the load of consumer group affinity out across all the cores in the cluster. So that's the first dimension. The second dimension is replication factor. So each of those partitions needs to be replicated. If you have one replica, you're susceptible to availability issues and data loss. But the more that you have, the more that you have protection of that data, but also the ability to recover quickly. So as soon as a node fails, a new replica is chosen. If you have five replicas, for example, then you have four choices on where you can recover to if a broker happens to fail. So those are the two dimensions you want to think about when you're configuring your consumer group topics.

Gotcha. So partitions help me with scaling out. Replication helps me with recovery. I have one question on just replication generally. Do Kafka consumers ever read from, um, non-leaders, like from replicas? Are they always reading from the leader, and the replicas are just there to fail more quickly?

So by default, every, uh, Kafka client is going to always go through the leader. And this is generally because we want to ensure very strong consistency, right? And so if we're all funneling through the leader, that's, that's how we ensure that. However, you're right, like there are cases where you may want to scale out readers, and maybe they're totally fine not seeing the very last message, right? And so this is, this reading from replica, I believe it's a KIP enhancement, and we are, uh, working on that, but it's not enabled in Red Panda yet. And with that, um, you know, you talked about consistency, and if, if I recall correctly, I read, maybe we're going too far off topic here, but Red Panda, I believe, like acts to all replicas by default, is that true? So like, would there be less of a consistency issue there with Red Panda based on its default behavior? Am I conflating different things a little bit?

So, yeah, X equals all is, you know, the strongest level of consistency that you're going to achieve. And that's somewhat independent of whether or not you choose to read from a replica, right? So you can still have X equals all, but you can have a client that has chosen explicitly to read from a replica, and, you know, they may not provide strong consistency. But again, that's not available yet in Red Panda.

Okay, everybody's getting strong consistency. Cool. Yeah. Well, I think one thing to note, uh, real quick, and I do want to get to static to static membership because, uh, we have to talk about the cool thing that that is actually an active PR right now. Um, when, when internally, when we acknowledge the rights. So remember that there is basically a series of, um, tunable parameters in the Kafka brokers. And, and so, so the tunable parameters for safety. We talked about it extensively on the dips and reports. I'm not going to go into too much of details. Actually, in fact, there's going to be a workshop with Kyle Kingsbury coming up. So save all of these questions of consistency, uh, and, and availability, and, and Kyle, Kyle does an amazing job, and it's going to be super fun. But sorry, just, just to, um, go back to consumer groups. Go ahead.

Just a quick two quick notes on that. Number one, I just dropped a link about that Jepsen and Red Panda webinar. It's gonna be Alex, it's gonna be Kyle Kingsbury talking about this stuff. So check that out. You can read the Jepsen report in chat. Um, one other thing on this chat, we've got a few people in chat, and one of them is really relevant right now. So I want to bring it up. So just noting, uh, you know, something following up to what we're talking about. Kafka has some design inefficiencies. You know, replicating for durability and replicating for high availability should be different. I don't know if that's something you don't want to riff on about that.

Yeah, yeah. So, um, yes, but sorry, first thing is, I need, I do need to give the context. So when we write internally to, to thread record. So at some point, um, this thing called like the consumer group coordinator, which is now, I was talking about affinity. Actually, let's go back to the drawing board so we can talk about it. Sometimes I've got some other questions I want to, I want to queue those up. But yeah, keep going.

All right. So let's see. Let me know where you need to go here. Okay, let's go down below, below that diagram. So assume that you have, um,

[Music]

Let's say you have three computers, um, three computers over at Panda, and they're all talking to each other. Um, and from, and let's, let's pick one to be the coordinator. So let's say, uh, oh, oops. Hey, is the Red Panda, uh, is the coordinator, right? Which means that, um, the computer that is, you know, A, remember that, like, you have, you have a cluster of computers. A is the coordinator for this particular consumer group. This is one I was talking about in terms of trying to understand some, some failure modes of consumer groups. So let's say this crashes, and now you need, um, you need either B or C to, to pick up, like, you know, where, where it goes. So what's, what's important is that, um, we use X equals all to replicate the consumer group state. So that when a consumer group acknowledges the state, the serialized byte format on disk is replicated to all of the other brokers and saved to disk. And so, um, this is, I think that's before the write goes back to the client, is that right?

Correct. That's before. Exactly. Exactly.

Okay. And so it's important because then, uh, and, and so this is slightly different from some, you know, other potential implementations that, you know, may save it to, like an in-memory part of the operating system and then sort of debounce the flashes. Um, in that we guarantee that, uh, you know, if, like, if we acknowledge that right, then it actually got written and persistent to disk. And so then what's cool about this is because it's using the, the raft leadership algorithm. And now, I think you, you can talk a little bit about how do we go about, you know, basically stopping the leadership of the coordinator, let's say a node crashes, and then like, how do we choose one? It's basically around this node affinity thing. But then the data is already on node B and C. And so we can say, okay, node B is now the new coordinator. Um, and then it'll just read the, the log that that was written, you know, localhost. And then there's this protocol for making sure that all of this, uh, all of these consumers, basically consume back and change their, their coordinator. Let's say on core 52 of computer B. Does that make sense? There's like, there's a lot of dimensions here that I think are worth highlighting. But, um, yeah, maybe Noah, you want to talk a little bit about that, like this coordination, why it's important to, to data to disk and so on?

Uh, yeah, so when, well, we need to persist data to disk so that we can recover. Yeah. Uh, you know, that's, that's the most important thing. Um, and, you know, when, when a fault does occur, um, raft itself knows how to, uh, deal with this situation. If a node goes down, there's a, there's a timeout that happens, and raft will elect a new leader. Now, assuming that there, you've, you know, decided, you've configured your system so that you have enough replicas, raft will very quickly decide on a replacement broker to act as the new coordinator for that raft group. And it'll go ahead and replay that log, like you mentioned. And that replay is very fast because it should already be almost already fully up to date. So it may just be replaying a few little bits and pieces right at the end of the log. But as soon as it does, um, it's going to have the recovered all the state and in memory. And then it's going to announce to the cluster, hey, I'm the new leader for this raft group. And at that moment, clients or consumers that are connecting the cluster are going to discover the location, um, very, like instantaneously, discover the new location of the leader and immediately start talking to the, um, replacement replica acting, uh, as the, the group coordinator.

Yeah, but I guess it's it's unclear to me and and I also want to answer Renee's question. Renee, it's good to see you here on the on the Twitch channel. Um, before we answer my next question, I think it's useful to to learn, I think for people in the audience, like, well, how do we actually, why is it then to me, I feel like the the leadership, coordination protocol, right? So like you have the heartbeat, the heartbeat expires, there's a new round of voting, and then you pick a new machine that's part of raft. But I think what's key for people to understand is like, how did we attach that to the partition coordination? Like, why is that true? And we basically just said, here, we made it so. But, you know, how does, how does the consumer group know he's like, hey, connect to this other computer? This other computer is is the new, the new leader. How does that work?

It's baked into the Kafka protocol. So there's really kind of like five key APIs, or maybe if you wanted to call it seven, but there's five key APIs for the core, uh, consumer group API, right? And one of those is the, is called find coordinator. And this is an API that a client can invoke against Red Panda, and it'll say, hey, I would really like to know which broker in the cluster is acting as the coordinator for for my special group. And any node in the cluster can answer this question. And so it will compute a hash of the, um, of the consumer group name, and it will, you know, then there's like a deterministic mapping from this hash across all the cores or the brokers in the cluster, and it'll just return back this answer to the client. And the client will then just immediately stop what it was doing and start directing all of its attention over to a new, new broker.

Yeah, and so just to highlight there, it basically, what happened is we transferred the leadership. And when we answered the next API, because the computer crashed. So it, so, so the consumer groups will reach out to all of the brokers that, one of the brokers that it has discovered, and it's like, okay, who's the leader for this partition? And then because internally we know who the leader is, then we'll respond. This is the new leader for that partition. But semantically, what it means is that we've attached the leader of the partition to the consumer group coordinator. That's that's basically the, the, the tie. And so we get it for free, right? Like, you can imagine multiple things, but what we said is, we already have this leadership transfer mechanism. It already rebalances across the cluster. We already have this partition count and blah, blah, blah, replicas and all of this. So we can recover correctly, deterministically. So all we have to do is figure out a map in between the consumer group coordinator and the partition leader. If we make them the same, then, then it's easy, right? Like, sort of the mental model for, for, for the person. Um, um, okay, I think that we have another question on the on Twitch. Um, uh, in sync replicas, is this one? Yeah. So, so, um, uh, does, uh, Red Panda use min in sync replicas?

Red Panda does not. We are a quorum system. And so this is actually explained in detail in the Jepsen report. Actually, was it Kyle when he did the Jepsen report? No, correct me if I'm wrong. Like a few years ago, wasn't Kyle the one that that suggested changing the min in sync replicas? It's like a change that caused an original Jepsen test on Kafka that caused two things, you know, leader election, a min in sync replicas to be a tunable. Is that, is that how I was saying?

Oh, it was a long time ago. I think that was one of the very first systems that Jepsen was, uh, testing against. Uh, I don't, I don't quite remember. But I do remember ISR having some kind of funny funniness in it. Yeah. Yeah. I don't know. So that's a good question to ask Kyle. Go ahead, Alex.

So, so just, just so I understand, in-sync replica is basically saying, hey, I, you know, I want to have this many rep, or if I have a partition, each one has some rough, because I don't want them to get too far behind. And at that point, if some of them do, my my leader will stop accepting rights. If that happens, is that right?

In, in basic, it'll, it'll block. So when you write, the partition leader will block until the minimum in-sync replicas have acknowledged that that right to the storage subsystem. It doesn't mean that it was written to disk, by the way. I think I think that's a caveat that people need to understand. It means it was written to the storage subsystem, and the storage subsystem is this totally separate subsystem that has its own tunable parameters. One is like, hey, you could write it to the page cache, and like, it's super fast, or you can actually durably write it to disk and have a different level of durability, right? And so I think that's when, when some people were thinking about durability and availability, right? Like, it's, I think it's very detailed, and computers are really fast, and I think that we are on the side of being safe these days. It's fundamentally different if we were operating with the spin and disk, but NVMEs are so fast. So anyways, this kind of philosophical difference here. But, um, but yeah, so that's what it means. It means that, uh, the partition leader accepts the right, then it, then it forwards it on to, actually, it doesn't forward it, it waits for the, for the ISR loop on the other replicas to pull the Kafka, because polling is like relatively efficient as opposed to a reactive system. And then when the min replicas act back on to the partition leader, it says, hey, I've written it to a storage subsystem. Then the producer is going to say, like, okay, here's the ACK. And, and so that's, that's the difference. And we, we don't do that. We are a quorum system. We are a raft, so we just simply ignore those properties. I mean, I think that's sort of, um, it's, it's, it's very easy to reason about a system when I can give you a math proof that says, this is what it means to have three out of five replicas up and running. This is what it means to have four out of five, or one out of five. Like, you understand either what it means to have unavailability or, um, you know, consistency. And, and which is to, to use a relatively, uh, uh, uh, consistent protocol for our replication, kind of is a rat hole. So let's pop it back here because I want to talk about static group membership.

Yeah, sure. Let's do it. Do you need the, the drawing out? Yeah, let's talk about the, so if you scroll up just a bit onto where I wrote the hundred thousand. Yeah, there we go. The original drawing. So, um, all right. So there is a ton of chatter. Noah, do you know, do we have any idea how to quantify the chatter? I guess in terms of bytes or or like hardware talking about Twitter? That's what I thought.

No, no, no. Um, I don't know how to quantify. So there's a lot of, you know, once you get through the initial round of the protocol, and I'm sure you're going to, you're going to get into this, right? Then it's a lot of chatter for heartbeats. And so it's, you know, it's some sort of mapping like heartbeat per second times size of the message. Message size is pretty small, but, you know, it adds up quickly. Yeah. And so assume that you have, you know, whatever, 10, 10 tens of thousands of of of consumers, this is this is a relatively, uh, chatty protocol. And, and I think what Renee was actually suggesting on his comment, like, oh, consumer group rebalance is well, when you have a really large group, it takes a while to to rebalance. And so to overcome this, I don't know, is this like a design limitation? Probably on the original design. Now, there is this notion of static group membership. Uh, and actually, Noah, you should talk about it because you're the expert here. Uh, we should, yeah, we should have Sri Lanka on. He's just put up a PR to to add static group membership into to Red Panda. So he's the expert. But basically, the issue is that when consumers are coming and going out of a consumer group, right? So it's like, you have a quiescent state, a new consumer arrives, or maybe a consumer crashes and its heartbeat times out, and the consumer group kicks that member out of the group. All of these cases trigger something called group rebalance, where the group all decides to, uh, re, um, recalculate who should be consuming from which partitions, right? And when you start to get up into, um, consumer group sizes where there's a lot of members, or even not a lot of members, but you have flaky consumers, then you spend a lot of time engaged in this rebalancing. And so what static group membership does is says, hey, you know what, we're going to pre-configure what this, uh, balancing affinity should be, which consumers are consuming from which partitions. And this means that we can forego a lot of this rebalancing, right? So it's like, if a consumer leaves or gets kicked out, then there's temporarily going to be this, um, this time when this nobody's consuming from whatever had been assigned to the, that consumer that left. But if it joins again, or some other new process takes its place, um, it's just going to take on, uh, the role of that lost consumer. And, and you're able to avoid some of that rebalancing overhead. So that is static group membership. The static is sort of referring to that, um, static assignments.

Interesting. And so is that for folks that have much larger clusters, huge consumer groups, things like that? They they'll say, hey, I'll take on a little bit of that burden of, you know, making sure my consumer comes back up and, and, um, doesn't fall too far behind, but they avoid all that rebalancing. That's what's going on there.

Yeah, I think it's important for for clusters where you, or sorry, for consumer groups where you have a large number of consumers. Um, but it's also important for situations where you have a lot of flake. You know, we might have like a much smaller set of consumers, but they're flaky, right? And so they're like, they're leaving and going, rebalancing.

Yeah, exactly. So what are the timeouts there that I think are useful for people to to know? And also, can you make a comment on underscore underscore consumer offsets? Not necessarily from an API perspective, but I think there's this misconception in the community that underscore underscore consumer offsets, it's like an API. And it's actually just for ecosystem compatibility with tools like LinkedIn Burrow and a bunch of other tools. But anyways, before we get into that, let's talk about like, what are some of this like, why, why do people care about this rebalancing? Like, what's the time impact? What is the latency? So let's say you have a trading application, and you have to, you know, be up and running within an SLA of, I'm making this up, like five seconds, right? So like, what are the potential failure modes that will cause a user to violate a particular like SLA? Like, why would something take long? Like, what, what are the, what are the points in the protocol that would basically cause an infinite timeout loop, which is which is often kind of hard to debug?

Um, I mean, I guess the, the easiest to sort of understand is that when consumers initially join into it, or sort of like arrive at a group, right? There's a debounce period, um, where this rebalancing step doesn't occur. Right? So if you have like 50 consumers that join very quickly together, like sort of within a few seconds, then that, there is no, um, rebalancing, even though like fundamentally they're sort of joining in serial. But if a 51st consumer joins, like right after that debounce period, then you have this rebalancing problem, right? So there's a tuning parameter there related to debouncing that's important. Um, and then you can have a lot of failure modes if you are in a situation where your network is a little flaky, or there's some, you know, there's a lot of, there's a host of things that could, uh, cause consumers to miss those, those, um, yeah, basically it's like a tumbling window thing where, like, let's say you have a window of time, and then something like, and this, this will happen in production. So for everyone tuning in, this is just a failure mode of the Kafka protocol. Uh, you know, so you have to understand this. Um, so let's say that you have, you have a window of time, and then you sort of, uh, come in, and like, right at the time of intersection, you trigger another timeout or rebalance effect, then it just extends the window again, and it continues to extend them. In the right. So like, you have to have, uh, you know, quiescent points, basically points where the system has reached some form of stability to make progress. Otherwise, from the point of view of the coordinator, you can't make progress. I, like, the protocol tells you, please wait this timeout, and it's like configurable. I'm going to make this up, like 10 milliseconds, or whatever it is. So if something keeps continuing to happen, let's say at the 7th millisecond, or 8th millisecond, or whatever it is, like, you're just going to have this thing, and you're like, uh, this, this, this is terrible. Like, wow, you know, you feel like Hulk is smashing your keyboard because nothing is making progress, and you're seeing all these thoughts. It's like, okay, well, something is broken. I don't understand what it is. And it's like, there's this critical timeout that will simply extend the, the point of which your system, you know, it's like, just looks like it's halted. And, and it's, and it's a very sad production environment. And of course, this happens when you push a bad configuration into production in a particular app. So it's like, you write an app, right? And so you're a trading application, and you say, like, oh, look, I need this SLA. And so you're hacking in your application, things are great. And then you happen to introduce, let's say, like, you have a bad packet, or, or like, whatever, bad something or other, and it just, it's just always, I feel like production systems will always find these particular edge cases. It's, it's kind of mind-bending the, the amount of like correlated failures that you have to have, but it will happen in production. And at some point, you're just like, I have no idea what is going on, but everything is crashing. Um, and it's very frustrating. And it's just part of the consumer group protocol. Like, people need to understand that this debouncing period is there to save on different edge cases. But this, this domain is a tunable parameter that people can set, which is how long do you wait? The longer you make it, the more, uh, quiet, like the fastest time that it takes to converge, because you're just like, giving people a lot of slack. The problem is, you're expanding your failure a window. And so at some point, it's just like, I have no idea what's going on, basically. It's, it's rough. Do I get that right?

Yeah, that was an excellent, uh, description of some of the more, some of the nuances. Yeah. I know we're coming up at the top of the hour pretty soon. Alex D. Yep. So, yeah, thanks for, thanks for doing this. We, we probably have enough content to do another consumer group or, or something else somewhat at the next time, but this was, this was great. Noah, thanks for, thanks for joining us today. Thanks for, you know, submitting that thousand-line PR a couple years ago. You met Alex and joined the team. Alex G, great to have you as always. I just want to call out again to everyone, May 25th, put it on your calendar. There's going to be a live webinar with Kyle Kingsbury, Jepsen, Alex G, talking about, you know, the Jepsen report that went through Red Panda. I'd say go look at that report, but then also go through in that webinar. They're going to walk through a lot of that stuff. They'll have some live Q&A. So, uh, check that out. Should be an interesting thing. I, I thought it was super interesting to see all the stuff that came out of, out of the Jepsen report. So anyway, Noah, Alex, thanks for, thanks for being here and thanks everyone for watching. We'll see you next time.

[Applause]