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

How to Convert Python Batch Jobs into Kafka Streams Applications with Rishi Dhanaraj

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

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

0:00 | 31:02

Zenreach is a company that makes tools to help retailers use digital marketing more effectively. If that sounds like a problem that only marketing people would be interested in, that’s because you don’t know what they do! There are all kinds of fascinating technology problems to solve by utilizing event streaming platforms to process data at volume. Rishi Dhanaraj, our guest today, worked at Zenreach as an intern, and took on a big pile of Python batch jobs, turning them into some really interesting Kafka Streams code. Listen in as he walks us through how he did it.

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_00

ZenReach is a company that makes tools to help retailers use digital marketing more effectively. Now, if that doesn't sound interesting to you, that's only because you don't know the kinds of problems they have to solve. Rishi Donaraj took on a big one while he was working as an intern there. Did some really cool things with Kafka Streams in the process. We'll dig into them on today's episode of Streaming Audio, a podcast about Kafka, Confluent, and the cloud.

SPEAKER_01

Thanks, Tim.

SPEAKER_00

You bet. Um, hey, tell us about yourself and uh what uh what you're up to.

SPEAKER_01

Sure. I'm a third-year computer science major at UC San Diego. And this fall I interned at ZenReach, which is this awesome startup in San Francisco. And there I got to rewrite their walkthrough detection system with Kafka streams. And I learned a lot there, and I wrote this blog post. I think that's why I'm here today.

SPEAKER_00

You did write a blog post. And just so you know, uh, yeah, everybody, we're recording this kind of in the middle of April uh 2019, and that blog post was uh earlier in the year, and it was um we won't get into a long discussion of confluent blog statistics because I think that's probably not exactly what everybody's really interested in. But it was one of the most popular posts in the first quarter. So uh it was really cool. And it's cool because it's a it's uh people in the real world doing things, you know. A lot of the time, particularly um, you know, in the developer relations team here, we spend our time talking about how to use stuff. And you know, we contrive examples to come up with good teaching scenarios. So it's always great, uh, Rishi, to talk to people like you who build stuff in the real world. Tell us about what ZenReach does. What if people don't know the company, uh what what's their business?

SPEAKER_01

Right. So ZenReach creates products to enable brick and mortar merchants to better understand, engage, serve their customers. And it all centers around our ability to quickly detect when a patron walks in and out of the store. And we usually do that through Wi-Fi. So we would have our customers would have routers. And by detecting when a customer is connected to or no longer connected to the router, we'd be able to tell whether there's a walk in or out of the store.

SPEAKER_00

Your customers are the stores, right?

SPEAKER_01

Yeah, our customers are the stores, and their customers are patrons.

SPEAKER_00

Or their patrons. Gotcha. Those are the words. Right. So if I run um a chain of artisanal goji berry shops, uh and I want to know things about uh you know who walks in and what they do and where they go, then then I would I would do business with ZenReach. I would I would buy the product. Exactly. And you said that works by Wi-Fi?

SPEAKER_01

Yeah, it does.

SPEAKER_00

Yeah. Can you I don't know what I don't know what of that is secret sauce and and what you can talk about, but that's sort of interesting. How is it is there anything you can describe about how that system works? Probably not.

SPEAKER_01

Yeah, not necessarily, but it's totally cool.

SPEAKER_00

So there are uh there are these radio things, and those radio things are around, and people have little radio transmitters they carry in their pockets willingly, and um you know somehow there's some science uh and some math and some computers that uh can give me in my artisanal goji berry shop information about uh who comes in and where they go. It used to be uh a little bell on the door in the olden days, and maybe a mirror somewhere and a security camera, but now we can use uh science and computers and magic to do that, which is pretty cool. So that's uh that's what Zenriach does. So what um uh certainly I smell events there, right? Uh that that sounds like a system that deals with events. So talk us through what you built. And by and by the way, um as a project an intern is working on, this sounds amazing. So uh you really scored there. Uh this is this really cool. Yeah. Uh so talk to us, and you know, to some degree, you you may be just talking about what was in the blog post, but I would love to go through that. So, what uh what did you build?

SPEAKER_01

Sure. Yeah, so thanks to Zenrich for giving me a lot of responsibility first. And in terms of what I built and I guess what the events and this architecture were, so the core of it is one event is when a customer walks into a store. And that's so that was something we'd call a walk-in. Okay. And that's something important. And the second kind of event is something we'd call a Zen Reach message. So this is basically whenever we have some kind of online engagement with the patron. Maybe we send them a marketing email, maybe we send them an SMS. So these are the two events. And we want to somehow connect Zen Reach messages to walk-ins, which is Zen reach messages potentially cause these walk-ins, and we want to try try and somehow figure that out. Did one of these cause the other?

