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

Hacking Kafka Streams with Sophie Blee‑Goldman | Ep. 15

Confluent Season 2 Episode 15

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

0:00 | 34:33

Tim Berglund talks to Sophie Blee-Goldman (Responsive) about her career in container orchestration and Kafka Streams. Sophie’s first job: interning at Google. Her challenge: helping a hyper-growth customer whose Kafka Streams app was about to hit partition-based scalability limits.


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_01

Today, from interning at Google to hacking Kafka Streams, this is Confluent Developer.

SPEAKER_00

We had a customer who joined us in the early days, and they basically were experiencing, you know, the kind of undound growth that a startup, you know, with a lot of AI customers might see.

SPEAKER_01

I'm gonna ask you what the hardest problem is you've you've worked on since the one. Right.

SPEAKER_00

Okay, no, that's scary. Yeah, exactly. Scary, if that's the right word.

SPEAKER_01

Hello there, everyone. I'm Tim Berglund, and welcome to Confluent Developer, where we explore the fascinating journeys of software developers tackling complex problems. In this episode, I'm interviewing my friend Sophie Blee Goldman. From her start as an intern at Google working on container orchestration to becoming a well-known figure in the Kafka community, Sophie shares a story of the time she and her team hacked Kafka streams in a really creative way to make it easier to scale in the face of Kafka partitioning limits. Now, there's plenty of Kafka and Kafka streams inside baseball in this episode, but Sophie really takes us deep into her solution, and I think she does a fantastic job explaining it all. Let's get to it. Sophie, welcome to the Confluent Developer Podcast.

SPEAKER_00

Yeah, thanks for having me.

SPEAKER_01

Yeah, it's good to see you again.

SPEAKER_00

You too.

SPEAKER_01

All right. Uh, for whoever doesn't know you, who what do you do? Who are you?

SPEAKER_00

Fair question, yeah. Uh my name is Sophie Blee Goldman. Uh some people might know me from the Kafka community, the Kafka Streams World, uh, other people I used to work here at Confluent. Um, and currently uh I'm a software engineer at my current company, which is responsive.

SPEAKER_01

Responsive. What do you do there?

SPEAKER_00

Um well software engineer.

SPEAKER_01

What do you work on?

SPEAKER_00

Yeah, yeah, that's a fair question. Interesting question. Uh so we're we're currently in the middle of a pivot, uh, and I don't want to give away all the exciting details just yet. Um, you know, I'll I'll tease people and say they can go to responsive.dev to sign up for uh you know a mailing list to get news and notifications. Um for now I'll just say we're we're pivoting to something in the observability space. Okay. Uh you know, we have a good team that's built a lot of deep infrastructure uh and has a lot of experience in that. And you know, as a startup, we're kind of always on the lookout for you know what's the what's the exciting opportunity where we can really have the most impact? Uh and we think we found something. So we're currently uh changing directions, but previously before that, we had a uh Kafka Streams platform. And you know, we're still all in favor of Kafka Streams, still love it as a technology, even if that's not our current focus right now.

SPEAKER_01

All right. Awesome. I look forward to hearing more about that. Um what was your first job before you were a software engineer and a prominent person in the Kafka world? What uh what was your first job?

SPEAKER_00

Uh fair question. I guess so. My first job that I was paid was uh was under the table working for my mom in her dean's office, sorting papers. Um not too exciting, you know, but uh things. You were paid by your mom or by the buttons keep that quiet.

SPEAKER_01

No, yeah, I I we know nothing. Uh yeah.

SPEAKER_00

This is not a public space, right? Yeah, um, the internet, but so yeah, my first uh above-ground uh above board job, uh I guess was my internship at Google. Um I know the classic tale, very software engineer. Uh a little embarrassing to look back, but no, it was it was a good time. Uh, you know, they released they spoiled their interns, which is exciting as an intern and a little bit embarrassing looking back, but it was a good time. Um I was I was a little bit uh isolated from most of the interns. I was up in, I think it's called Crittenden, the area. Uh it's like way off the normal Google campus. You have to take a bus to get there. It's kind of near the LinkedIn campus or where it was.

SPEAKER_01

Okay, but it's in the Bay Area.

