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

Apache Kafka Fundamentals: The Concept of Streams and Tables ft. Michael Noll

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

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

0:00 | 48:52

If you’ve ever wondered what Apache Kafka® is, what it’s used for, or wanted to learn about Kafka architecture and all its components, buckle up! In today’s episode, Michael Noll (Principal Technologist, Confluent) and Tim Berglund (Senior Director of Developer Advocacy, Confluent) discuss a series of fundamental questions: What is Kafka? What is an event? How do we organize and store events? And what is Kafka Streams? 

Over the course of this episode, Michael covers an in-depth look into Kafka technology and core concepts: the process of reading from a topic, differences between tables and streams, mutability, and what ksqlDB is and what its event streaming database features accomplish. If you've ever wanted to get a better grasp on how Kafka works, this episode is for you!

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

When I first learned the difference between streams and tables in Kafka, I learned it from Michael Mole. Today I talked to him about that, but to get there, we take a journey through all of Kafka, from events through topics, all the way up to KSQLDB. It's in today's episode of Streaming Audio, a podcast about Kafka, Confluent, and the Cloud. Hello and welcome back to another episode of Streaming Audio. I'm your host, Tim Berglund. And I'm joined in the virtual studio today by my colleague, Michael Knoll. Michael, welcome to the show.

SPEAKER_01

Hello, everyone. Thanks, Tim, for having me.

SPEAKER_00

You bet. Now, Michael, you are uh a member of the Office of the CTO, uh, sometimes abbreviated Octo. For the longest time I thought that meant opposite of a CTO, and now I realized it's Office of the CTO. I was so surprised. Anyway, seriously, what uh a lot of us are familiar with that title, but uh, what do you do in the office of the CTO?

SPEAKER_01

Yeah, let me answer that question by going on uh on a slight detour. So I joined Confluent very early on, one of the first few employees. And uh most of my years uh at Confluent were spent being the product manager for one of our um product lines, which was stream processing in my case. So I've been you know working on Kafka streams and KSQLDB, etc. And I did that until last year, when I moved to the office of the CTO at Confluent. And what I'm doing now is um you know taking a step back from the more operational role on the product side, and I'm now focusing more on longer-term product and technology strategy at Confluent, plus a few other things here and there, you know, like doing a podcast or speaking at a conference and some other um functions that I uh do for Confluent. But mostly it's about you know this forward-looking work in our technology space.

SPEAKER_00

Got it. Uh another thing that you do is you occasionally write blog posts. And there will be linked in the show notes of this podcast, a series of four blog posts that you recorded, uh that you wrote, rather. Um you don't really record blog posts, you record podcasts, you write blog posts is usually how you say that. And they uh what's what's remind us of the title of the series.

SPEAKER_01

So the series is called you know, streams and tables in Apache Kafka and what every software engineer should know about them. Uh essentially it's a whirlwind tour through the Kafka fundamentals, you know, going from the processing uh layer to the storage layer, um, concepts uh and fundamentals on on that front, etc. Essentially, everybody who is interested in Kafka should find something useful in that block series. Totally agree.

SPEAKER_00

And that is why I wanted to have you on the show, because it was such a good introduction. And uh what I what I love about it is you your goal is to talk about streams and tables. And it's funny how much I personally associate you with that notion because uh you were the person when I was new at Confluent who first explained that to me. And every time I'm explaining, you know, uh Kafka streams or KSQL or something like that to someone, it's um that thing, this this concept seems to be one that's sort of sticky. It's it's hard to get into people's heads. Like everybody gets what a stream is if they know how to think about Kafka, but when you introduce tables, it gets difficult. And so um I remember my own difficulty getting those concepts to stick at first, and you're the ones who you're the one who made it clear to me, but thank you. And you set out to do that in this series of blog posts, and I'm like, well, Michael's the streams and tables guys, that makes sense. But the way I love the I just love the way you did it. You started at the very beginning, and I just want to walk through that whole process. So for those of you who like to listen to people talk instead of read, this is your opportunity to have Michael Knowle explain streams and tables to you starting at the beginning. So let me ask you the first question. Um, what is Kafka? And I know everybody who's listening probably knows, but sometimes it's helpful to get a guy like you to give an explanation anyway.

SPEAKER_01