SPEAKER_00

Okay, yeah. You want to attribute you're you're doing this, you know, marketing, communication, messaging people, uh, and and can you attribute uh your reaching out to the customer or prospect with them actually coming in? That would be a good thing to be able to tell me, my goji berry, my goji berry shop, artisanal goji berry shop, is that um you know what what you're doing, your your pinging of my prospects is actually resulting in them walking in.

SPEAKER_01

Exactly. Exactly. And yeah, so the exact term that we use to call this is a walkthrough. So that's the the number of walk-ins that happen as a result of one of these ZenReach messages. Got it, got it. Right. So previously, our walkthrough counts that we would have a Python script back from the early days of the company, just do this more crudely, where it would periodically run read walk-in data from Cassandra and read ZenReach message data from MongoDB, and then write walkthrough counts to MongoDB. And from that, all of our other Zenrich services would query MongoDB for that data. But the problem was that that Python script was on just one server, and our goal was to move to Kafka streams so it could run in a distributed banner. And we also wanted to reduce the load or the read pressure on MongoDB on Cassandra. And for those reasons, we wanted to switch to Kafka streams. But everything was around making sure we can generate walkthrough counts from walk-ins and from Zenriach messages.

SPEAKER_00

Got it. That's the that's really the primary business metric, or at least for this part of the business, that's the primary metric.

SPEAKER_01

Exactly.

SPEAKER_00

How many walkthroughs are happening as a result of a message. Um, which is another way of describing the effectiveness of the marketing campaign. Um which your uh which you know, stores, your customers are are are paying for that service. That makes sense. Okay. So uh what was that like? I mean, I guess there's there's uh there's two things I'd like to talk about. One is um what what's the stream processing topology like? Like what what are the operations you performed to the extent that you can go into detail there? Um and then for you as a newcomer to Kafka Streams, and if you were a newcomer to Kafka, uh you could tell us that. What was that process like? I imagine there was some things that there were some things that were pleasant and some things that were uh friction, and it's always good to talk about those. So two kind of two goals. One is as much as you can talk about stream processing, and then we'll uh walk through what it felt like.

SPEAKER_01

Sure, yeah. So I'll talk to the first one, which is you know, what was the topology, what was the structure of the architecture? So the this specific project we kind of split into two services. So one, we had a walkthroughs generator service that would take in a walk-ins topic and a Zenriach messages topic, and then output a walkthroughs topic. And then we had another service called walkthroughs Mongo Exporter that would take that walkthroughs topic and basically record that in MongoDB. So the second one is not necessarily that interesting. But to focus on the first one, the walkthroughs generator. Right. That one was interesting because I think it uses a system that a lot of that seems a bit ad hoc, but it seems like other users of Kafka streams on Stack Overflow also seem to be using it. And it seems to be an ad hoc kind of batch processing system where you store the incoming events into a state store and go ahead and process those. So for us specifically, we had a message processor that would wait for incoming Zenriach messages, and we would store them in a state store and then process them. So it wasn't necessarily that we we were able to process them just as they came in because potentially events could arrive out of order, and specifically messages would for sure all the time arrive out of order. And as a result, we weren't able to process them immediately, which is I think ties to one of the first kind of lessons I learned as a beginner to Kafka streams is that you there are some real-world constraints that you know make it make writing this code a little bit more icky. And so as a result, we did have to use this kind of ad half batch processing instead of just being able to use you know some kind of k-table that would automatically update. Um and so yeah, to back up, um that's that's one lesson. But that that's just the first part of the walkthroughs generator service. You have we have this message processor that takes in ZenReach messages, updates this ZenReach message state store. The second part of it is that there is a walk-in transformer that takes in walk-ins and reads to from that Zenreach message state store. Also goes ahead and reads from a walk-in state store, which is what it's updating when walk-ins come in. And basically by reading from the ZenReach message and walk-in state stores, we're able to generate walkthroughs. Uh and basically what we do is in both of those, the message processor and the walk-in transformer, we would go ahead and update the state stores, and then we'd periodically run punctuate in walk-in transformer to actually generate walkthroughs.

SPEAKER_00

Gotcha, gotcha. Yeah. How I'm curious, how late could the data arrive? It was uh walk-in data that you put into a state store. So you like sort of um uh I guess you said ad hoc, you built this ad hoc K-table-like thing, but how late arriving could the data be?

SPEAKER_01

I can't remember off the top of my head right now, but it was enough so that we couldn't just rely on Kafka streams to wait for it. It would it could potentially arrive quite late. And we would have to send a correction event afterwards. So it wasn't necessarily something simple.