SPEAKER_00

Yes, yeah, it was in the the Mountain View uh campus. Um but yeah, that's where they stuck all the uh the infrastructure folks. So, you know, all the the deep nerds of the the nerds got they sent them up the hill to live alone and work on uh Docker and containers and all those fun things. Nice.

SPEAKER_01

Yeah, people external to our world don't understand there is a hierarchy of of nerds and uh infrastructure is is pretty high or low, depending on which way you're sorting, in that uh in that ordering. What what kind of stuff were you working on as a business?

SPEAKER_00

Yeah, so they they really thrust me to the deep end uh as you know, an intern uh I must have been like 19 at the time or so. Uh and I was working on containers, container migration, um their system Borg at the time, which basically is the precursor to Kubernetes. All exciting stuff that I I think looking back, I probably didn't really understand fully at the time. Um, but it was very exciting. I was there when they were planning their their DockerCon, which just seemed like a big deal in a whole new world, which I'm sure was very exciting, even though I didn't get to go to it.

SPEAKER_01

Yeah. And is now, and that's a major event now.

SPEAKER_00

Yeah, exactly. So I got to see there when it was uh it was just first getting started, I guess.

SPEAKER_01

Aaron Powell That's kind of cool. And that I guess that was a time when Docker was established, like, okay, this is how we're gonna do containers, but how are we going to orchestrate containers? There was really a struggle for control of that.

SPEAKER_00

Exactly. Yeah. Uh and that, you know, I think probably at the time Kubernetes was at least not in my vocabulary, uh, if in anyone's.

SPEAKER_01

No, it was not clear at that point how that was gonna go. And there were things like Mesos and like what kind of approach are we gonna take to this, and and uh I I think the judgment of history is that Kubernetes won.

SPEAKER_00

Yeah, that seems to be the case. Yes. Uh non-controversial take there. It is weird to look back on what it was like before. Uh and it was, you know, it was all kind of a mess. Uh you know, there was a lot of different thoughts in the area. Um I specifically was working on uh something called container migration. Um it was actually an open source project, which is interesting. Um I have history with open source projects, I guess. Uh it was called CRIU, so the acronym C-R-I-U. Container something with an R. Container something with an R, you know, in user space. Um and basically you you know it, I'm sure. No.

SPEAKER_01

No, I and so I I'm I'm gonna ask you what the hardest problem is you've you've worked on. Is this the one?

SPEAKER_00

Uh no, probably not. I mean, this was hard in a different way. So this was interesting, you know, maybe the technological details were not the interesting of it.

unknown

Right.

SPEAKER_00

Um, but it was it was a weird dynamic because I was, you know, I was at Google on the, you know, basically the Docker team, the infrastructure team, uh working on an open source project with the the major contributors being in Russia. So setting aside current politics, you know, it was this was five, ten, something years ago.

SPEAKER_01

Yeah, they're a different geopolitical. Yeah, yeah.

SPEAKER_00

But still, there's a time difference, right? Uh and so you know, it was it was very hard to make progress as an intern when you're having these open source discussions, which you know tend to be, tend to drag on a lot more, and there's a lot more bike shedding on any open source project uh than compared to you know how move fast, uh break things companies tend to be on the inside. So, anyways, uh there was an interesting dynamic where on the open source side of the project, we were having these long back and forth discussions every PR, every proposal, uh with people that were, you know, basically a full time zone ahead of us. Uh so every single conversation had a relay time of you know 24 hours.

SPEAKER_01

Right, right.

SPEAKER_00

Which was a struggle. Um I think in the end, we we kind of had to drop any sense of real progress on the open source side, and Google had its own internal fork uh where we merged things, and I guess you know, once the interns left, then they would they would clean up the mess, or I'm not really sure uh what happened after we left.

SPEAKER_01

Yeah.

SPEAKER_00

Uh but all I know is you know, uh a lot more of my my PRs were merged into the fork internally to Google uh than the open source project. And so sometimes you can only make progress on the inside.

SPEAKER_01

Yeah, they move forward to their fork. Yeah, I mean collaboration across time zones and agendas and motivations and incentives and everything is is is actually hard to to do.