Yeah, so my answer to your question, Tim, is you know, Kafka is an event streaming platform. And I think before I continue with that, I think it is important that you know we're in an industry where a lot of things are actually pretty complicated. And most of the people that you know are using Kafka in practice are not like you know, we at Confluent who are developing Kafka, you know, and dealing with Kafka 100% every single day. But they have a lot of other things that they worry about. You know, people want to deploy things, so nowadays they have to learn about Kubernetes. Previously they did it on-prem, now they're moving to the cloud. So like Kafka is only like one share of everybody's you know mindset every single day. So that's why I think it's very important to uh help people get up to speed with things like Kafka in like a podcast or blog series. And I think it's always great when we see people in the community, whether it's for Kafka or other open source projects, you know, give back, maybe not necessarily by contributing uh a very complicated uh pull request, for example. Um, sometimes it's much more helpful to help others understand what you just learned yourself. So uh I'm using, so to speak, your podcast as a means to encourage people that might be listening in that you know, you could be the next person helping out your colleagues or somebody on the internet in and in the community to make more sense out of technology. So uh that's just something I want to bring up here before starting our you know deeper dive into Kafka, because I think everybody as at someone is a learner. So uh need to keep that in mind.

SPEAKER_00

And and really good point that everybody has, you know, at Confluent, we think about Kafka, we think about Confluent Platform, we think about Confluent Cloud. Uh it's like event streaming all day, every day, because that's the kind of technology company we are. You know, suppose you work for BMW, you have to know things about cars and all the related infrastructure. Or uh, you know, you work for Walmart, you have to know about retail. Um, and developers there need to be specialists in those domains. And a lot of their mind has to be focused on that. And not hopefully all of it is going to be focused on infrastructure. That's that's what infrastructure people do. So yeah, that's a good reminder. So yeah. So it's an event streaming platform.

SPEAKER_01

Going back to your question. You said no, it's an event stream platform. So then typically people ask me, okay, like what is an event? Like, is that like Netflix or YouTube? Is that what the streaming is about? I said, no, no, no, it's not that. So let me try to tear this apart. Now, that's a bit more difficult in a podcast than if you have a whiteboard, but I'll try my best. So there's event, streaming, and platform. So what Kafka allows you to do is it allows you to publish and subscribe to events, it allows you to store these events for as long as you want, and it allows you to process and analyze these events. And that sounds already like a very compelling uh value offering. But then again, like what is an event? You know, you can publish it, you can subscribe to it, et cetera. But what is an event? And simply speaking, an event is just something that happened in the world. That could be anything depending on you know the company you work for. It could be something that you know a good was sold, it could be you know a financial payment that was being made from one person to another at a specific uh date and time, etc. So that is simply an event, what we're talking about here. And uh Kafka allows you then to work with these events in a variety of ways. So I hope that is a good introduction of what Kafka is in a nutshell, without going too much into the tech side of things.

SPEAKER_00

Absolutely. Yeah, no, it's good. Especially if if somebody brand new to the topic who's just moving past the stage of awareness that there is a thing called Kafka, that's all you want to say at this point. Uh we we want to build this a piece at a time. So we have these events. And how do we organize and store these events?

SPEAKER_01

So, what what Kafka allows you to do is then you know to continuously capture these events and then do something to these events as they are occurring. And in Kafka, this is being done through streams of events. So a stream of events is simply the history of what is happening in the world, and you just capture those sequentially as into these aforementioned events. And the reason why I think a lot of people are nowadays so interested in Kafka is that you know a lot of things have changed in the past 10 years. So, you know, the advent of mobile is one, um, the move to the cloud is another, AI machine learning is another trend. So particularly with mobile, I think we've seen that you know, also as consumers, like an average Joe consumer, you're expecting your bank or your you know retail store to give you as great of an experience that you might get from you know downloading something from Apple, from the App Store. And this is really up the ante on what a company has to do now in order to satisfy their customers. And a lot of that has to do with speed. You need to be timely, you need to respond to things as they occur in your business, whatever that business might be. And that is why so many companies have started to use Kafka to model their infrastructure and reshape it so they are fit for the 2020s, so to speak.

SPEAKER_00

Yes, indeed. So um, and there's a there's a particular term uh that I'm looking for, the, the, the place, the place events go to get stored, we call that a topic. So tell us about topics, storage format and topics, how topics scale, just give us a brief tour of topics.

SPEAKER_01

So if we're starting at these events that we then capture into a stream, then um you know, starting already to look a little bit under the hood, that needs to be stored somewhere. Because in Kafka, which is different to traditional messaging systems that people might have used, like Red RabbitMQ, for example, Kafka stores this data permanently. And it does so by storing these event streams into topics. And the topic and a stream is very much alike. Um the main difference is that a topic does not have a schema. So when uh an event is stored in Kafka, specifically by a Kafka broker into a topic partition, then the event needs to be serialized into bytes because everything that the Kafka cluster knows about your data is just raw bytes, has no friggin' idea what is in the data, other than you know, some random, random binary um information. And then maybe your other question then is well, once it is in a topic, how do you get it back again into readable format? And that is where then you apply a schema to your data in the topic. And this allows you then to get back, you know, like payments or you know, user location updates, etc. etc.

SPEAKER_00

Yeah, you've got you've got some domain objects that you're thinking about at the application level. You're you're not thinking about bytes at the application level, hopefully. There's some sort of domain object, and that those get serialized and deserialized by a layer, which we'll wave our hands about right now. We'll just consider that to be magical. Um, but that happens to get the things into bytes in the cluster.