SPEAKER_00

Gotcha, gotcha. That'd be uh interesting to dive into in more detail at some point because it it seems like uh it seems like there are uh facilities in the API, you know, windowing facilities in the API that uh will accommodate late arriving data, but uh that was not that. Right. Oh, okay.

SPEAKER_01

So to clarify, so this late arriving data, it one of the things is that the timestamp that Kafka Streams would be put would put on it would might be also is that is not necessarily the time would be used that would be used for calculating whether this was when the walk-in or Zen reach message actually happened. So for example, um Wi-Fi routers would give us different Wi-Fi routers or ZenReach messages, they might have different timestamps we could potentially use. And sometimes we wouldn't be necessarily using the Kafka streams timestamps. So as a result, we would have to, we wouldn't be able to rely purely on just windowing. Got it, got it.

SPEAKER_00

And that timestamp was available in the message?

SPEAKER_01

In the message, yes.

SPEAKER_00

Got it, got it. Okay, so you built this state store and then you had um I mean used the state store API uh ingesting those are the well, I forget now. Are those the messages or the walk-ins that are in the state store?

SPEAKER_01

So there's two state stores. One of them holds the Zenerge messages, the other holds the walk-ins.

SPEAKER_00

Got it, got it. Yeah. Um and what is the operation that you do? And and and again, and I any of these questions I feel like I'm potentially asking secret sauce things, so obviously, you know, just tell me. But uh then what do you do? You have those two state stores. Uh, what constitutes a walkthrough, and where in the API or or where where in the topology uh is there an event that you respond to that that lets you do that computation.

SPEAKER_01

Right. So the major sort of things we use, we would first we would re-key it. Um because the the way it's done right now is in reach messages and walk-ins are keyed by arbitrary things for different parts of the system. So first we would we would re-key, then we would go ahead and filter because some walk-ins, some Zen reach messages aren't even necessarily relevant to to a walkthrough conversion.

SPEAKER_00

Gotcha.

SPEAKER_01

And then those, and then we would also uh map some of those values to different values just for the for ease of comparison. But then I I guess where stuff is actually compared would be in the walk-in transformer. And in there, what would run periodically is that is punctuate and would go through some of the the walk-ins and state stores that have accumulated and you know, look at the timestamps of the Zen Reach messages in the message processor, look at a few other things, and would decide whether something was a walkthrough or not. And also would have to decide which of the events would be um, which of the ZenReach messages should be tied to this walkthrough, which also requires some thinking. So I I guess the the filtering, the select key, that kind of stuff is that kind of d that's that stuff that's used with the DSL is pretty straightforward. And then inside the walk-in transformer in punctuate, we would kind of loop over the relevant loop over the relevant walk-ins and see um and run some algorithm to see if Zenrich messages would cause them.

SPEAKER_00

Got it. Got it. Um yeah, and the re-keying, by the way, is um is uh par for the course. Yeah, you're uh it it it seems like topics never have the right key. Yeah. Um but that's one of the nice things about them being filled with immutable data is that you can just make copies of it and uh make a topic with the right key. Exactly. And yeah, yeah, yeah. I guess in this case a stream with the right key. Uh cool. And so then in that punctuate, you made the uh you know, you you did that kind of secret sauce calculation and emitted that into a new topic.

SPEAKER_01

Right. Yeah, the walkthroughs topic, yeah.

SPEAKER_00

Got it. Um and that you said ended up in Mongo. And was that with Kafka Connect, or did you just have a consumer uh consuming that and writing to Mongo?

SPEAKER_01

We did not use Kafka Connect because at the time Kafka Connect didn't support Pro Debuff payload. And also we didn't necessarily want to provision additional infrastructure. So we just wrote that ourselves.

SPEAKER_00

Okay. Yeah. Was that in the Streams app, or was that you write to a topic and then have another another consumer process that does that?

SPEAKER_01

Uh we the the latter.

SPEAKER_00

Yeah, cool. Cool. Um now what um what what kind of code? You're a you're a third-year computer science student when you got this internship last semester, which again hats off to you. Uh what kind of code had you written before this? Like what was your software development experience going into this?

SPEAKER_01

So I suppose well, I had had um two internships beforehand, but they were more in you know, front-end web development, more more front-end or more product development, client development, that kind of thing. But this is the first time I had delved fully into some kind of back-end where I had to, you know, actually understand architecture, or actually, I think that actually had to understand something deeply. Whereas previously, you know, if you're reading some kind of front-end documentation, you can kind of just read what it says once and you can kind of understand how to use whatever library. But I think it was different back when you had to think about overall architecture, you had to debate the design, which is pretty different. So I did have you know programming experience and all, but I hadn't had like back-end, like like proper back-end experience before.