SPEAKER_00

It is.

SPEAKER_01

Um yeah, and PRs are magical, but they don't make it so that people all think the same thing.

SPEAKER_00

No, that that is very true.

SPEAKER_01

Um so in your career so far, most interesting problem you've solved, what do you think?

SPEAKER_00

Yeah, so um there's a little bit of a tie-in to the the last thing in terms of uh kind of managing expectations and realities around uh open source projects.

SPEAKER_02

Uh-huh.

SPEAKER_00

Um but I'll start with the backstory. Uh probably the the most interesting uh challenge that uh I had was uh IRSponsive when we were still working on our Kafka Streams platform. Um we had a customer who, you know, they they joined us in the early days and they basically were experiencing you know the kind of unbound growth that a startup, you know, with uh a lot of AI customers might see, right? Okay, as you can imagine. Um and so their uh their Kafka Streams application, which was running on our platform, was of course uh constantly growing in its needs in terms of resources and throughput. Um and for those who don't know, you know, with Kafka, everything is divided up into partitions, just to give a quick little overview. Uh and the partitions are basically the maximum limit that you can have on the parallelism. And uh those partitions derive from the consumer groups, and Kafka streams derive from the consumer groups. And so basically with Kafka streams, the most you can scale out is limited by the number of partitions that you have in your upstream topics.

SPEAKER_01

So for the the non-Kafka folks listening, if you've got, say, five uh five partitions in a topic, the topic is like a big log and you split it into five partitions, you can have five processes doing the compute on that. You can't have six. Right. You can have five.

SPEAKER_00

You can have six, but it'll just it'll sit there and it won't have anything to do.

SPEAKER_01

Until KIP932. But that's not what this podcast is about. Yes. Keep going. Go ahead, guys.

SPEAKER_00

Well, as you can see, it's a kind of a fundamental problem. And so a lot of work in Kafka has been done around how to solve this problem. Uh the issue being that, you know, once you create a topic, the partitions, you know, are are set in stone.

SPEAKER_02

So yes.

SPEAKER_00

You know, you don't want to over-partition because you're then you're wasting a lot of resources, perhaps up front. But if you under-partition, then you're gonna hit this maximum limit down the road.

SPEAKER_01

And and you can change it, it's just that you'd rather not. It's not a great way to live.

SPEAKER_00

Very complicated, uh, which maybe will be s saved for another podcast. Yes, yes, yes. Um so uh interestingly, actually, uh the you know, our this customer that we were working with, they had already increased the number of partitions, but it was a huge ordeal, and uh it's not something that they wanted to do again. And so as they were growing and growing, they were kind of coming up we could predict and project outwards, you know, at what point uh is their Kafka Streams application gonna be unable to keep up with the throughput given the number of partitions and the number of processes that can actually be doing active work.

SPEAKER_01

There's your limit, you're gonna run out of gas in you know 30 minutes or whatever.

SPEAKER_00

Yeah, exactly. Luckily we had more than 30 minutes, but you know, we were something on the timeline of like a month, right? So things are going fast. Scary. Yeah, exactly. Scary, if that's the right word. Uh, you know, exciting but scary, like all challenging problems uh start out. Um but yeah, so we had this problem. We had you know pretty much a hard deadline approaching in like the horizon of about a month. Um, and we were trying to solve this problem with fundamentally uh a platform that was on top of open source Kafka streams. So we didn't have our own fork. Unlike Google, we didn't have our own fork of Kafka streams.

SPEAKER_01

Right.

SPEAKER_00

Uh we built something on top of and around Kafka streams that that plugged into um the actual open source code rather than replacing it or anything like that. So the you know, the struggle there is that we were really beholden to what was in open source Kafka streams and what could be in open source Kafka streams.

SPEAKER_01

Right, because you're you at this point, at this point have not forked. Right. That would be uh big decision to do, and that could affect kind of your go-to-market and like adoption. Oh, it's not for Kafka streams. I don't want to, you know. Yeah, big deal. You don't want that. Now a quick word from our sponsor. Confluent Developer the Podcast is brought to you by Confluent Developer the website, which has everything you need as a developer of data streaming systems. And it's completely free. We've got curriculum, hands-on exercises, executable tutorials, the online data streaming engineer certification, also free, a way to find a meetup near you, those are free. Everything is there. I really want you to be successful in your journey as a data streaming engineer, and this is the site that has what you need. Check it out at developer.confluent.io. That's developer.confluent.io. Now back to the show.