SPEAKER_01

Exactly. And there's there's one important thing that I want to point out here is that sometimes people think, well, the the broker doesn't really do all that much, right? Like it doesn't know what's in the data. It's mostly like the producers and the consumers that do the real work. And that is true. And that is true by design and by intent. Because this approach of dump brokers, and I put dump you know, in quotation marks, means um the broker can scale very, very well. It's because they don't have to deserialize and re-serialize messages as they you know come and go. Um, they're saving a lot of horsepower that they can do uh use for other activities. So, in other words, this design decision of Kafka, the Kafka Prokurs only knowing about by it is a key part, a key reason why Kafka scales so well.

SPEAKER_00

Yeah, and I I sometimes uh by the way, I agree uh with the fundamental genius of that decision. Um it takes a lot of discipline over time to keep brokers dumb because as you're adding features to the platform, there's all these reasons why you could push smarts into the broker. And sometimes I'll get people saying things to me like, well, uh consumers have to poll, and polling's bad, and if brokers could push, uh consumers wouldn't have to poll. Yeah, that's a lot of state that goes into the broker uh all of a sudden. You know, there's a there's you're gonna pay for that in scale if you if you want brokers to be able to push, or you know, we'll get into stream processing in a minute. Um why stream processing doesn't happen on the broker. You know, we still have to serialize that data and move it over a network at some point to do computation on it. Wouldn't it be better to avoid that network hop and put stream processing on the broker? Well, you know, you're gonna get some pretty big Kafka clusters that way. So it it um I think it has taken discipline for the the core Kafka team uh to maintain that broker stupidity over time, but it's so important um and a good thing that we can't emphasize enough. How about partitioning? Uh did we talk about partitioning the way we we break uh the way we scale topics?

SPEAKER_01

Yeah, partitions are for me always something I really like to talk about because I think a lot of people don't yet fully understand how important the concept of partitioning is in Kafka, because it's something that you know is ubiquitous one way or another when you're working with Kafka. And um just you know, to wet your appetite a little bit, when data is read from or written to Kafka, when it's being processed, when you're joining data from two different streams, for example, or from a stream to a table, when data is being stored, when the stored data is being replicated for fault tolerance, and uh when we talk about ordering guarantees, all of this is based on partitions. So we're talking about fundamental aspects of the storage layer of Kafka, which is like the Procurs that we just talked about, but also uh the processing layer, you know, things like KSQDB, Kafka streams, and other tools in the open source um world, you know, like Fling, Spock, etc. etc. Everything happens there fundamentally because of these partitions, and it works the way it does, because the data is being partitioned. So that's why this is super important to understand what these partitions are and how you make the best use of them.

SPEAKER_00

Excellent. Uh and we'll come back to that as we uh as we get into stream processing a little bit, a little bit down the line. So we've got events, uh events are stored in topics, uh they're stored just as bytes in topics. There's some serialization and deserialization that happens to get domain objects in and out. Topics are partitioned, which is our way of of really scaling how big topics can be. And you you stressed this topics are durable. Kafka stores messages uh in a durable and potentially permanent way. It's not like an ephemeral cue where we're just trying to hold things for a little bit until we can get back to them. But tell me about consumers. I I wanna I want to set up this idea of streams and tables and what happens. So just talk to me a little bit about um without any sort of stream processing API, what life is like uh reading from a topic, just being a Kafka consumer.

SPEAKER_01

So when when you're reading from Kafka, and what I'm trying now is I try to paint a picture in your mind. Now, for those of you who listening in, um, because otherwise it might be a bit difficult to understand. Um imagine um you have like you know a canvas in front of you and you're dividing it like in two halves, a left half and a right half. On the left half you have a topic, and on the right half you have your consumers. That could be you know a Kafka Streams application, could be KSQLDB, it could be a plain Kafka consumer that you're using in your microservice, etc. Now on the left side where the topic is, the topic is partitioned. Let's say it has four partitions, you know, color-coded, maybe like green, blue, red, and yellow. And now your application on the right side wants to read from that topic. Now, what is happening behind the scenes when you're you know creating um your application that wants to consume, it will read from that topic through a consumer group. And inside that consumer group, one or more consumers will be active and actively reading from those topic partitions, you know, those four that we have on the left side with the colors. Now, as you're adding more consumers, which could be, for example, if you have an application that is containerized, if you're starting with one container, there um is just this one instance of your application running. If you're a second, a third, and a fourth container, now there are four of them simultaneously processing and consuming data inside the same consumer group. Then what we can do now is we can draw a line from the partition on the left side, let's say the green partition, we can draw a line to the right side to the green consumer inside that consumer group. And this is the way that Kafka allows you to read data in parallel and lets you scale the processing, which includes consumption, but also like producing new events that you derive from your original data. And as machines come and go, or as containers come and go, these lines can, you know, from left to right, these lines can move from one container to another container, you know, as these containers are being added or removed from the picture. I hope that was a good semi-visual overview of what is going on when data is being consumed from Kafka. And the fact of that is data is being read in parallel, and even in the face of failure or in a situation where you want to have elasticity, like adding or removing capacity, Kafka will make sure that things are staying busy all the time in a positive manner.