SPEAKER_00

Right. And uh uh you know, a moderately complex API for a uh you know, somewhat difficult um uh category of infrastructure software. Right, sure. So yeah, that's that's uh that's quite a thing to be thrown into on the internship. So what kinds of things did you learn? And this this doesn't I I'd love to hear, you know, Kafka lessons, Kafka streams lessons, and really just development, because it's uh you know, it's not every day we have somebody on the show who um you know who's who's writing their first big project. So what was that like?

SPEAKER_01

Right. Yeah, so I think the uh Kafka aside first, the big part of the internship that was new was the amount of design that went into it. That it wasn't necessarily that you know I could sit down immediately and write code. And it was also that to actually come up with a good design, I had to debate with my my teammates, specifically Eugene, who is my mentor. There was a lot of times that we'd actually, and for for once, since we were I was actually designing an algorithm, you know, not just maybe some web page. We actually had to use the whiteboard, which I thought was really cool, which I didn't necessarily always get to use in an internship, and debate what the pros and cons of different approaches might be. And sometimes my approaches were wrong. It wasn't that there was one one right answer. There was you know also an alternate solution that Eugene proposed. So it just seemed very collaborative. There were a lot of little things to be considered. I think that's another big part is that there's a huge attention to detail that was needed that I didn't realize would be needed. You know, just forgetting about one detail. Or for example, I didn't think about say out-of-order events at first. Um, I think this is said on the blog post. I wanted to just use a DSL, which was you know pretty and functional and all that, but then I had to end up using punctuate. It just made sense because of the out-of-order events and some other reasons. So overall, it just there was a I had to pay a lot of attention to detail, and I had to you know defend my design and collaborate. And both of those things were really fun.

SPEAKER_00

Awesome. I'm glad you like them because that is a great deal of what a career in software development is like.

SPEAKER_01

Good to know.

SPEAKER_00

Um I mean, yeah, there's plenty of heads down uh just writing code and trying to learn a new API and and reading through docs and writing tests and all that kind of stuff. I mean, that's that's the the bread and butter of the quiet time of coding. Um, but a great deal of the time is not quiet. It's it's actually talking to people and uh you know, debating ideas and and you know, putting forth your proposal and you know, kind of over time developing the wisdom to know. When uh a criticism of your favorite approach is a valid one and you should give up the thing that you thought was beautiful, um and when to say, well, no, this is really better, we need to we need to go this way. But that uh kind of whiteboard uh debate, there's a lot of that. It's so important. And that's that's where um uh that's where as developers we make new knowledge.

SPEAKER_01

Of course.

SPEAKER_00

Yeah. Glad uh glad you got to experience that in a positive way.

SPEAKER_01

Yeah, I'm glad too.

SPEAKER_00

Cool. Um what other uh uh what other lessons like uh were you you're new to Kafka, I assume. Was there did you have any coursework or anything?

SPEAKER_01

So that was interesting too, because I didn't necessarily have follow an explicit course. I just had to you know look up the blog post. In hindsight, I think a course could have been or not like a full-fledged course, but sometimes when you're just learning things through blog posts, you don't know what you didn't learn. And that was why it was really helpful to be to have, for example, Eugene, my my coworker on my team, was knew already knew about Kafka. So if I describe my algorithm to him, he'd be like, he would be able to point out something, something wrong with it often. And I would be able to, he would explain this, the reason the way Kafka works, and I'd be like, oh, well, I should probably go understand that. So I'd go go back to the internet, understand it. So that was the feedback loop, but in hindsight, I wish I spent I I think I maybe rushed a bit through the design aspect and specifically about understanding Kafka streams deeply, or maybe that was just part of learning is that you know, the first read-through, you don't grasp everything. But I did miss parts of Kafka streams, and I didn't realize how important the way the architecture is is. For example, even just small things. One one thing is that when you know, and in Kafka streams, when a consumer fails and restarts, it consumes messages from the most recently committed offset. And so our setup and most setups use the at least once setting where it'll just make sure to emit that message once or multiple times. And so as a result, your your service has to be idempotent and handle duplicate messages. But my original solution never accounted for that. Um Eugene didn't didn't necessarily notice it either. And so that's something I had to go back and fix at the end, even though I had read over that text that described that you know Kafka streams has this special setting. I'm like, I'm pretty sure my eyes just glossed over it. I'm like, uh hopefully it's not relevant to me. Scroll through.

SPEAKER_00