SPEAKER_00

And especially because our clients were, you know, really this was this were applications that were running in their environment, so it wasn't the same as like, you know, Google internally using their own fork. This was this would be asking uh our customers to use a fork of Kafka streams, which is a whole other to build against some new dependency and redeploy.

SPEAKER_01

Right. Yeah. That's mean.

SPEAKER_00

Big hassle. Um so that was really kind of off the table. Okay. Um but you know, there wasn't really uh the this this limit of partitioning and maximum parallelism is really like a fundamental limit in Kafka Kafka streams, which derives from the limit of the partitions in Kafka.

SPEAKER_02

Yes.

SPEAKER_00

Uh so there wasn't, you know, at least an obvious thing that we could sort of build on top of it that would necessarily help. But of course the problem is um, you know, I'm I'm a committer in the Kafka streams, uh, which means like we we could have done a um we could have contributed code to Kafka Streams, right? And that was kind of it's naturally our way of thinking, right? We're working with an open source project, we want to do something that will be beneficial to everyone. But the problem was this this one-month time limit.

SPEAKER_01

Yeah, I was gonna say, like you learned in that internship, right? Coordinating that kind of change uh um amongst an open source community, many of whom sort of don't care about your problem, right? They're they got on the motivations. That's tight. I don't see that happening.

SPEAKER_00

Yeah, I mean, you know, and and the reality was like this was uh a dearly needed feature for Kafka streams. Um, and I'll get to what this feature is, of course. Uh so you know, but no matter, you know, even if you have everyone on board with the direction of a proposal, you know, there's just like the realities of the open source uh project itself, which is there's a process, right? So you have to uh draft a proposal, uh, a kip Kafka improvement proposal. Uh you have to publish that, it has to go through discussion and basically have the entire community agree on it line by line, you know, without uh any any little like single sentence of conflict, uh, and then vote for three days. And that whole process, you know, that's usually weeks to months in itself. Yeah. Even for something that's rushed, even for, you know, even for little projects sometimes can take a long time. And for for big ones, especially, even if everyone agrees that it's a necessary uh improvement, how exactly to approach it is is always going to be something that just takes a while to hammer out.

SPEAKER_01

Because you've got this consensus mechanism, you don't have a unitary executive who can say, this is the thing we're doing. That's not how Apache projects are governed.

SPEAKER_00

Yeah.

SPEAKER_01

They're governed by the consensus of a group. Always slower to make decisions that way. Not worse, uh, but it's it's definitely slower. And you didn't you had a month.

SPEAKER_00

Yep. Yeah. Uh-ish. You know, it's it makes for high quality code, but when you're kind of looking at, you know, down the barrel of the gun with one month deadline, that really unfortunately wasn't an option. So now we're at these crossroads where you know we can't or we won't fork the project. Um, we can't fix the thing on time by going through open source Kafka. So we really had to find a way to fix this with our uh our sort of our wrapper, our platform that was built on top of open source Kafka. Um but the problem there is that we're the thing we're fixing is really the like fundamental architecture of how Kafka Streams works, which is that you you know you have a a thread, a stream thread, which is assigned the partitions of uh the Kafka topic, and that's how all the work is done. So really we were talking about kind of rewriting the threading model of how Kafka Streams works in order to be able to work on uh a partition in parallel with multiple threads, which is kind of hard to do from outside of a project, as you can imagine.

SPEAKER_02

Yes.

SPEAKER_00

Um so that I mean that was the real challenge. Uh kind of came down to how do we uh how do we you know fiddle with the way that uh threads are you know used and operated and managed in Kafka streams without changing Kafka streams itself.

SPEAKER_01

And by default, that uh pool of threads is uh not the same as the number of partitions. You've got number of partitions and your maximum number of processes you can scale. Uh number of threads is within one of those processes and they're different things.

SPEAKER_00