SPEAKER_00

Yeah, trying to make sure all the consumers in the group are the the work is assigned to them as fairly as possible, as evenly as possible. And yeah, that is that that's a good description. It is sort of a fundamentally visual thing, but you've got some number of uh we'll say call them containers or you know, application instances, all instances of the same consuming application uh doing the consuming work, some number of partitions, and those partitions get mapped to those consumers. Those uh those assignments get made automatically by the cluster, right? They're the consumer group and the cluster cooperate to make that that happen. It's not like a thing you have to worry about in application code.

SPEAKER_01

Exactly. And and I caught myself, Tim, trying to use some words that we haven't talked about yet. So it's hard. That's why I I hesitated a little bit when I was you know trying to paint this picture. I can't use the this term yet because we haven't even touched it yet. Right.

SPEAKER_00

It and you know, any just the whole Kafka platform uh I've found, it's kind of funny that you say that. There are always forward references. Uh you know, you you're you're trying to build it a piece at a time. And I've got a talk, maybe you've seen me give it, but it's it's this same sort of sequence where I start with events and I end up with KSQL on the other end. You know, you try to build it a piece at a time, and inevitably you have to say, I can't really explain this yet, but here's this thing, and it'll make sense in a little bit. Um, and that's you know, that's okay. Understanding always proceeds that way. You know, we're we're willing to uh uh have some ill-defined concepts sort of hanging in our mind for a little bit and then flesh them out later.

SPEAKER_01

Yeah, and and to add to what you just said, Tim, and this is maybe for the listeners here, like how the sausage is being made. My blog series and subsequently this podcast, this started by a very innocent email I got from someone in Tim's team said, Hey, you wrote this article about you know streams and tables on your personal blog, you know, more than a year back. Can't you just write a short blog post also for the Confluent website? So it was like, you know, how long could it take? It takes you maybe like a day or two, and then we have it. And this turned then into this you know four-part series. Uh, and we spent quite a lot of time on that, you know, to get it right and get the flow that you know, you're not using forward references, as you just said. So um, yeah, the torches is it's being made in a way that you always you not always all want to see how it's being done. Right, right.

SPEAKER_00

And there's so many beautiful diagrams in it. A lot of work went into this blog series. But this is the Confluent Office of the CTO. Uh, can you write a short blog post? The answer, you know, on on something. The answer is no, but they will write some very good blog posts. They just won't be short. And uh this that's what that's like I said, that's what this uh uh is all about. It's a great series of blog posts. And so I hope uh as you listen to this, uh also refer back to those in the show notes. There'll be some pictures and things like that that'll that'll help. Okay, so uh we have a consumer group. Now uh the way the way I always think about this is some things happen when you write consumers. There are there are things that you're going to do. It doesn't matter whether you're in finance or manufacturing or retail or defense or whatever, right? Whatever kind of business you're in, um you have invented data, uh, you've chosen to model your system around events, and so you're using Kafka, and you're gonna write consumers, you're gonna read data from topics, you're gonna read events from topics. And there are just some patterns that that present themselves. There are things you do in consumers that that well, uh for example, you might want to compute an average of something. And if you're gonna compute an average, if you just think of that like you you would in terms of a of database, um You're going to group by something. And once you take events and you group them by something, you don't really have a stream anymore. You have a table. So I guess for the first time, but tell us we've talked about what streams are, but tell us what tables are. And if you could spend a little while differentiating between the two. Sure.

SPEAKER_01

So let's let's recap very briefly what a stream was. So for a stream, we said it records the history of what happened in the world as a sequence of events. And sequence in the sense of like an ordered list, an ordered sequence of events, and it can be could can potentially be unbounded data. So it keeps just growing and growing and growing. Example could be a sales ledger. Or uh maybe my favorite example, you record the moves in a chess match. Yes. And now a table represents not the history of the world, but the state of the world at a particular point in time. And typically that point of time is now. So in the chess example, the stream would be the chess notation. No, no, pawn moves from E4 to E5, etc. So even and the table represents the boards. And what you see at you know right now on the board, you know, with all the cheeses having moved a few times.

SPEAKER_00