Exactly once, yeah. Whatever, whatever. So I think that that did you end up did you solve that by enabling turning on exactly once semantics?

SPEAKER_01

Uh yeah. So the semantics were already on, but I just had to make sure my service would be able to handle duplicate messages.

SPEAKER_00

Gotcha. Yeah, gotcha.

SPEAKER_01

Um I'm not sure if you said the at least once setting, but that's what we use.

SPEAKER_00

So then I just you used at least once. Okay. What did you did you consider exactly once as a as an option? Uh there's a there's a configuration parameter. Right, right.

SPEAKER_01

You can also use um exact once. Uh I don't remember exactly the reason why, but our company wanted to stick to at least once, so I had to change my service. Yeah.

SPEAKER_00

Gotcha, yeah. It may change your service so that it's idempotent and it's effectively the same. Yeah. Yeah. Uh cool. Yeah, it's good that you could make that make that uh adaptation. So here's a very self-serving question uh for me. What would have been for you, like if there was an online resource or a thing that you knew about at the beginning of this that told you what you needed to know about Kafka and what you needed to know about streams, what would that look like? Now that you know, just looking back a few months, what would the ideal uh you know, derp I'm new at this kind of thing be?

SPEAKER_01

I mean, well, I think to be honest, I did use most of the conflow documentation. I forget what it's called. It's called like a tutorial or like there's some table of context on it. So that was the the majority of it. And I think there were only a few other posts that I used. But I think the the problem was that I was I didn't realize that was the source that was the most thorough. So then at the beginning, when I just looked at other blog posts, they weren't as thorough and I would like miss things. For example, I totally missed the or I think we as a company totally missed for a while the topology test driver, which Kafka Streams provides, which is totally detailed in that tutorial. And we used we were using other things, we used mock streams, we used um embedded Kafka for a bit, and when we used embedded Kafka for a few tests, very few, I have to say, but we had to use thread.sleep calls in our code, which is pretty terrible. Um but I but then in the tutorial I noticed topology test driver, and I was like, oh wait, why are we not using this?

SPEAKER_00

Yeah. Yeah, let's let's not use thread.sleep in the bad. But I mean, hey, sometimes the first first pass through, and you know, come on.

SPEAKER_01

Yeah, no, that was yeah, that was not good, but yeah.

SPEAKER_00

Um good. And anyway, I'm glad that you uh glad that you came across topology test driver.

SPEAKER_01

Yeah, yeah. Um yeah, no, so to answer your original question, um I think something that actually like delved into the architecture as would be relevant, um which I'm sure the confluent documentation did too. But yeah, so that one that goes into architecture, but also you know presents things on a basic level that's digestible and is thorough. Nice.

SPEAKER_00

Got it, got it as an architecture introduction. Otherwise, shout out to docs.confluent.io. I appreciate that. That's uh glad that's been a helpful resource for you.

SPEAKER_01

Of course, yeah, that was awesome.

SPEAKER_00

Yeah, good deal. Any other uh any other lessons learned as you look back, it's a few months later now, or what is what is how how do you think about the experience?

SPEAKER_01

I I think that's this is just I think the overall lesson is that when I'm learning some kind of back-end system, whether it be Kafka streams or something else, I would hope that I actually try to understand how it works under the hood because sometimes that affects how you actually develop your solution on top of whatever it is. And there were plenty of debugging things, bugs I had to debug that were as a result of not understanding it deeply. And I didn't perhaps that could have been avoided. And so if I had to go back and redo my internship, then I would spend more time understanding Kafka streams at the beginning, and that potentially could have saved me a lot of debugging time. So I think one one funny one is uh uh I only some of the walkthroughs were being generated when I first got everything to to work, and I was really unsure about why that was for a few days. But that's because I thought I should have known this too, but I assumed that the case dream map method forced a repartition by key, but it didn't. So when my code would try to process a walk-in, it wouldn't be able to find corresponding Zenreach messages. And so that took me a while to figure out as well, like a I think a few days. And if I took the time to understand, you know, which methods would repartition, which wouldn't, also could have been avoided.

SPEAKER_00

Gotcha. Lots of uh lots of little details your first time through. Yeah. Well, our guest today has been Rishi Donaraj. Rishi, thanks for being a part of Streaming Audio. Thanks, Tim. And there you have it. I hope that was helpful to you. If you've got questions, you can ask me at TL Burgland on Twitter. That's T-L-B-E-R-G-L-U-N-D. Or you can leave a comment on any of our YouTube videos. Your question might be featured on the next episode of Streaming Audio. And feel free to subscribe to our YouTube channel and 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 and just generally helps us get the word out. We appreciate your support. See you next time.