Aaron Powell Yeah, yeah. So uh usually, you know, with people that have uh you know correctly allocated their resources uh and you know don't have a lot of unforeseen growth, you might have uh one stream thread which could be processing multiple partitions, but in the default in the default Kafka streams, you you can't have one thread, uh multiple threads working on one partition.

SPEAKER_01

Partition, yeah, that's not how it works.

SPEAKER_00

Yeah, uh, you know, just as the general rule of you know parallelism and concurrency.

SPEAKER_01

Ordering all the assumptions that everything in Kafka makes is just a no.

SPEAKER_00

That's a big one. You know, everything in Kafka is it's a log, so everything is by definition in order, and all the processes are are built on top of that. Uh, for example, like correctness, uh processing guarantees like exactly once, they are built on this assumption that records are processed in a partition in order and they're committed in order and they're flush in order, and all the stuff that comes out of that, not to mention trying to you know write to state stores or you know, uh forward records that are being executed in parallel and all the other concurrency issues that you that arise from that.

SPEAKER_01

Right. And uh just as an aside, at the time of this recording, there's the this big news of a of a thing that's in sort of in preview in Kafka 4.1 that I'm sure you know about it, but audience might not, that is a mechanism to get around that. If you need multiple processes uh consuming from a single partition, more like a classical queuing system, there's this new way to do that. It's called cues for Kafka, kip 932, if I'm getting that number right. And it's all this extra stuff that's taken years and a big kip and lots of work and investment, and that's not gonna happen in a month.

SPEAKER_00

Yeah, and that was a big deal, you know, and that was that was just for uh for kind of plain Kafka and not for Kafka streams. Not even for Kafka streams, right, right. Exactly. So uh, you know, we again we we had to figure out something to do that would we would be basically uh kind of injecting a solution into Kafka streams from the outside. Uh so how do we do that? Um basically we we ended up writing a wrapper. We kind of use our few uh access points. You know, we we looked at the the API that Kafka Streams offers, and we thought, like, where can we actually inject something into the API?

SPEAKER_01

What can you hook?

SPEAKER_00

What can we do? Exactly. Um and so the I mean the main API of Kafka Streams is you write your application, you write uh basically what the the processing topology is, which is just how our record is actually being processed.

SPEAKER_01

Yeah, get stuff from here, do things to them and put them on the phone.

SPEAKER_00

Basically, yeah. Uh you know, you can think of it as like a graph, and the records just flow through the graph. Um, you know, call the topology a graph all the time. Uh so it's the best way to visualize it. And you you know, you write some code to define the the individual nodes in that graph, and then you just write the actual application, uh, which you know, your application is calling things like start and stop. But there's not much going on in the application layer.

SPEAKER_01

But you define your topology, which feels like the code that you're actually writing, but the execution-wise, you've just handed off some data structures to this engine, and then you call start.

SPEAKER_00

Yep, basically. And once you call start, that's when you know all the magic happens. Kafka streams that creates its threads, those you know, spin off and operate separately, uh, and your your sort of main application process. People do different things, but usually it the main application thread is just kind of sitting there and waiting for Kafka streams to you know trigger a shutdown or something like that.

unknown

Yeah.

SPEAKER_00

Not much going on there. So we really had to look at like the stream threads. How could we kind of get into the stream threads and hijack them, for lack of a better word, uh, to do what we want. Um, and really, you know, the the best way or perhaps the only way to do that was going in through the processor topology.

unknown

Okay.

SPEAKER_00

So what the stream threads do is uh you know, each stream thread it gets its task, uh essentially its uh assignment of a partition corresponding to the input topic, and it it grabs records from that uh that partition and then it executes these uh graph nodes, these processors that the user has defined. So we thought, well, a good way to get control over the stream thread is when the stream thread is actually executing these processors.

SPEAKER_01

Okay. So the processors are things if I'm using the API that says like, you know, uh filter and and group and things like that, those are creating processors. There's a graph of these, records are flowing through the graph, and the processor gets a record, does stuff with it, puts it to its output.

SPEAKER_00

Yeah. Kind of most simply a processor is just it's uh an interface and it has um it has you know an open and a close basically methods, but it's one important method which is just process, which just takes it a record and it's a little life cycle and now do your thing. Processor has a process method. I know it's it's crazy.