Right. Let me stop you on the chess notation because if you don't know chess, I don't know chess particularly well, but you can imagine, like you've heard in a movie or something like that, where there are like the two super geniuses who are just saying kind of like weird words back and forth to each other, playing a chess game in their minds, you know, they'll just say, like, you know, uh King to Rook 7 or whatever. I I don't I don't know what that means, but it's this you could do it as dialogue. So just picture the the movie scene where people are saying these words back to each other as state back and forth to each other as statements describing the unfolding of the chess match. So each one of those, and you don't need to know chess well, and you don't need to know what those weird words mean that that that happen in that sort of cliched movie scene, but those are events, those are descriptions of mutations of the chess game moving forward kind of one discrete step at a time. So just picture those weird chess words and do not worry if you don't know chess. Back to you.

SPEAKER_01

And Tim, I can I can also offer you something contemporary.

SPEAKER_00

Good.

SPEAKER_01

We have now a global pandemic at the door. We do at the time of this recording. You can also imagine. Yeah, exactly. And you can also imagine that you know one stream of events could be reports from hospitals and doctors about people that have been uh infected by the uh you know COVID-19 virus. So it could be you know tests that have been done with a positive result, tests that have been done with a negative result, and all of that information from you know across the globe or you know, your home country is being collected, and then you can compute the current state of the pandemic. So how many people are currently infected? You can also rewind the table in time by going back to that recorded stream of information to say, how did the uh the pandemic look like a day ago, a week ago? What was the difference between, let's say, the United States and um China or South Korea? So you can slice and dice your data um as needed because the stream will not forget. All the information, you know, ordered based on the timeline is available in your event stream. And then you can compute your tables above that. And the table could be number of affections per country, for example, which would be one way to group the data, where the grouping would be by country, and then the aggregation would be the total number. And uh what you can then do with Kafka is not just what I just described, but you can also have this table being continuously updated 24-7 as soon as a new piece of information arrives in the underlying stream that populates the table.

SPEAKER_00

So that table is always a snapshot of the state of the system right now. Um you gave the pandemic example, or even just the chessboard. I mean, everybody understands chess well enough to know that like look down on the top of the chessboard and the position of the pieces is the state of the game. So the table is that current snapshot, whereas the stream is this record of the events that have unfolded to produce that state. Talk to us about mutability. Uh I'll just I'll just put this out. Events are immutable, right? You do a thing, you can't undo the thing. That's just fundamental. And hey, as I always say, sometimes sad. Uh right. You wish you could unsay words, you can't unsay them. Um but uh events are immutable things, but talk to us about uh streams and tables in terms of mutability.

SPEAKER_01

As you said, Tim, we could think about a stream as something that you know it's like a you know a financial ledger. So whatever you write into that stream, you won't touch it again. And it is immutable, which means you can't change it. So that is the stream side of things. So uh if you tr try to draw a comparison to you know the operation database world, we could say, you know, you can uh it's like an append-only table. You can only insert, append new events, but you can't update or delete existing events. And a table is a bit different to that. A table allows you to mutate data. So you can do inserts, updates, deletes, and so on. And then in Kafka case, it's uh the key of an event, and we haven't talked about that yet, but um, it's a key, the key of an event determines which row is being affected by a statement. And uh what is a key? So in Kafka, um an event, also called a message in the documentation, has a number of fields, and the two most prominent ones are you know the event key and the event value. Event key could be something like you know, the username and the value could be what the user just did, and that was recorded as the event. And there are a few other attributes for an event, including like a timestamp, like when did that event occur in the real world, for example, um, other things like headers, but let's ignore that for the moment because it's not super relevant to understand how all of that you know ties together. So recap stream, think about this as something you write into, but you will never change it again. And a table is something that you can derive from a stream, and here you know things can change. So as you know, infections in a pandemic go up and down, you see that the table is updating along the way, whereas the stream just keeps crowing and crowing with more information that is being fed into it.

SPEAKER_00

Got it. How about data contracts? You we talked before about how we have domain objects and Kafka doesn't care, Kafka's just bytes. How are how are data contracts maintained? I mean that that that sets us up for a difficult problem because our view of these objects, of course, evolves. Uh you know, the the definition of a domain object evolves over time.

SPEAKER_01