SPEAKER_01

It's hard to name things, and it seems like that was a win.

SPEAKER_00

Oh yeah. Uh we actually we had to create a new quick aside, a new processor API to fix some problems, and we just couldn't think of a better word than processor. So we had to change the whole package and structure and everything so we could maintain compatibility from the processor, keep processor API. It's other package. Yeah, exactly. Okay. It's just too perfect of a lot of people. Well, it was a it worked. Yeah. That was a controversial kip.

SPEAKER_01

Which processor is this?

SPEAKER_00

Well, it's it's yeah. I won't say it didn't cause any problems.

SPEAKER_01

Say hello to fully qualified class names.

SPEAKER_00

The word processor was it was hard to give up. It's too good. It's too good. Yeah. And so it's processing, it processes these records. Uh, and you know, that's when the stream thread is doing its thing. Stream thread calls, you know, grabs a record, it calls process on you know the first process and then processor, and then they kind of define how it percolates through the graph. But basically it's just calling process a bunch of times.

SPEAKER_01

Okay.

SPEAKER_00

So we thought, well, we'll we can just inject our own implementation. Of a processor using the processor API, basically how you define a topology. And that's going to put a wrapper around the user's processor. So basically we have this little window, this little moment in time between when the stream thread is picking up a record from the partition and handing it to the user's processor. We just we you know tear open a little hole and we injected our code right there, uh right before the user's processor is invoked. And that's kind of how we got control over the stream thread. Um so that's just that's you know, in some sense that was the hard part, but that's also just where it's where it all begins. Um after that came all the actual challenging coding. But uh at a high level, what we did was gave each stream thread its own separate thread pool that it could use for the actual execution of these methods. So rather than uh the stream thread just forwarding the record to the user processor, it would actually stick it on a queue and then move on to the next record without calling process at all.

SPEAKER_01

Okay.

SPEAKER_00

And then the thread pool, uh, we have you know a set number of threads. They're just sitting there, they're picking up records off this queue, and they're the ones calling the actual uh processor. Process method. Yeah, exactly. And they're just they're doing that in a loop.

SPEAKER_01

Um and backing up, this lets you have more of those threads in that secondary thread pool than you had partitions.

unknown

Yes.

SPEAKER_01

Because you're not gonna have more Kafka stream stream threads than you're gonna have partitions. But if you have, I'm gonna guess, more cores on the instance than that, uh, this is how you use those cores. Exactly.

SPEAKER_00

Or even at the kind of the limits, you might have like one instance uh with one, you're like a an instance, the process of Kafka streams for every partition, in which case you would have only one stream thread on that instance. And they're, you know, generally most uh modern machines have more than one core. So there's a lot, you know, there's a lot of work that's kind of being uh neglected or uh picked up. So the idea, the you know, the high-level idea here is really well, let's let's uh get multiple threads working on one partition at the same time, rather than each thread being responsible for uh you know iterating over every record in order for the partition.

SPEAKER_01

So basically, at this limit, you're compute bound. I mean, that this is that's the whole that's how this whole uh surprisingly dramatic story started. I mean, you've got me up like, okay, how how are you gonna do that? And I I kind of get your solution just without with just listening to you without seeing code, you know, that makes sense.

SPEAKER_00

Yeah, I mean that's that's kind of how we came to it, right? It's just you talk through the problem and uh eventually it starts to to take shape. Um so what happened? Yeah. Well, so uh I mean uh I'll cut the uh the drama short and I'll just say that you know in the end it did work out. Uh we we did not manage to run up against this throughput limit and crash and you know have to hide hang our heads in shame or anything like that. Um excellent. But you know, but it was it was a lot of work and there were a lot of interesting technical challenges there. Uh I did a talk on it actually at I think the last the Kafka Summit London 2025. Uh um so if anyone is more interested on the the the technical details, uh they I'll refer them there.

SPEAKER_01

So I I I should have known that, and I'm just gonna confess uh to the audience. I didn't know that, but link in the show notes.

SPEAKER_00