So Kafka allows you, which is also a fundamental concept, to decouple the parts of your architecture that produce events and those that consume them, which is great. Like for a larger organization in particular, you know, the team that is responsible for generating some data, maybe like a car and you know how the car is moving or how the engines are working, etc., that team can operate independently from the team that is using the data and to do other things with it with the information. And it could be, of course, you know, more than one team. Now, what Kafka um has in its design is that the producer and the consumer have to agree on the schema and the format of the data, so to speak. Because as we talked about before, the broker has no idea what's in your data. It's just bytes. So the broker can't help you there. And in practice, what does that mean then for Eula? How can you know the team that, you know, the car team, so to speak, how can the car team help the you know predictive maintenance team so that you know both can understand each other? Right. They both have to have a common deal. There are various ways. Well, by the way, I'm I'm uh German native, I have no idea how a car works, I'm not into soccer, and I also don't like beer particularly much. So uh I'm I'm not representative here, but I'm trying my best to use the car analogy still. Totally. So there are a couple of ways you can do that in a company. You know, sometimes it's you know, I'm I refer to this in my blog series. Sometimes it's you and your fellow colleagues sitting in the you know company canteen and you know scribbling uh on the back of a napkin what the data informat should be like, and you just agree on it, like out of band, just a convention. You promise that, hey, this is how I write my data into Kafka, so you know how to use it. So that is one way to do it. It's more like a free-for-all wide rest approach. Um there are more structured and more disciplined ways to do that. And um typically what we see then you know, as soon as you know Kafka is gaining adoption inside a company, and people are really depending on the data. So it's not like a nice to have side project, but it's like a mission-critical use case for the company. Like um, we have a customer of ours is you know running a stock exchange, like a pan-European stock exchange on top of Kafka. If this doesn't really work, you're into trouble when you're responsible for running this infrastructure for your company. And uh what we see these people do is they are adopting data contracts that are more explicitly um defined. And one way is you know by using your confidence schema registry, afro, protobuffers, etc., where you use programmatic rays to ensure that what is being produced into Kafka can also be consumed again from Kafka and people can make sense of it. And something that we've added recently on the confluence side is um a feature that allows you to do a server-side validation. So the program comes somehow into play here. I I won't go into the details, but um thus far, um, and tell me Tim whether I'm going on a tangent here or not, but thus far, even if you do conference schema registry, it is still like an agreement on both sides that you know you're following the convention of using schema registry and so on. And with this new feature that that we added um you know earlier this year, we allow the server side, which is the broker side, to reject a message when someone tries to produce to a topic if that message isn't um formed um you know well formed according to your defined schemas, which is very nice because then you don't have to rely on some teams to you know do it for you, but you can really enforce that, which is great if you have you know security policies in place, you have you know regulations that you have to adhere to, etc. So um, in other words, to wrap this up, um, because there were quite a few things I touched here, typically Kafka is a schema on read setup. So you're getting biased as a consumer, you determine how you want to read it, but ideally, you want to really know that you're reading it in the right way. And you know, schema registry and so on helps you to have this communication with those that produce the data, how you should do that. And with the server-side, you know, schema validation or schema enforcement, if you want to call it like that, you can really enforce that no one can write junk data to your Kafka cluster, even if people try to you know skip the convention, so to speak. Right.

SPEAKER_00

Now, uh, and that's that's helpful because that that you know, we were intentionally being a little cavalier before saying, oh, there's domain objects, they get serialized, it's fine. That is actually a thorny problem. I mean, that's the problem of schema evolution, which is hard everywhere. Uh schema registry makes it makes it a lot more manageable. We're also talking about consumer groups and how you can have multiple instances of a consumer and how partitions of topics get assigned to them, and so work gets done fairly when we move up a layer into uh Kafka streams, and I I I think I should ask you to define Kafka streams, uh that that whole assignment of work gets more complex. I realize as I'm asking that question, tell us about Kafka streams first. Uh you used a BPM for this, but not not very many people in the world know it better than you. Uh what is it?

SPEAKER_01

Sure. Before I do that, let me just add one thing to what we just discussed because maybe it is a bit easier for listeners to understand what I was trying to get at. You know, we talked about schema on read. I didn't mention schema on write, but that is what you get with server-side enforcement. Just imagine the following. Imagine you have a team meeting in your company, it's a national company, you walk into the same meeting room, but no one agreed that English is the business language that everybody should speak in that meeting. Imagine everybody would speak with their native language, and you can imagine pretty easily that no progress would be made in that meeting.

SPEAKER_00

That works in Star Wars though. You notice everybody speaks their native language in Star Wars and everybody else understands it. But in the real world, that would be terrible.

SPEAKER_01

Yeah, my son still can't understand how Han Solo can understand Chewbacca. Right. Right. But there you go. Yep.

unknown

Yeah.

SPEAKER_01

So going back to Kafka Streams, your your question there. So Kafka Streams allows you to build event streaming applications. And it allows you to do that by giving you a library for the JVM, so you can use it from Java, Scala, Clojure, etc., to build applications that read from Kafka, write to Kafka, communicate also with other applications, you know, without Kafka in between. And what Kafka Streams gives you is elasticity, fault tolerance, and so on, and you know, quite a few batteries included, so you don't have to build a distributed application from scratch. And that is a great value proposition if you're working either with you know a little bit of data, so not a large volume, but still mission-critical data, or very large data volumes, because it gets rid of most of the things that you would have to worry about that are very complicated to understand and deal with by giving you that information, um, that functionality built in. So you don't have to worry about that. Example could be a fraud detection application that reads data from an event stream with payments. Does it join against a table of customer profiles? So you know context information about every transaction. Has this customer ever paid with their credit card from Argentina? No, was always inside Europe. Probably this Argentinian uh transaction is fishy. So let's let's alert on that. So that would be an example of what you can build with Kafka streams.