All right, yeah, there you go. Um and yeah, and the code is um it's up, you know, anyone can can take a look at it if they want to, they can try it out. Um and yeah, we call it so we call it our async processor. Uh I know there's there's some controversy over the name as certain other entities think it should be called uh a parallel processor.

SPEAKER_01

Um Yeah, okay.

SPEAKER_00

Which, you know, the because it is it is doing parallel processing, but the way that it does that is by processing uh these records asynchronously within a partition.

SPEAKER_02

Yes.

SPEAKER_00

So we don't have to get into any uh any bike shitting over the names. That's we'll leave that for the kit discussions.

SPEAKER_01

But absolutely. We do not. Um I but I'm just uh like in my mind, I'm going, well, okay, yeah, I get why they're mad because there's threads, and if it was async, there'd be one. And but so I'm I am literally doing that in my head right now. I'm gonna stop. You know, it's it's fine. Uh like we said, we had the the biggest naming win in the world with processor, and right it can't all be like that.

SPEAKER_00

Exactly. I mean we we still stand by our async processor uh but you know I will acknowledge that to some people, you know, they might they might wonder why it's called that. And it's because each each thread is doing uh is processing asynchronously these records.

SPEAKER_01

Um and I want I I gotta take you back into there for a second because I'm I'm thinking of something. Is there an ordering problem in there? Um or oh no, your Q.

SPEAKER_00

Well, I mean, no, that's that is a good question. Um so so yes, I mean there is there's fundamentally like it all comes down to ordering. That's how correctness, that's how processing guarantees, all those good things are guaranteed. Um and there's kind of like one one key observation that we made uh that is generally applicable, even if it might not fit every single case. But usually you have uh a partition which includes records of many different keys. So you might have, you know, uh you might have letters in your key, you know, A, D, G, whatever. They all map to partition one. Yeah.

SPEAKER_01

So you have multiple keys. Big set of keys will map to a single partition.

SPEAKER_00

Yep. You'll have a few different keys in your partition, or many, many of them. Um and fundamentally, you only really care about uh setting aside you know things like committing and uh the the Kafka mechanisms that rely on ordering for a second. We'll go back to those. But in terms of correctness, you only really care about uh maintaining ordering within a key, right? So let's say you're you're like processing the um computing the number of times uh each letter appears in some document.

SPEAKER_02

Yeah.

SPEAKER_00

You you know, it doesn't really matter if you uh you know you process the counts for D, you know, for or for let's say the E and the before you process the T, because they're tracked individually.

SPEAKER_01

And all the as long as all of the E's are in the order that the E's arrive, exactly. It's okay if they interleave with the H's.

SPEAKER_00

Yeah. So you know the So you don't have a guarantee about those. Exactly. The insight there really is just that like while Kafka partitioning kind of gives a natural way to uh break up and and maintain and track things within a single Kafka topic, it's not sort of the fundamental limit of parallelism. It's more it's uh a physical limit because of Kafka, but in terms of like the theoretical limit of what can be processed within a topic or within a partition, it's really um the maximum limit is set by the number of keys that you have, the number of distinct keys.

SPEAKER_01

Okay.

SPEAKER_00

And so that's what we are really trying to hit is to be able to scale up to handle uh you know, potentially up to the number of keys in terms of your parallel parallelism rather than being bound by the number of partitions.

SPEAKER_01

Okay. No, that makes sense. That makes sense. So depending on the cardinality of the keys and threads in the instance, uh cores in the instance, um, that could get that could get very exciting. I see how it all fits together.

SPEAKER_00

Yeah. So our our our little uh our innocent little queue was was not just a simple, you know, first in, first out kind of thing. Uh it really what it did was track um tracked separate queue for each key that we saw. So in that way, we made sure that uh you know records could be could be pulled from any of the uh non-blocked keys. So in other words, if a thread is looking to pick up the next record, um if another thread is currently processing uh you know something with the letter E as the key, another thread can't pick up and process something else with the letter E, because then those would conflict. There would be ordering issues, there might be issues when they tried to uh if they tried to write to the state store at the same time or read from the state store at the same time. You can kind of imagine how keys might uh conflict with each other. Yep.

SPEAKER_01