SPEAKER_00

Awesome. And it's it's providing for you a lot of framework that you'd build yourself if you if you don't use Kafka streams, you just program the consumer API, you're eventually going to build kind of some, as I put it, a buggy partial implementation of Kafka Streams inside your own consumers. So it's it's all that framework. Um and of course Kafka Streams is a part of Apache Kafka. This is not a constant thing, uh, this is uh part of the ASF project and available just as a as a as an API uh that you can use. So with that, rem recalling uh the question I asked a minute ago, how partitioning and uh consumer groups interact, how how partitions are assigned to instances and consumer groups, sometimes when people start to think deeply about streams and tables in in the context of Kafka streams, they'll think, wait a second, uh what if I have a table that's really big? You know, it's this state, and you tell them, yeah, it's kept in memory and it's also on disk, and some people start to freak out. They say, wait a bit, wait a minute, what do you mean it's in memory? How can I scale that? So what if I have a really big stream that's got a lot of state, like a lot of unique keys in it, and I I want to make that into a table, you know, like it's a 1 million by 1 million square chessboard or something. Maybe make it bigger than that. That that's hard to fit into memory. Uh so how are how do tables scale?

SPEAKER_01

Yeah, that would be a pretty big chessboard. It would be. So what what Kafka Streams allows you to do, and and the same applies also to KSQL D DB, which is built on top of Kafka streams, is is the following. So and and let me let me take a step back real quick, Tim, because I think we we need to cover like one or two additional things so people can better understand what we're talking about. Um if you want to work with a stream in Kafka and something goes wrong, no, that's no big deal. No, everything is recorded in order, so you can just rewind, you know, like an old cassette deck or your backup tape, and you can start again. That's no big deal. So dealing with situations where you know something something crashes or you know you made a mistake in your application, there's some bug, etc., that's not a big deal. You can just rewind your topic, rewind your event stream, and can start again. What is a bit more difficult, um, but this is being handled for you by Kafka streams and KSQLDB. So that's one of the cool things. But I'm explaining now what happens behind the scenes. And that is a bit difficult. When you're doing certain computations, you know, like aggregations, joints, these are computations where you need to remember what you have processed before in order to produce a new result. So if you're counting the number of infections of COVID-19 in the US and you want to increase it by one because you just got a new update from a hospital, you need to remember what the previous total was, right? And then your question, Tim, was what if if that previous total is actually not just a single number, but maybe like a huge amount of information that you need to keep track of? And what Kafka streams and KSQDB do here is the following. They um use local, kind of like a local database. Um, in Kafka streams, it's called like local state stores, to you know track this information and maintain these, you know, this state, um, these totals, etc., um, for you. And whenever there is a new event, you know, they will look up their local data, you know, increment it or do whatever they need to do to the data and you know, process the next event. That is um also something that happens when you work with a table, because a table is you know nothing but state in this context. So whenever you're doing aggregation from a stream into like you know, average values, et cetera, or you're joining a stream in a table, or table on a table, there is some state that is being maintained behind the scenes for you. And that state is you know stored locally, but it is more like you know a temporary local storage, um, as I just described. Now, and that is probably what you're getting at then is you know, things can break. You know, my machine could die, Kubernetes could decide to just kill my pod and launch a pod elsewhere. So you know, we can't really rely on local data to be persistent and durable. Right. And what Kafka Streamstone does is it will use this concept of what a stream-table reality is, we haven't really talked about yet, but it simply means that you can turn a stream into a table and a table into a stream, etc. It will do like a streaming backup, so to speak, of all the local data into Kafka. And if need be, it restores from Kafka again. And that is, depending on your data, like a fairly quick process. In other words, whenever you're doing aggregation, joints, and so on, so the more complicated operations like to implement in a distributed system, what Kafka Streams does is it uses Kafka as the source of truth for everything that it computes. So if data needs to be restored because a container died, um, which would be like the failure scenario, or because you just added containers and you need to migrate some of the processing from the existing containers to the newer ones, this migration operation will also happen through Kafka behind the scenes, where the state data is being moved from A to B. Because moving the computation is very simple, right? You when you're starting your container, your computation logic is already inside the new container. But what the container misses is well, the data that I have to work with that this other container already processed and maintained over the past few hours. I don't want to do this from scratch again. Just give me what you had. And this is something that Kafka Streams does behind the scenes for you by using Kafka as, so to speak, a very fast backup and restore mechanism.

SPEAKER_00

Right. Which is huge, by the way, because otherwise maintaining that state yourself uh is a complex problem to solve. And that puts you in the position you thought you were an application developer, now you're also a distributed state infrastructure developer, and that's that's not money well spent. So that's a huge Yeah.