If you're following along and you were worried about that, because I a little bit was, uh now you know.

SPEAKER_00

Yes. My goodness. So that's that's our you know, that's our kind of correctness guarantee. Um there's also the you know, I said we'll set aside the issue of flushing and committing and the actual mechanisms that Kafka uses to guarantee the correctness and the you know exactly once delivery and all this sort of stuff.

SPEAKER_01

Um you had you had all those problems also to solve or defer.

SPEAKER_00

Yes, exactly. Exactly. So you know, kind of the the fundamental problem there is that Kafka uses, you know, offset commits to track what has been processed and and what hasn't. Uh and that is done per partition. So if I commit offset five, that means that I have finished processing fully every record up to offset five. Which means that uh if we want to kind of maintain these same uh correctness and processing guarantees that Kafka offers sort of natively in Kafka streams by being built on top of Kafka, we had to make sure that when an offset was committed, every single record that was uh processed or sorry, every single record that existed before that offset has been processed. Okay. Um which is hard to do because you have uh a thread pool and everything is working in parallel, and you might have uh you might have you know Kafka streams sort of trigger or call for a commit to happen while you have a bunch of these threads off processing different records.

SPEAKER_01

Things in flight, and you're like not going to be able to do that commit. Yeah.

SPEAKER_00

Exactly. Yeah. So uh, you know, that sort of the fundamental fix there, right, is is pretty obvious, which is just before you commit, you want to wait for everything that's in flight to to finish up.

SPEAKER_01

Um which is fine because you kind of I mean, and and you know, when and how often and everything you commit, and that's a tunable thing, but um you do sort of wait around to do that. And so asking that process to be like, actually just wait a second longer because we're doing something. Exactly.

SPEAKER_00

It's kind of just like a checkpoint. Um and you know, the the the fundamentally like challenging part there is more just how do we do this and how do we hook into the commit. Um, because in Kafka, in open source Kafka streams, it's obvious everything, you know, the the commit is uh a call from a Pi that says finish up what you're doing and then stop and we're gonna commit everything and wait before we move on.

SPEAKER_02

Right.

SPEAKER_00

Um but of course our thread pools are just kind of off to the side. They don't really have a direct alignment communication to Kafka streams because of course we have built this outside of Kafka streams or on top of it. Um so that was another big challenge. Um luckily we you know we knew a little bit about the kind of fundamental implementation of these sort of offset commit and other protocols.

SPEAKER_01

Yes, because you know you're all people who are deeply aware of the underlying stuff. So built some of it. Yeah.

SPEAKER_00

Yeah, yeah. So uh, you know, uh the solution you know presented itself pretty pretty quickly, uh, which is rather just you know, more hijacking. We we injected our own version of the Kafka clients. So we'd have our own custom version of the Kafka consumer and the Kafka producer. Um, I say that as if it's a big deal. Really, they were just wrappers around the sort of default Kafka producer and consumer client. Okay. But what it did was it would intercept any call to commit or flush or things like that.

SPEAKER_01

And wait on the appropriate synchronized things. Exactly. Doesn't that require a change to the client libraries on the client side?

SPEAKER_00

Fortunately, Kafka Streams has thought about this before, uh, and it already has a uh Kafka client supplier config, basically, uh that you can use to inject your own clients.

SPEAKER_01

Okay. So uh config fully qualified uh class name, jar file lying around.

SPEAKER_00

Something like that, yeah.

SPEAKER_01

And you're good. Uh okay.

SPEAKER_00

So that part, you know, was uh was f fairly straightforward uh in the end. And basically, just by waiting to uh to finish up all the in-play records until we commit, we get all of the same correctness and processing guarantees that Kafka offers, which is pretty nice.

SPEAKER_01

Pretty great deal. I you know, it's it's fascinating with this this question of what's the most interesting problem you solve, just the different kinds of answers people take to that. Uh this has been a great deep dive into how this stuff works. It's a really, really interesting problem, and I appreciate you taking through taking us through all of it.

SPEAKER_00

Yeah, thanks.

SPEAKER_01

You bet. My guest today has been Sophie Blee Goldman. Sophie, thanks for being a part of Conflict Developer.

SPEAKER_00

Thank you.