SPEAKER_01

You don't want to do that unless you really, really have to. I mean I I like these kind of things. I'm I like to read up on the concepts, you know, reading papers on distributed systems. But if I have to do that as part of my day job, uh it's great to know about it, but you don't want to deal with it yourself, like the nitty-critty details every single time you're building your application. Right.

SPEAKER_00

Unless unless it's literally what you do, unless you're that infrastructure developer, you you shouldn't be touching. Exactly.

SPEAKER_01

So and one thing I want to mention here, just real quick, um, sorry to interrupt Tim, because a lot of people ask me that question after I gave them a brief pitch here. Um, a lot of people uh think that the data that Kafka streams you know manages for for your state, for your tables, et cetera, behind the scenes, that needs to fit into RAM, like in memory. That isn't the case. As I said, the the data that is being stored temporarily uh locally inside your container or your machine, etc., um it doesn't need to fit into RAM. Like whatever local disk space is available, it will use that. So it can be gigabytes of data uh easily. And uh that has the other side effect that you can uh get away with um using cheaper instances in the cloud, because typically you know, more memory means you're paying like a An extra uh on top of uh the normal cheaper instances and you don't have to do that with Kafka streams and KSQL.

SPEAKER_00

Nice. And you said KSQL again. So coming up against time here, maybe in conclusion, top this cake off with uh KSQL. What is what is KSQLDB?

SPEAKER_01

So we said at the beginning Kafka is an event streaming platform, and KSQLDB is for us, and that's something that you know we're happy to you know to have introduced into this space, uh, it's an event streaming database. And what do we mean by that? I think a particular reason why I personally really like KSQL is that it lowers the bar to using event streaming and you know, by extension, using Kafka. Because it's still, and we talked about things that Kafka Streams does for you, but still you need to you know understand how you write a Java application, you need to package it, you know, put it into a container, run it, you know, these things. With KSQL, it's just much easier to get on and running. You can use you know SQL, like a SQL dial like that we have to build your applications in a few minutes. And that can be done by people that don't know how to write Java code, which is a great value proposition. And uh what we also added to KSQL is the difference between push and pull queries, as we call them. So a pull query is like a normal classic database query. Let's say you have a table of users and you're gonna say, please give me um Tim and everything you know that this table knows about Tim. And you get a result. That is what you've been using MySQL for, Oracle, um, whatever. But what KSQLDB also allows you to do is push queries. And push queries are more like these you know continuously running streaming queries where you're getting uh an input event stream and you're doing something to your data, like you know, create an aggregation, which would result in a table that is continuously being updated, or you're taking the stream and you're filtering sensitive data from it to have like a lower, uh lower volume stream coming out of it that is sanitized, so to speak. And this is something that you can do uh both with KSQL. And that is pretty cool because in many use cases that we see in the wild when people are starting to use Kafka, you need both of these functionalities. So it's very rare that you only do, let's say, the pull queries, you know, these classic queries that you know fetch like snapshot information, so to speak, from a table. Or applications that only do the push queries. Oftentimes you want to have both. And um, originally we started to add such functionality to Kafka streams under the name interactive queries, but with you know KSQL, it's just much easier to use. Because we wrapped all of that behind the scenes in a slightly different implementation. And so that's events all the way up through SQL. Yeah, absolutely. Including what might be cool for those of you who are listening in and just want to now play with it after this podcast. Um, there's also an integration with Kafka Connect. Uh, what does that mean? So Kafka Connect, in case you have not heard about it, Kafka Connect allows you to hook other systems into Kafka. So you to get data from let's say a database like you know, Postgres into Kafka in real time through event streams. And then once it is in Kafka, maybe get it out of Kafka again in real time to something like Elastic for dashboarding purposes. There's confluent hub, I think the URL is like hub.conflunt.io, where you can you know download these ready-to-use connectors and you can you know have dozens, I think it's more than 100 by now, of uh ready-to-use connectors that you can then use in order to move the data in and out of um you know in your KSQL setup as you're playing with it, even like even when you're doing it for the first time on your local laptop.

SPEAKER_00

My guest today has been Michael Noll. Michael, thanks for being a part of Streaming Audio. Thank you, Tim. Thanks everyone for listening in. And there you have it. I hope this podcast was helpful to you. If you want to discuss it or ask a question, you can always reach out to me at TLBurgland on Twitter. That's at T L B E R G L-U-N-D. Or you can leave a comment on a YouTube video or reach out in Community Slack. There's a Slack sign-up link in the show notes if you want to register there. And while you're at it, please subscribe to our YouTube channel and to this podcast wherever fine podcasts are sold. And if you subscribe through iTunes, be sure to leave us a review there. That helps other people discover the podcast, which we think is a good thing. So, thanks for your support, and we'll see you next time.