π Apache Kafka Crash Course With Spring Boot 3.0.x | @Javatechie β Transcript
Full transcript
- 0:00[Music]
- 0:14[Music]
- 0:26[Music]
- 0:50[Music]
- 0:57hi everyone welcome to Java techi as I
- 1:00announced before I'm going to start
- 1:02capka series from beginner level to
- 1:04advanc level and this is the part one
- 1:06video of capka series and asend of this
- 1:08tutorial to understand what is kapka
- 1:12where does kapka come from why do we
- 1:14need kapka how does it work with high
- 1:17level overview okay so without any
- 1:20further delay let's get
- 1:22[Music]
- 1:28started
- 1:30[Music]
- 1:42so let's start with what is capka if you
- 1:46open the official page of kapka you will
- 1:48find this definition Apachi kapka is an
- 1:51open-source distributed event streaming
- 1:54platform what does it mean so let's
- 1:58break down this words to understand
- 1:59understand in better way when I say
- 2:02event streaming it points to two
- 2:04different task create realtime stream
- 2:08process the realtime stream let me
- 2:11explain this to word with an example
- 2:14okay so I hope everyone uses PTM in this
- 2:18digital world if anyone does not know
- 2:20PTM PTM is a Epi payment method let's
- 2:24assume I am using that PTM application
- 2:27to do some payment or I am booking the
- 2:30flight ticket or I'm just booking some
- 2:32movie ticket okay because PTM providers
- 2:35feature to do n type of transaction now
- 2:38when I'll do any transaction that event
- 2:41will go to the capka server but I'm not
- 2:44the only one PTM user who is doing the
- 2:47transaction at this time since people
- 2:49uses this across Globe This capka Server
- 2:53receives million or billions of event in
- 2:56each minute or each second or even in
- 3:00each millisecond right so sending the
- 3:03stream of continuous data from PTM to
- 3:07the capka server is called creating
- 3:11realtime event stream or generating
- 3:14realtime stream of data now once capka
- 3:18server receives the data he need to
- 3:22process it right so the PTM team or PTM
- 3:25developer created one client application
- 3:29which will read the data from the kka
- 3:31server and do some process for example
- 3:35let's say PTM want to restrict marks 10
- 3:39transaction per day I mean a user can
- 3:42only do 10 transaction per day using PTM
- 3:46method okay if it exceed the limit then
- 3:50client application wants to send a
- 3:52notification to the user in such
- 3:55scenario my client application needs to
- 3:59continuously face the data and need to
- 4:02do the validation to check the
- 4:04transaction count for a specific user in
- 4:07each and every seconds or millisecond
- 4:10right so this continuously listens to
- 4:13the capka messages and processes
- 4:15processes them is called processing Real
- 4:19Time Event stream okay so if you combine
- 4:23these two term this will give you the
- 4:25answer what is event streaming in simple
- 4:29wordss continuously sending the message
- 4:32or even to the capka server and reading
- 4:35and process them is called realtime
- 4:38event
- 4:40streaming now let's move to the next
- 4:43word called distributed in microservices
- 4:46World distributed means distribute
- 4:49multiple computers to different node or
- 4:52region to balance the load and to avoid
- 4:55the downtime similarly as you know capka
- 4:59is a distributed event streaming
- 5:01platform we can also distribute our
- 5:04capka server if you can observe here we
- 5:07have three Capa server running in
- 5:09different region to perform event
- 5:12streaming operation in case if any
- 5:15server goes down another server will
- 5:18come up to pick up the traffic to avoid
- 5:21application downtime okay hope you
- 5:24understand what capka is or why it is
- 5:27called a distributed event streaming
- 5:31platform now let's move to the next
- 5:33point that is where does capka come from
- 5:37capka was originally developed at
- 5:40LinkedIn and was subsequently open
- 5:43sourced in early 2011 now this capka
- 5:46software comes under Apachi software
- 5:49Foundation okay now let's move to the
- 5:51next Point why do we need capka so let
- 5:56me walk you through an example to
- 5:58demonstrate this particular Point let's
- 6:00say I have some parcel on my name so
- 6:03Postman come to my door to deliver the
- 6:06parcel unfortunately I was not there in
- 6:09my home and went for vacation so he
- 6:13returned back again next time he came
- 6:15but I'm not there in home he tried two
- 6:18three attempt to deliver the parcel and
- 6:20every time I am not there in my home
- 6:23after some day he might forgot about the
- 6:26parcel or he return the parcel to the
- 6:28main office in this case I lost the data
- 6:33or I was I will not able to receive the
- 6:35information what Postman bring to my
- 6:38door that parcel could contain some
- 6:41important information or money related
- 6:44stuff but I missed it because I was not
- 6:46available during the period when the
- 6:48postman came to my door this could be a
- 6:51huge loss for me Isn't it how can I
- 6:55overcome this no worries I'm very smart
- 6:58I have installed a letter box near to my
- 7:01door so when next time parcel Postman
- 7:04brings for me and if he found I am not
- 7:08there in home or I'm not available in my
- 7:10home then simply he can put that parcel
- 7:13to my letter box so whenever I will be
- 7:16back to my home it can be a day after or
- 7:19a week after I can go to my letter box
- 7:22and I can collect the parcel of the
- 7:24messages what Postman dropped in my
- 7:27letter box in that case I will not lost
- 7:31the data data will be there in my letter
- 7:33box until I pick it up here the letter
- 7:37box acts as a middleman between the
- 7:39postman and me how cool is this isn't it
- 7:44let's try to relate this understanding
- 7:46with one real time example let's say I
- 7:50have two application application one and
- 7:52application two Now application one
- 7:55wants to send the data to application
- 7:57two but if if application 2 is not
- 8:01available to receive the data then again
- 8:04he will lost the data like me which
- 8:07might impact to the business of
- 8:09application to so to overcome this
- 8:12communication failure we might need to
- 8:14install something similar to the letter
- 8:17box between these two application right
- 8:20that is where capka comes into the
- 8:23picture you can add a messaging system
- 8:26between application one and application
- 8:28two and that messaging system can be a
- 8:31capka or it can be a rabbit mq or it can
- 8:34be a redies okay but we are going to
- 8:37focus on the capka messaging system so
- 8:41in case application 2 is not available
- 8:44he can collect the message from capka
- 8:47whenever he will come online this capka
- 8:50again will act as a letter box between
- 8:52application one and application two that
- 8:56is how he won't lost the data right hope
- 9:00now you understand why we need capka or
- 9:04messaging system now let's understand
- 9:07the need of kka with more complex
- 9:10scenario okay let's say you have four
- 9:14application who wants to produce
- 9:16different types of data to the database
- 9:20server this looks simple what is the
- 9:22problem here nothing as of now but in
- 9:26future your application can grow you
- 9:28might have n number of service to
- 9:30communicate with each other in that
- 9:33situation it's really tough to manage
- 9:36these many connections between the
- 9:39services there could be a lot of
- 9:41challenges what could be the possible
- 9:43challenges data format connection type
- 9:47or number of connection when I say data
- 9:50format this front end app may be want to
- 9:54produce some different type of data to
- 9:57each different type of app application
- 10:00maybe front end want to send some
- 10:02payload structure to database server and
- 10:04some different payload structure to
- 10:06security system similarly Hardo also
- 10:09might to do or might to send some
- 10:11different type of data there could be a
- 10:13complexity to handle the data format or
- 10:17schema and next connection type so there
- 10:20could be different type of connection I
- 10:23mean when I say different type of
- 10:25connection it could be a HTTP connection
- 10:27or it could be a TCP connection
- 10:29connection or it could be a jdbc
- 10:31connection like the connection type will
- 10:34be complex to maintain with multiple
- 10:37Services okay and second is number of
- 10:40connection when I say number of
- 10:42connection if you observe carefully
- 10:44front end connect to the five different
- 10:47destination
- 10:48Services hard connected to different
- 10:52five database SL connected to five and
- 10:54chart server also connected to five
- 10:57different services so so if you count
- 11:00the number of connection here the total
- 11:02count will be 20 so to just
- 11:06maintain left side is four application
- 11:09and right side is five application so
- 11:11just to maintain nine application or to
- 11:14just communicate between the nine
- 11:17application we need to manage the 20
- 11:21Connection in Enterprise application
- 11:24really it's kind of a bottleneck
- 11:26situation to handle above challenges
- 11:29then how to overcome it that is where
- 11:31this messaging system like capka came
- 11:34into picture so if You observe this
- 11:37particular diagram carefully now front
- 11:39end Hardo database lab chart server
- 11:42whatever the data type they want to send
- 11:45or whatever the schema structure they
- 11:46want to send or which type of connection
- 11:49they want to make they'll simply send
- 11:52the payload or information to this
- 11:54particular Kafka server or my messaging
- 11:58system now which kind of data this data
- 12:01server need he will go to the capka
- 12:04server and he'll directly get it from
- 12:05this kapka server similarly let's say
- 12:08security systems want to face some data
- 12:11given by either Hardo or database slave
- 12:14he can simply go there and check the
- 12:16data what he was looking for is there or
- 12:19not if it is there he can simply get
- 12:22those data and he can play with those
- 12:24data in this approach we are maintaining
- 12:27a centralized or similar to the letter
- 12:29box all the frontend hardw database laab
- 12:32chart server they will drop the messages
- 12:35and other five Downstream services will
- 12:39go to those kavka server and they will
- 12:41pick the messages as for their need and
- 12:43in this approach we are also reducing
- 12:46the connection count Can You observe
- 12:49here the total number of connection here
- 12:511 2 3 4 1 2 3 4 5 total nine Connection
- 12:55in previous approach we found 20
- 12:57connection but when we centralized using
- 13:00the capka we are reducing the connection
- 13:03count as well okay so this is the
- 13:06advantages of using the messaging system
- 13:09like
- 13:10capka I hope you understand why do we
- 13:14need capka server or capka as a
- 13:18messaging system now let's move to the
- 13:20next point that is how does it work I
- 13:24will just give high level overview to
- 13:27understand how this capka or messaging
- 13:29system works in real time okay so this
- 13:33specifically works on pops sub model
- 13:36when I say pop up pop stands for
- 13:38publisher sub stands for subscriber
- 13:41model okay so usually there should be
- 13:43three component publisher subscriber
- 13:46messaging system or you can say message
- 13:49broker okay publisher is the person who
- 13:52will publish the event or messages to
- 13:55the messaging system like capka okay and
- 13:59the message will go and sit in the
- 14:01message broker now the subscriber will
- 14:05go to that particular message broker and
- 14:07will ask for the messages or subscriber
- 14:10will simply listen to that message
- 14:12broker to get the messages okay so this
- 14:16is just 50,000 high level overview of
- 14:19the pops up model don't worry we'll
- 14:22understand in details how the message is
- 14:25getting processed inside the broker and
- 14:27what all components helps to process the
- 14:30message in my coming session okay so
- 14:33this is just a heads up tutorial to just
- 14:36get basic information about capka and
- 14:40its need in real
- 14:48time this is part two video of our capka
- 14:51Series in this tutorial we are going to
- 14:54discuss about capka architecture and its
- 14:57components okay all right if you
- 15:00understand the capka terminology and its
- 15:02purpose then you can picturize
- 15:05internally how a message is going to
- 15:07flow from producer to the consumer via
- 15:10broker or capka server okay so without
- 15:14any further delay let's get started so
- 15:17these are the capka components or you
- 15:19can say core concept of capka here each
- 15:22component play a different role in the
- 15:24capka ecosystem no or is we'll discuss
- 15:27one by one okay so let's begin with the
- 15:30first one that is producer if you
- 15:33understand the popup model in our
- 15:34previous video the producer is the
- 15:37source of data who will publish the
- 15:39messages or events the consumer act as a
- 15:44receiver it is responsible for receiving
- 15:47or consuming a messages but they won't
- 15:50directly communicate with each other to
- 15:53process the messages from producer to
- 15:55the consumer there should be one
- 15:57middleman between them that is what
- 15:59called capka server or broker in pop sub
- 16:03model usually producer will publish or
- 16:06push the messages to the broker then
- 16:09broker will store that messages and then
- 16:11consumer will consume it from the broker
- 16:14in other word you can say a broker is
- 16:17just an intermediate entity that helps
- 16:20in message exchange between a producer
- 16:23and a consumer okay so this is what the
- 16:27basic key component of your pops sub
- 16:29model and producer consumer and broker
- 16:33let's move to the next one that is
- 16:35cluster so I believe everyone heard the
- 16:38word called cluster if not a cluster is
- 16:42a common terminology in the distributed
- 16:44computing system it is nothing but just
- 16:46a group of computers or servers that are
- 16:49working for a common purpose since capka
- 16:52is also a distributed system it can have
- 16:55a multiple capka server or Brokers
- 16:58inside a single capka cluster so there
- 17:01can be one or more capka Brokers inside
- 17:04a single capka cluster let's assume your
- 17:08producer is publishing huge volume of
- 17:10data then a single Capa broker may not
- 17:12able to handle the load right so you
- 17:15might need to add additional capka
- 17:16server or additional broker so you can
- 17:19group multiple broker using a capka
- 17:22cluster fine so next let's Deep dive
- 17:25inside the Brokers to understand what
- 17:28all core components it has to store the
- 17:31messages okay so let's move to the next
- 17:34component that is topic so before I
- 17:37explain the theory what is the topic
- 17:40let's find out why it is there inside a
- 17:42broker let me walk you through the same
- 17:45PTM example as you understand the PTM
- 17:48application will keep sending different
- 17:50type of activity to the capka server and
- 17:54The paytm Client app will continuously
- 17:56listen the messages from the capka
- 17:59server right so here the left part is my
- 18:04producer and the right side is my
- 18:06consumer and the middleman is the broker
- 18:09or you can say capka server so broker
- 18:13receives different type of messages or
- 18:15events here it can be a payment
- 18:18transaction or it can be a ticket
- 18:20booking or it can be a specific to the
- 18:22mobile rearch related message right now
- 18:26consumer needs to read all the messages
- 18:29to process them so the consumer ask to
- 18:32the broker hey give me all the messages
- 18:35you have received then the broker checks
- 18:38his storage and replies to the consumer
- 18:41hey consumer I have a lot of messages
- 18:43with me which one you need again
- 18:47consumer ask hey broker give me all
- 18:51payment specific
- 18:52messages again broker attack with the
- 18:55counter question hey buy I have
- 18:58different type of payment messages with
- 19:00me which one you need now it's gone
- 19:04right the consumer Hurts by broker
- 19:06counter question this is obvious right
- 19:09then how we can overcome this back and
- 19:12forth conversation that is where topic
- 19:15came into picture to categorize
- 19:18different type of messages so how topic
- 19:22can help here you can simply create
- 19:26multiple topic to store different type
- 19:29of messages even you can rename the
- 19:32topic payment topic booking topic and
- 19:34insurance topic if producer sending
- 19:37payment specific messages then the
- 19:39message will push to this particular
- 19:41payment topic If the message type or
- 19:44producer is sending booking related
- 19:46messages you can keep those messages
- 19:49inside this booking topic or similarly
- 19:52if producer is sending some mobile re
- 19:55specific information you can keep it
- 19:56separate topic this is how you can
- 19:59segregate different event type to
- 20:02different topic so here consumers have
- 20:06no need to have a back and forth
- 20:07communication with Brokers whatever
- 20:10messages consumer want he will simply
- 20:13needs to subscribe that specific topic
- 20:15for example I have a consumer who needs
- 20:18only the booking specific messages then
- 20:21I'll tell him hey consumer you go and
- 20:23check directly in the booking topic you
- 20:25no need to ask to the broker to give me
- 20:28booking specific messages you simply go
- 20:30to the booking topic or you simply
- 20:32subscribe to the booking topic and get
- 20:35all the messages what you need this is
- 20:37very simple right alternatively you can
- 20:40consider the topic act as a database
- 20:43table in Capa ecosystem so usually we
- 20:47store employee information in the
- 20:49employee table and payment information
- 20:51in the payment table in database world
- 20:54right similarly when producer send
- 20:57employee specific event or messages we
- 21:00need to store them in the employee topic
- 21:02and similarly if producer is publishing
- 21:05payment specific information then we
- 21:07need to store it into the payment topic
- 21:10right now if you have multiple consumer
- 21:13then they can go to a specific topic and
- 21:16consume the messages as for their need
- 21:19for example if consumer one wants to
- 21:22read all employee messages he can go and
- 21:25collect the messages from the employee
- 21:28topic Direct ly he know need to ask to
- 21:30the broker hey broker give me the
- 21:31employee specific messages he know that
- 21:35employee topic contains all the employee
- 21:37specific messages consumer one can come
- 21:40to the employee topic and can get the
- 21:42info similarly all the consumers can
- 21:45subscribe to the corresponding topic
- 21:47what information they want to listen
- 21:50Okay so if you look into its actual
- 21:52definition topic specifies the category
- 21:55of the messages or the classification of
- 21:58the message
- 21:59listeners can then just respond to the
- 22:01message that belongs to the topic they
- 22:03are listening on so consumer one want
- 22:07the payment info he can go and listen to
- 22:10the payment topic that is what the need
- 22:13of topic or this topic play a critical
- 22:15role in capka ecosystem okay so let's
- 22:19move to the next key component that is
- 22:22partition now we know that PTM producer
- 22:26sends data to the broker and the broker
- 22:29store the message inside a topic now
- 22:32consider you have a huge volume of data
- 22:35when I say huge volume of data let's
- 22:38assume your producer is publishing
- 22:41millions or billions of messages per
- 22:43second to the topic then will my topic
- 22:46be able to handle those many messages
- 22:49obviously no there could be a storage
- 22:52challenges or it can be very challenging
- 22:54for the topic to store data on a single
- 22:57machine right then how to handle this
- 23:00situation solution is very simple as we
- 23:04already know capka is a distributed
- 23:06system in such a scenario we can break
- 23:09the capka topic into multiple parts and
- 23:13distribute those parts into different
- 23:16machines this concept is called topic
- 23:20partitioning and each part is called
- 23:23Partition in capka terminology okay so
- 23:26this partitioning will give give you
- 23:28better performance and high availability
- 23:31because if You observe we split our
- 23:34topic to multiple partitions so when
- 23:37producer publish bulk messages then each
- 23:40partition concurrently accept the
- 23:43messages that will definitely improve
- 23:45the performance right even in case if
- 23:49any partition goes off then other
- 23:51partitions are available to handle the
- 23:54loads without any application downtime
- 23:57this is the reason topic partitioning is
- 23:59the key aspect in the capka ecosystem
- 24:02and we can decide the number of
- 24:04partitions for a topic during capka
- 24:07topic creation okay no worries we'll
- 24:10hands on with these terms in our coming
- 24:12videos so now let's move to the next
- 24:15component that is offset so far we
- 24:19understand a capka cluster can have one
- 24:21or more capka server then a capka server
- 24:25can have multiple topics and each topic
- 24:28can have multiple partitions right now
- 24:32when producer send a message then it
- 24:35will go and sit in any partition inside
- 24:38a topic that we don't have any control
- 24:41on it it will work based on the round
- 24:43robin principle so if You observe as
- 24:46soon as message arrives in a partition a
- 24:49number is assigned to that message that
- 24:52is called offset 0 1 2 3 4 5 these are
- 24:56the sequence number as send to each and
- 24:59every messages that is called offset so
- 25:02the purpose of offset is to keep a track
- 25:06of which message have already been
- 25:08consumed by the consumer for example
- 25:11let's consider after reading four
- 25:14messages from the partition consumer
- 25:17went down okay then when consumer will
- 25:20back to the online then this offset
- 25:23value would help him to know exactly
- 25:26from where the consumer has to start
- 25:28consuming the messages now Here If You
- 25:31observe consumers should start reading
- 25:33message from offset 4 not from the zero
- 25:37again right this is what the key role of
- 25:41offset fine now let's move to the next
- 25:44component that is consumer group so I
- 25:49believe we all are clear up to this
- 25:51producer will push huge volume of
- 25:53messages to the capka topic and message
- 25:56will again split to different different
- 25:58partitions and inside partition each
- 26:01message have their own identity with
- 26:04offset number so now let's go one step
- 26:08ahead who will listen to the messages
- 26:10from the partitions the single statement
- 26:13is consumer right but if You observe
- 26:16here I have three partitions and a
- 26:19single consumer is reading from each and
- 26:22every partition which will definitely
- 26:24leads the performance issue because
- 26:27there is no concurrency you can see here
- 26:30a single consumer is playing with all
- 26:32the three partition so definitely it
- 26:35won't give us the better throughput
- 26:37right then how we can overcome this
- 26:39situation it's very simple just share
- 26:43the workload how can I share the
- 26:45workload so what I will do simply I'll
- 26:48Define n number of consumer instance now
- 26:52I can group all the three consumer what
- 26:54I Define into a single unit by specific
- 26:58speing the group name that is called
- 27:00consumer group now payment consumer
- 27:03group have the three consumer instance
- 27:06now since I have multiple consumer with
- 27:08me I can divide the workload to each and
- 27:11every consumer to achieve better
- 27:13throughput so simply now my consumer one
- 27:17can read from partition zero consumer 2
- 27:20can read all the messages from partition
- 27:22one and consumer 3 can read from the
- 27:25partition two then in this approach all
- 27:28the three consumer instance read
- 27:31parallely from each and every partition
- 27:33which will definitely give me the better
- 27:35performance or we can get the better
- 27:39response right just keep a note we can't
- 27:42guarantee about consumer and partition
- 27:44order any consumer can talk to any
- 27:48partition that will be decided by the
- 27:51coordinator okay so you can say consumer
- 27:54one might listen from partition two or
- 27:57any order okay we don't have any control
- 28:00on it cool now you might have a question
- 28:03hey I have three partition so three
- 28:06consumer so they will talk to each other
- 28:08that's correct but what if I have a
- 28:11fourth consumer added to my consumer
- 28:13group then what will be the behavior of
- 28:15this consumer 4 what he will do right
- 28:18there is no change that consumer 4 will
- 28:22sit ideal because there is no work for
- 28:25him now since all the three partitions
- 28:28assigned to three different consumer
- 28:30instance there is no partition left for
- 28:33him so he'll simply sit ideal but in
- 28:36case if any consumer instance is
- 28:39rejected or goes off then C4 will get a
- 28:42chance to connect with any partition
- 28:45this concept is called consumer
- 28:47rebalancing don't worry I'll cover this
- 28:50part in the coming session okay fine now
- 28:53let's move to the next component that is
- 28:57Juke keeper
- 28:58so zuke keeper is a key component for
- 29:01capka capka is a distributed system and
- 29:04it uses zuke keeper for coordination and
- 29:08to track the status of your capka
- 29:10cluster okay it also keeps the track of
- 29:13capka topics partitions offset and all
- 29:17the component of the capka ecosystem so
- 29:20in single word you can say this
- 29:22particular jeeper will act as a manager
- 29:26for your broker or for your capka
- 29:29cluster okay so this is what all the key
- 29:33component of capka we have seen the
- 29:35concept of capka architecture moreover
- 29:37we discussed capka components and basic
- 29:40concept so furthermore for any query
- 29:43regarding architecture of capka please
- 29:45feel free to ask in the comment section
- 29:48we'll discuss more internal of each
- 29:50components in our coming session
- 29:57okay
- 29:58so before we install capka let me tell
- 30:01you that we have different kind of capka
- 30:03server available in Market with
- 30:05different type of flavor like Apachi
- 30:08capka which is open source commercial
- 30:11distribution managed capka service so
- 30:14let me give you some heads up about
- 30:16above three capka flavor so that you can
- 30:18able to understand which capka service
- 30:20need to choose based on the situation or
- 30:23based on the scenario okay so the first
- 30:26one is open source version of Apachi
- 30:28kapka this you can easily download from
- 30:31Apachi portal and then you can install
- 30:34and use it in case if you face any
- 30:37operational issue or open bug on capka
- 30:40then you need to manage it by your own
- 30:42or you might need to upgrade if needed
- 30:45so maximum Company still use this open
- 30:48source with proper infrastructure setup
- 30:50and also they have dedicated developers
- 30:53or experts to handle any kind of infr
- 30:56related issue so the second one is the
- 30:59commercial distribution that comes with
- 31:02a lot of tools and utility to perform
- 31:05your day-to-day capka operation this
- 31:08will add cost to your organization or to
- 31:10the project and conflent capka is one of
- 31:14the best commercial distribution to use
- 31:17which simplifies connecting data source
- 31:19to the capka building streaming
- 31:21application as well as securing
- 31:23monitoring and managing your capka infra
- 31:26so today if if you will observe conflent
- 31:29platform is used for a wide array of use
- 31:32case across numerous Industries and
- 31:35conflent capka also provide a Community
- 31:37Edition for developers which is
- 31:40completely free don't worry I will show
- 31:43you how to use the community edition of
- 31:45confluent capka next is managed capka
- 31:49service this comes with everything you
- 31:52need you just create instance as per
- 31:54required configuration let everything
- 31:57managed by the the Capa provider all
- 31:59infrastructure will be ready for you and
- 32:02it really easy to scale when needed as
- 32:04similar as Cloud infra so Amazon msk and
- 32:08conflent are best managed capka service
- 32:10provider to use So based on the project
- 32:13requirement and budget availability you
- 32:15need to decide which capka service need
- 32:17to choose okay but as part of this
- 32:20course we'll use the open source apachi
- 32:23kapka but don't worry I will also
- 32:25demonstrate how to use the community
- 32:28edition of Kafka confluent okay so let's
- 32:32begin with the installation steps so
- 32:34I'll walk you through the steps to
- 32:36install Apachi capka which is open
- 32:39source also we'll install the community
- 32:41edition of confident capka then at the
- 32:44end we'll install capka upset Explorer
- 32:47which will helps us to monitor our capka
- 32:49messaging system okay so let's begin
- 32:52with the first one that is installation
- 32:54of apachi kapka so we are going to
- 32:57download the open source of apachi kapka
- 32:59so for that go to the Chrome
- 33:02browser and then you can search
- 33:05here
- 33:07capka go to the official page of apachi
- 33:11kapka then you can find the option to
- 33:13the download click on download capka so
- 33:16the latest version you can find your
- 33:193.4.0 okay so you can download any of
- 33:22them this is the binary distribution so
- 33:24let me click on the first one
- 33:29the file size is 102 MB so it will take
- 33:32few second to
- 33:33download So once the download will be
- 33:36completed go to that particular
- 33:38folder so you need to extract this
- 33:41particular
- 33:44folder you can see here the version we
- 33:47downloaded
- 33:493.4.0 now go inside this particular
- 33:51folder you can find bin config leaps
- 33:55license and side docks okay so inside
- 33:59the bin you'll find all the Sal script
- 34:02file for mac and all the batch file for
- 34:05Windows so if I open this you can see
- 34:08here all the sh extension is being used
- 34:12for the Linux or Mac operating system
- 34:14and in the windows you can find the
- 34:17batch file to start the Juke keeper or
- 34:20to start the capka server or to start
- 34:22the topic or to create the topic all the
- 34:25things whatever we have discussed before
- 34:27you can can play by executing this batch
- 34:29file or by executing the Sal script file
- 34:32so once you will download the binary
- 34:34distribution you will find the support
- 34:36for both operating system this is for
- 34:39Windows and this C script is for Mac or
- 34:42Linux okay and if you'll go inside the
- 34:45config it contains the properties file
- 34:48okay so you can see here server.
- 34:51properties which holds the capka
- 34:54property specification and you can see
- 34:56here jer. properties okay so these are
- 34:59the properties file you will understand
- 35:00once we'll play with those properties
- 35:02file now if you'll go here inside the
- 35:06leaves there are couple of Jar that's it
- 35:09okay you just need to download this
- 35:11binary distribution then all set in
- 35:14upcoming session I will show you how you
- 35:17can play with the capka component by
- 35:19executing this script okay so let's move
- 35:22to the next one that Community Edition
- 35:26of conflent capka okay so to install the
- 35:29confluent capka what you can do you can
- 35:32go to the confluent
- 35:37doio then you can click on get started
- 35:42free and you want to install the
- 35:45software which is community edition of
- 35:47confluent capka we don't want any
- 35:49managed capka service so navigate the
- 35:52tab to the
- 35:53softwares next you can fill this form
- 35:57and then just acknowledge these two
- 36:00click on start
- 36:04free here we go so if You observe here
- 36:08there is a local install confident
- 36:11platform on a single local Machine by
- 36:14using G tar archive or Docker images so
- 36:16we don't need this what do we need we
- 36:18need the Community Edition okay there is
- 36:21something called distributed confident
- 36:23for
- 36:23kubernetes and there are few more so we
- 36:26need to focus on this community which is
- 36:29completely free so you can see here use
- 36:32it for free forever okay so we want to
- 36:36install the Jeep or we just want to
- 36:38download the Jeep just click on this
- 36:43download this is the size 46 MB okay so
- 36:47let's wait it to complete so once it
- 36:51downloaded just show in the folder and
- 36:55again you need to extract this
- 37:03you can see this right the lest version
- 37:06of Community Edition of conflent capka
- 37:09is
- 37:107.3.2 fine now go inside this particular
- 37:13folder you can find the same structure
- 37:16that is B and you can find the file
- 37:20which is the Sal script for mac and
- 37:22Linux and for Windows you can file this
- 37:25batch file okay and and if you will move
- 37:28further inside this Etc as I already
- 37:32mentioned once you are playing with the
- 37:34conflent capka apart from capka you'll
- 37:37find lot of utility when I say utility
- 37:40you will find kapka rest kql DV rest
- 37:44utils schema registry so we'll
- 37:47understand about this particular utility
- 37:49in our upcoming session okay now if
- 37:51you'll go inside the capka you'll find
- 37:54the properties file Okay so
- 37:58whatever files you downloaded through
- 38:00the Apachi kapka open source and
- 38:03conflent each file is similar only the
- 38:06structure is changed so let me compare
- 38:09this so that it will be easy for you to
- 38:11understand okay let me open the capka
- 38:14and meanwhile I'll open the another
- 38:18window
- 38:19for confluent capka
- 38:25okay now let me
- 38:29so this is the capka and this is the
- 38:31conflent capka folder now if we compare
- 38:34these two folder if you'll go inside bin
- 38:38you'll find the script file and Windows
- 38:40File and again if you'll go inside this
- 38:43you'll find the same okay now if we go
- 38:46back and if you'll check about the
- 38:48config you'll find all the properties
- 38:50file but in conflent capka since it
- 38:53supports for multiple utility methods in
- 38:57inside the ETC you will find a folder
- 38:59called capka inside that capka you can
- 39:02find the properties okay and these are
- 39:05the utility fine this is the only
- 39:09structure changed you will find between
- 39:11these two capka and conflct kapka okay
- 39:14so as I mention I will just show you how
- 39:18you can play with both okay that is my
- 39:21next session I will just demonstrate
- 39:23using the command line interface how we
- 39:26can play with all the capka component in
- 39:29open source Apachi capka and in
- 39:31confident capka fine now the next we
- 39:35need to install the capka upset Explorer
- 39:38to monitor our capka messaging system so
- 39:41what you can do for that go to the
- 39:44Chrome and just search here capka upset
- 39:48Explorer the first link you can find
- 39:50something called capka tool okay so just
- 39:53open this
- 39:55link and then let me Zoom this for
- 39:59you you can see here right offset
- 40:02Explorer this is the caput tool is a GUI
- 40:05application for managing and using aaji
- 40:07capka cluster okay so just click on the
- 40:11download you'll find the option to
- 40:14download for Windows and Mac so download
- 40:18Based on your operating system so since
- 40:20I'm using the Mac I'll download this
- 40:23particular Mac
- 40:25OS so those size will be 60.2 MB so once
- 40:31it will be download you just need to
- 40:33double click to install it okay so let's
- 40:36wait it to
- 40:38complete so it downloaded just show in
- 40:42the
- 40:43folder and you just need to double click
- 40:46on it okay so that it will be installed
- 40:48to your
- 40:49machine drag and drop I have already
- 40:53installed it so it is asking me whether
- 40:55you want to replace or you want to keep
- 40:56the book
- 40:57so I'll just stop for now okay because
- 41:00already I installed it so I'll also
- 41:02remove for
- 41:03now fine so this is what about the
- 41:06installation steps in my coming tutorial
- 41:09I will show you how you can play with
- 41:12the capka component using the common
- 41:14line interface we'll create the consumer
- 41:16and producer we'll publish the event to
- 41:18the topic and we'll add the partition
- 41:20and we'll show the various type of
- 41:22behavior in this particular producer and
- 41:24consumer flow okay
- 41:33so we already downloaded kka binary
- 41:35distribution right now what next first
- 41:38we need to create or start a capka
- 41:40ecosystem for that we need to start a
- 41:43capka server or broker okay but to
- 41:46manage This capka Server we need a
- 41:48manager which is Zookeeper so the steps
- 41:51we need to follow start the Juke keeper
- 41:54then next we need to start the capka
- 41:56server or broker once you have the capka
- 41:59ecosystem ready with you then we can
- 42:01start producer and consumer since we are
- 42:04playing with common line interface we'll
- 42:06use the command to create producer and
- 42:08consumer for now but in my upcoming
- 42:10session I will create two different
- 42:12application for producer and consumer
- 42:15okay cool next producer wants to send a
- 42:19message to the consumer but they can't
- 42:21directly communicate with each other so
- 42:24we need another component called topic
- 42:27so that producer can push the message to
- 42:30the topic and consumer can consume it
- 42:33okay once you have created the topic
- 42:35next you can Define n number of
- 42:37partitions to distribute the load coming
- 42:39from producer for example let's say my
- 42:42producer wants to send 10 messages then
- 42:45those message will distribute to all
- 42:47three partitions assigning with upset
- 42:49number so while creating a topic you
- 42:52need to specify the partition count you
- 42:55need to take the decision based on the
- 42:57load coming to your application how many
- 42:59partition will be perfect fit for your
- 43:02topic okay that decision you need to
- 43:04take while creating the topic once a
- 43:07message reach to the partition then
- 43:09immediately whoever consumer listen to
- 43:12it they can happily consume it so this
- 43:14is the typical producer and consumer
- 43:16flow we'll follow the same steps to Deep
- 43:19dive further okay so using the common
- 43:22line interface we'll start the Juke
- 43:24keeper capka server we'll create the
- 43:26topic then I'll will show you how you
- 43:28can Define the partition then also we'll
- 43:30understand how the bulk message is being
- 43:33distributed to the multiple partition
- 43:35and how it is being consumed by the
- 43:37consumer okay so to play with the
- 43:39command line interface you need to go to
- 43:41the folder where you install the capka
- 43:44so it is there in my download folder I
- 43:47have created a folder called softwares
- 43:50inside this I have kept these two capka
- 43:53okay so I'll show you how you can do
- 43:55using capka also I will demonst
- 43:57how we can do using the confluent capka
- 44:00Okay cool so first step what we
- 44:03understand first we need to start the
- 44:07Juke keeper right then you need to start
- 44:09the capka server so what you can do go
- 44:13to the folder then go to the capka open
- 44:16a new
- 44:18terminal fine then I will also create
- 44:22couple of
- 44:26terminal
- 44:29just go
- 44:31here new terminal fine let's go with
- 44:35this three let me minimize this
- 44:43okay fine so first I need to start the
- 44:48Juke keeper so how can I start the Juke
- 44:50keeper you need to fire the command so
- 44:53there is a cell script file you will
- 44:55find inside the cap B directory so if
- 44:58you'll go inside kapka if you'll go
- 45:01inside B you can find your something
- 45:04called Juke keeper server hyphen
- 45:07start.sh okay so what we can do we'll
- 45:11run this particular command so this is
- 45:14inside B folder Juke
- 45:17keeper start
- 45:20server
- 45:21start.sh okay now to start the Juke
- 45:24keeper it needs to read the zuke keeper.
- 45:27properties file so again if you'll go
- 45:30inside the config you'll find a property
- 45:34called Juke keeper. properties can you
- 45:36see here this properties it needs to
- 45:39start the Juke keeper server so you just
- 45:42need to Simply give the path of it it is
- 45:45inside
- 45:46config let me Zoom this for you it is
- 45:51inside
- 45:53config and the file name is Juke keeper
- 45:56do properties file now just enter
- 46:02it it will take few second to start the
- 46:04Juke keeper and the default Port it will
- 46:07run
- 46:082181 okay so I can show you that let me
- 46:12filter it
- 46:142181 can you see here the default Port
- 46:17is
- 46:182181 now as per the presentation we
- 46:22understand after start the Juke keeper
- 46:24we need to start the capka server so so
- 46:27go to another terminal let me use the
- 46:29second one here you need to start the
- 46:32capka server how can you start the capka
- 46:35server again if you'll go and check in
- 46:37your B
- 46:40folder you can find something called
- 46:43capka server start.sh this is the Sal
- 46:47script file if you're using Windows you
- 46:49can go insert the windows and you can
- 46:51search
- 46:53for capka where is the capka server
- 46:59yeah this one you can run this capka Hy
- 47:02server start. batch file okay since I'm
- 47:05in the windows sorry I'm in the Mac I'm
- 47:08using the s file so just go to
- 47:12the
- 47:14terminal and just run it is inside bin
- 47:19okay capka hen
- 47:21server
- 47:23start.sh and to start the capka server
- 47:26it needs to read the properties file
- 47:28called server. properties so if you'll
- 47:31go inside the properties file or config
- 47:33folder you'll find something called
- 47:35server. properties okay so you need to
- 47:38give this properties file to start the
- 47:40capka server so again give the path
- 47:43config
- 47:44SL server. properties just enter
- 47:53it yeah it take few second to start your
- 47:57capka
- 47:58server now to check which Port This
- 48:01capka Server is running you can filter
- 48:04something
- 48:069092 can you see here this is the
- 48:09default Port of capka server so just
- 48:13keep a note the default Port of Juke
- 48:22keeper is
- 48:242181 and capka
- 48:28server or
- 48:31broker is
- 48:339092 okay so that's it now our capka
- 48:37ecosystem is ready now what is the next
- 48:40step if you understand correctly then
- 48:42the next step producer and consumer need
- 48:45to communicate with each other using the
- 48:47topic so you need to create a topic and
- 48:50we need to specify the partition and
- 48:52replication Factor okay so to run a
- 48:55topic again I'll take the help from this
- 48:58common line
- 48:59interface okay now if we go inside the B
- 49:04folder you'll find something capka Hy
- 49:08topic.sh okay so what I can do I'll just
- 49:12run here I need to go inside bin then
- 49:15capka hyen topic. SS also I need to tell
- 49:21here while creating the topic I need to
- 49:24give my bootstrap server host and put
- 49:27when I say bootstrap server it's about
- 49:29your capka server okay now I'll
- 49:33Define
- 49:36bootstrap
- 49:37stap do server my bootstrap server is
- 49:42running on host Local Host and Port is
- 49:469092 this is where my capka server is up
- 49:49and running right so let me make this
- 49:53white so that we can type in a single
- 49:55line okay I have ran the capka hyen
- 49:58topic.sh and I giving where is my capka
- 50:00server is up and running the host is
- 50:02Local Host and Port is
- 50:059092 now what action you want to do I
- 50:08want to create a
- 50:11topic okay and what is the topic name
- 50:14you want to give let's say I want to
- 50:17give JT topic or let's give a valid name
- 50:23Java underscore or
- 50:28topic okay now how many partition you
- 50:32want to Define so if you want Define
- 50:34anything it will take the default
- 50:36partition as a one but you have option
- 50:38to Define n number of partition so as of
- 50:42now I just want to Define three
- 50:44partition
- 50:46okay
- 50:48partition count equal to three now how
- 50:52many copy you want to keep that is
- 50:54called replication Factor right since we
- 50:56are running on a single broker so let's
- 50:59keep the replication Factor as one so
- 51:01you can
- 51:03Define
- 51:06replication Factor equal to 1 so there
- 51:10is nothing complex to remember this
- 51:12command or this particular command
- 51:15whatever you have type here I'm just
- 51:17running my topic by giving the
- 51:19information where is my Capa server is
- 51:21up and running by defining the host and
- 51:23port and what action you want to perform
- 51:26I want to create a topic and the name of
- 51:28the topic is Java hen topic and the
- 51:31number of partition I want to assign as
- 51:33three to this particular topic and I
- 51:35want to keep a single clone of this
- 51:38particular broker so that is the reason
- 51:40I have defined replication Factor as one
- 51:42now click on
- 51:45enter so the mistake we have done we
- 51:49know need
- 51:51to add the dot it need to be hyphen okay
- 51:55so let me Zoom this for
- 51:58you now let me drag this yeah now let me
- 52:03enter
- 52:06it can you see here it created the topic
- 52:10Java hyen topic you can Define n number
- 52:13of topic if you want so let me create
- 52:16another
- 52:18topic topic one
- 52:25okay now now we have created two topic
- 52:28so to list down all the topic available
- 52:31inside your capka server what you can do
- 52:35you can simply run bin I mean same
- 52:38command capka topic.sh in this
- 52:41particular server just give me the list
- 52:43of topic okay so I can happily copy this
- 52:48command
- 52:49then you can run it here and just ask
- 52:53him to give me the list of
- 52:55topic now if you'll enter we'll find the
- 52:58two topic javate topic javat topic one
- 53:03now if you want to describe this topic
- 53:05to know what is the number of partition
- 53:08replication Factor who is the leader so
- 53:10if you want to describe this topic that
- 53:13is also straightforward you can run a
- 53:16command so what you can do run the same
- 53:19command and ask him to describe the
- 53:24topic name so I just want to describe
- 53:27javat hyphen topic okay this is what the
- 53:31topic name we have right I just want to
- 53:33describe this run
- 53:39it what is the
- 53:42error okay I need to Define that I want
- 53:46to describe a topic right so describe
- 53:50then you can tell him what you want to
- 53:51describe I want to describe
- 53:55topic
- 53:59there is something
- 54:00wrong I might be doing some spelling
- 54:03mistake let's see kka event topic.sh
- 54:06bootstrap server
- 54:099092 okay I don't want this list because
- 54:13I want to describe so this is the
- 54:15mistake I mean you need to understand
- 54:18this particular command okay you just no
- 54:20need to remember it just understand what
- 54:23you want to perform okay now just enter
- 54:26it
- 54:28can you see here okay let me take the
- 54:31background as
- 54:34different fine now if You observe
- 54:37here the topic name is this and the
- 54:41number of partition we have Define three
- 54:430 1 and two okay and replica is zero so
- 54:49you can see here partition count equal
- 54:51to three replication Factor equal to 1
- 54:54there is no additional config given so
- 54:56this is empty okay you can describe the
- 54:59topic if you want and also we have
- 55:02defined the number of partition so we
- 55:05are done up to the creation of topic and
- 55:08defining the partition now we want to
- 55:10play with the producer and consumer so
- 55:13since we are playing with the common
- 55:14line interface I will run the consumer
- 55:16and producer application using the
- 55:19common line interface or using the
- 55:21command okay so we'll run the producer
- 55:24we'll push some message to the top pick
- 55:26and we'll see whether it is being
- 55:27consumed by the consumer or not next
- 55:30step we'll push bulk number of messages
- 55:33from the producer to the topic and we'll
- 55:36see whether the bulk messages is being
- 55:39getting distributed to the different
- 55:41partition or not how we can monitor
- 55:44those things we have installed the
- 55:45offset Explorer right so what I can do I
- 55:49will open the offset
- 55:55Explorer
- 55:58so what you can do here very simple step
- 56:02you just need to create a connection and
- 56:05to create a connection you just need to
- 56:07give a name of your cluster and where is
- 56:10your zookeeper is up and running and
- 56:12Port of your zookeeper so to not confuse
- 56:16you what I'll do I will create a new
- 56:21connection let's say JT I'll give the
- 56:25cluster name
- 56:26JT hyen con okay Local Host 2181
- 56:31everything it will take by default okay
- 56:33so I'll give Java techy hyphen new
- 56:37that's it click on
- 56:40test
- 56:42yes now if I open this I have the
- 56:46connection to the cluster so that I can
- 56:48monitor everything going on this
- 56:50specific cluster okay now if I expand
- 56:53this how many topic we can see here Java
- 56:57topic and Java topic one so this is the
- 57:00two topic we created and we can able to
- 57:03see inside this particular Juke keeper
- 57:05okay now if you'll go and click on the
- 57:08data and if you'll hit this particular
- 57:10thing there is no data is getting
- 57:13filtered if you'll check the topic one
- 57:16there is no data right we want to push
- 57:19the data and we want to verify whether
- 57:21it is going to the topic or not if it is
- 57:24going to the topic then we'll verify
- 57:25whether is being consumed by a consumer
- 57:27or not okay so what we can do quickly we
- 57:32just need to start a producer and then
- 57:34we just need to start a consumer fine so
- 57:37what we can do to start a producer
- 57:40you'll find something called Capa
- 57:42console producer. given by the capka
- 57:45binary go inside the bin and if you will
- 57:49search
- 57:50here can you see here there is something
- 57:53called capka console producer. to just
- 57:57kicked up by producer application and
- 57:59capka console
- 58:01consumer. to run this command to behave
- 58:03as a consumer okay so we'll run these
- 58:06two command okay so what you can do go
- 58:09to the
- 58:11terminal this
- 58:13guy now I'll clear everything what I
- 58:17want to do here I want
- 58:19to let me what okay I'll minimize these
- 58:23two okay we started the Juke keeper and
- 58:25we started the capka server let's
- 58:27minimize it now to make it clear I'm
- 58:31just running the producer application
- 58:33here so you just need to run go inside
- 58:37bin
- 58:40capka console right console hyen
- 58:44producer Dosh and then you need to give
- 58:49all the list of kka server available
- 58:51okay so you can Define broker hyen list
- 58:55current ly we have a single broker which
- 58:58is running on Local Host
- 59:019092 but in real time you can find n
- 59:04number of Brokers right so that is the
- 59:06reason you need to Define them inside
- 59:08this list now to which topic you want to
- 59:11push your messages from producer you
- 59:14need to Define that topic name so our
- 59:17topic name is javat
- 59:19hyphen topic okay that's it just enter
- 59:25it
- 59:27now it is asking you to push some
- 59:29messages so we can push the messages and
- 59:32we can verify in the uh capka upset
- 59:36Explorer but meanwhile let me start the
- 59:38consumer application so that parall we
- 59:41can push the messages from here it will
- 59:43go to the topic then consumer will
- 59:45consume from it so I want to demonstrate
- 59:47the complete pops of flow so what I'll
- 59:50do I'll open a new terminal to start a
- 59:54consumer so go go here and I just want
- 59:57to open a new
- 59:59terminal to start a consumer what you
- 1:00:02can do go inside B then simply let me
- 1:00:06Zoom this for you go insert bin then
- 1:00:10run
- 1:00:12capka console consumer Dosh okay now
- 1:00:17where is your bootstrap to whom you want
- 1:00:19to listen you can
- 1:00:22Define
- 1:00:23bootstrap hyen server
- 1:00:28which is running on
- 1:00:29Port Local Host
- 1:00:349092 and to which topic you want to
- 1:00:36listen I want to listen from the topic
- 1:00:41javat hyphen topic whether you want to
- 1:00:44listen from the beginning if yes just
- 1:00:47Define I want to listen from
- 1:00:54beginning okay don't do any spelling
- 1:00:57mistake while typing the command
- 1:00:58otherwise it will give you the error
- 1:01:01fine now let's run this
- 1:01:06consumer there is no messages being
- 1:01:09showing here so what I'll do to make it
- 1:01:13simple let me split into the
- 1:01:17two terminal so that we can see whether
- 1:01:21what we are sending we are able to
- 1:01:23receive it from the consumer or not so
- 1:01:25let me split these two
- 1:01:28tab
- 1:01:30cool now the left side is my producer
- 1:01:34and the right side is my consumer and
- 1:01:37they both can connected using this
- 1:01:40particular topic javat H topic okay
- 1:01:43producer is publishing to this topic and
- 1:01:46my consumer also listening to this
- 1:01:48particular topic if you'll do the
- 1:01:50spelling mistake in the syntax or in the
- 1:01:53command it don't work so let's see
- 1:01:56whether we are correct or not let me
- 1:01:58send some messages
- 1:02:00hi can you see here we're able to see
- 1:02:02the messages in right side let me Zoom
- 1:02:05this for
- 1:02:06you let me send something
- 1:02:09hello we can see here let me send
- 1:02:12something Java
- 1:02:15TIY we can see here let's send
- 1:02:22welcome let me some random number okay
- 1:02:27so this is my producer application from
- 1:02:29producer I'm sending the messages okay
- 1:02:33and this is the right part is my
- 1:02:35consumer let me minimize right part is
- 1:02:38my consumer he's able to consume it now
- 1:02:41whether this message is being passed
- 1:02:43through the topic or not how you can
- 1:02:45verify go to the offset Explorer now
- 1:02:49just refresh your topic okay we have
- 1:02:53these two topic Java topic and topic one
- 1:02:56but we are dealing with this javat topic
- 1:02:59and there is a consumer offset okay I
- 1:03:01mean it will just assign the offset and
- 1:03:04it will just set up to what your
- 1:03:06consumer read okay now what I want to
- 1:03:09show you I want to show you the data in
- 1:03:12this topic Java topic just run it can
- 1:03:17you see here how many messages we have
- 1:03:19send five messages so the offset count
- 1:03:23being getting increased by one 0 1 2 3 4
- 1:03:27and all the messages what we have been
- 1:03:30sent to this particular topic now from
- 1:03:33this topic consumer is able to listen
- 1:03:36the messages and we can able to see here
- 1:03:40right now if you'll Deep dive one step
- 1:03:43ahead this particular topic we have
- 1:03:45created the partitions how many
- 1:03:47partitions we have created three 0 1 and
- 1:03:51two now if you check the
- 1:03:54messages in the
- 1:03:56topic let me show you the messages
- 1:03:59inside the
- 1:04:02topic I'll go to the
- 1:04:04data all the messages went to partition
- 1:04:08zero but still we have two more
- 1:04:10partition partition one and partition
- 1:04:12two but why all the messages is going to
- 1:04:15partition zero that is not in our hand
- 1:04:19that will be taken care by Juke keeper
- 1:04:22and coordinator okay now we are playing
- 1:04:25with the
- 1:04:26very less messages so all the messages
- 1:04:28is going to a single partition now let's
- 1:04:31say I want to push thousand of row in
- 1:04:34single sh from the producer then maybe
- 1:04:37the Thousand row could be distribute to
- 1:04:40different partition okay this what we
- 1:04:43understand from the theoretically now
- 1:04:45let's prove it what I will do I will
- 1:04:48send a bulk file CSV data through the
- 1:04:52capka and we'll see whether the data is
- 1:04:54being uh pass to the different partition
- 1:04:56or not so for that what I can do I'll
- 1:05:00will cancel the producer okay Let It Be
- 1:05:06My Consumer is keep listening I'm not
- 1:05:08sending anything from the producer that
- 1:05:10is fine you won't find anything but once
- 1:05:12I will start publishing the message will
- 1:05:14keep my My Consumer will keep listening
- 1:05:16to it so what I want to do I want to
- 1:05:19publish a bulk messages so what I have
- 1:05:23done in my download section I have
- 1:05:26downloaded a user. CSV file if I'll open
- 1:05:33this it's loading if you'll observe
- 1:05:37here the total number of row present
- 1:05:39inside the CSV file is th000 okay this
- 1:05:43is triple9 and this is th000 including
- 1:05:46the header it is th1 okay so I want to
- 1:05:50push this particular
- 1:05:52messages from the Capa producer I want
- 1:05:55want to push this user. CSV file how you
- 1:05:58can do that that is very simple so you
- 1:06:01can simply
- 1:06:03Define the path of your file it is there
- 1:06:06in my download folder so I just need a
- 1:06:10path of my download folder so what I can
- 1:06:12do I'll just take the path of
- 1:06:17it
- 1:06:24PWD close it I want to send the enter
- 1:06:30users. CSV file from my producer just
- 1:06:34enter
- 1:06:38it can you see here if I zoom My
- 1:06:41Consumer we can see all the Thousand
- 1:06:44messages in my consumer section it means
- 1:06:47the CSV file the bulk CSV file is being
- 1:06:50loaded to my topic can you see here all
- 1:06:53the messages including the header we can
- 1:06:55able to see the total number of row
- 1:06:58total number of row is 1,1 so let's
- 1:07:01verify it in our upset Explorer now we
- 1:07:05are seeing only the five messages I
- 1:07:07didn't refresh it Let me refresh it can
- 1:07:11you see here there are so many messages
- 1:07:13see these are the CSP data ID first name
- 1:07:17last name email and this is the IP
- 1:07:19address that is what we have there in
- 1:07:21our uh CSV file okay now
- 1:07:26all messages went to single partition or
- 1:07:28different partition let's verify it
- 1:07:31partition
- 1:07:34zero the
- 1:07:37data okay now let's go to the partition
- 1:07:40one it didn't receive any data the count
- 1:07:43is zero partition two also didn't
- 1:07:45receive anything okay so all the
- 1:07:48messages went to the partition zero but
- 1:07:51that is not always the same behavior
- 1:07:54what we can verify I will send it
- 1:07:57again or I have another file called
- 1:08:02customers. CSV okay it is not there let
- 1:08:05me send the same users.
- 1:08:11CSV fine all the messages again being
- 1:08:14send here we how we can
- 1:08:17verify let me go to the
- 1:08:21topic fine now if you'll see in the
- 1:08:23partition zero the to total number of
- 1:08:26record received by partition zero is
- 1:08:301546 so first we published five messages
- 1:08:34then we publish thousand messages then
- 1:08:36again we publish next th000 messages so
- 1:08:39partition zero all total received 1546
- 1:08:43messages now if you open the partition
- 1:08:45one it receives 461 messages okay
- 1:08:49partition two didn't receive anything so
- 1:08:52whoever partition will get the messages
- 1:08:54is not in our our hand this will be
- 1:08:56handled by the Juke keeper or
- 1:08:58coordinator okay now partition zero and
- 1:09:01partition one have the messages so don't
- 1:09:04worry we'll demonstrate this behavior
- 1:09:06from our application once we'll create
- 1:09:09the producer app and consumer app this
- 1:09:11is what we are just trying to understand
- 1:09:13the behavior of partition and
- 1:09:15distributing the messages to multiple
- 1:09:17partition through the common line
- 1:09:19interface okay so I believe we are very
- 1:09:22clear about the producer consumer
- 1:09:24workflow and we have discussed
- 1:09:26everything what we have covered in this
- 1:09:29particular slide we have started the
- 1:09:30Zookeeper we created the capka server we
- 1:09:33create the topic Define the partition
- 1:09:35and also we understand how the load is
- 1:09:37getting distributed to the partition how
- 1:09:39the offset is being assigned how the
- 1:09:41message is being published from producer
- 1:09:43to the topic from the topic to the
- 1:09:46consumer the entire workflow we have
- 1:09:48covered just now okay the next whatever
- 1:09:52we have discussed everything on open
- 1:09:54source apach AP now the next I will walk
- 1:09:57you through a small steps how you can
- 1:10:00play with the conflent capka okay so
- 1:10:03which is the Community Edition we have
- 1:10:05installed or we have downloaded
- 1:10:06yesterday right so if you'll go to the
- 1:10:09desktop if you'll go to the
- 1:10:12downloads software we have this conflent
- 1:10:16right conflent capka now to start this
- 1:10:20confident capka you need to stop the
- 1:10:22open source aaji capka because theuk ER
- 1:10:25server and capka server both are using
- 1:10:28the same port the configuration is same
- 1:10:30right so you need to stop them first
- 1:10:34what I'll do I'll stop everything contrl
- 1:10:36C I stop my producer I stop my consumer
- 1:10:41which I ran in open source API capka
- 1:10:45then I need to stop this capka
- 1:10:48server then I need to stop my jke
- 1:10:54keeper
- 1:10:58where it
- 1:11:04is this
- 1:11:06guy yeah just stop it fine so I'll clear
- 1:11:14everything anyway I need to close this
- 1:11:16particular terminal because we we want
- 1:11:18to play with the confr capka now so what
- 1:11:21you can do go inside this folder just
- 1:11:25open a new terminal atart folder
- 1:11:27everything is same each and every
- 1:11:29command is same only you need to give
- 1:11:32the different path of your config okay
- 1:11:36so what I mean I'll show you we'll
- 1:11:39follow the same first we'll start the
- 1:11:41Juke keeper then we'll start the capka
- 1:11:43server then we'll create the topic by
- 1:11:45defining the partition then we'll see
- 1:11:47the messages so for that what I can
- 1:11:51do I will just go to the folder and I
- 1:11:55will open a
- 1:11:57terminal new terminal let me open couple
- 1:12:00of
- 1:12:04terminal
- 1:12:08okay fine I have created four terminal
- 1:12:12now first let me start the bootstrap
- 1:12:14server sorry zookeeper right so to start
- 1:12:17the Juke keeper I will run the same
- 1:12:20command see here this is what we have
- 1:12:21run in let me Zoom this for
- 1:12:23you the same command we have ran for
- 1:12:27open source Apachi kapka only thing you
- 1:12:30need to Define instead of this config
- 1:12:33you need to give the actual path of your
- 1:12:36Juke keeper. properties file so if you
- 1:12:38check here the Juke keeper. properties
- 1:12:41file is inside
- 1:12:43Etc capka and then the property file
- 1:12:47right so you need to give that path
- 1:12:50instead of config just change it to
- 1:12:54the
- 1:12:55it is inside Etc hyen capka and jer.
- 1:13:00properties just run
- 1:13:03it
- 1:13:05okay so I'm already there inside ban no
- 1:13:09right okay so what I'll do I'll clear it
- 1:13:14let me copy the command okay I've
- 1:13:16already documented the command so let me
- 1:13:19run this particular command I mean see
- 1:13:21here the command is same I'll Zoom this
- 1:13:24for you only we are changing
- 1:13:26the path of your config properties okay
- 1:13:29jer. properties so let's see what is the
- 1:13:33mistake we are
- 1:13:36doing great I've run the same command
- 1:13:39there is no change okay now all good
- 1:13:43Juke keeper server is up and running now
- 1:13:45what is the next step we need to start
- 1:13:47the capka server so to start the capka
- 1:13:50server the command is same here see be
- 1:13:54hyphen b/ capka server start only we're
- 1:13:58changing the path of server.
- 1:14:01properties just copy
- 1:14:03this open a new
- 1:14:07terminal just run
- 1:14:16it cool it started the capka server now
- 1:14:20we need to create a topic so what I will
- 1:14:23do I will copy the same command I mean
- 1:14:26the same command you can use which we
- 1:14:28used in open source Apachi Capa to
- 1:14:30create the topic only the Juke keeper
- 1:14:34and capka server changes to this
- 1:14:36particular config path rest all are same
- 1:14:39because on top of bootstrap or on top of
- 1:14:42your broker you are creating the topic
- 1:14:44so what I'll do let's copy the same
- 1:14:47command anyway it is same right just
- 1:14:49copy
- 1:14:50this go to the terminal any terminal
- 1:14:55just run it we're creating the topic
- 1:14:58with name new topic one okay all good
- 1:15:02we're able to create the topic now if I
- 1:15:05list down all the
- 1:15:07topic just see the command see already I
- 1:15:10explained in open source apka purpose of
- 1:15:13each word right so I'm even I'm using
- 1:15:17the same again for the conflent capka so
- 1:15:21let's run it I just want to describe the
- 1:15:24top iic I mean I just want to list out
- 1:15:28all the topics available in this
- 1:15:329092 can you see here new topic one
- 1:15:36we're able to see this Java topic and
- 1:15:38Java topic one as well because they both
- 1:15:40are using the same 9092 and host I mean
- 1:15:45earlier where it started my open source
- 1:15:48Apachi kapka and this conflent Capa they
- 1:15:51are using the same configuration okay so
- 1:15:53that is the reason you can able to see
- 1:15:55but don't focus on that can you see the
- 1:15:58topic which you just created new topic
- 1:16:00one cool now if you want you can
- 1:16:03describe your new topic again the
- 1:16:06command is same let me copy this and go
- 1:16:10to the terminal ran it okay I did a
- 1:16:16mistake contrl
- 1:16:21C I just want to describe it
- 1:16:25can you see here there is a one
- 1:16:28replication factor and we have defined
- 1:16:29the three partition we can able to see
- 1:16:32that okay all cool now let's start a
- 1:16:35producer message I mean producer
- 1:16:37application and consumer application
- 1:16:39then just publish a message to visualize
- 1:16:42whether we are able to push the messages
- 1:16:44in conflent capka or not okay so I'll
- 1:16:48just copy this
- 1:16:50command to start a producer go to the
- 1:16:54terminal I mean I can use the same right
- 1:16:57I can use
- 1:16:59this control sorry command
- 1:17:04K now meanwhile I need to start a
- 1:17:13consumer I don't want to read it from
- 1:17:15the beginning I want to read it because
- 1:17:18we have not sent anything right now I
- 1:17:21will go to the new
- 1:17:23terminal this one this will my consumer
- 1:17:27now this is my producer cool just run
- 1:17:33it both are pointing to the same topic
- 1:17:36he will publish to the topic one I mean
- 1:17:39this is the producer will publish to the
- 1:17:40topic one consumer will read from the
- 1:17:43topic one I'll send
- 1:17:47something welcome to
- 1:17:52javat or I'll open my post man I'll take
- 1:17:55any payload I'll send in the form of
- 1:17:57Json okay I mean Json
- 1:17:59string you can send any type of data
- 1:18:02okay so let's open
- 1:18:06it I'll take
- 1:18:11something let me take this
- 1:18:15okay I'm sending some Json
- 1:18:21string
- 1:18:23fine some random
- 1:18:26messages can you see here now if you'll
- 1:18:30verify we have sent three messages to
- 1:18:33this particular topic right if you'll go
- 1:18:36and check the topic is not listed here I
- 1:18:39need to refresh
- 1:18:41it refresh the
- 1:18:43topic can you see here we can able to
- 1:18:45see our new topic and we have three
- 1:18:47partition 01
- 1:18:502 0 1 and 2 okay go to this new topic
- 1:18:56one run it can you see here all the
- 1:19:00messages I mean whatever the messages we
- 1:19:03send in the Json string each line is
- 1:19:07published as a individual messages so
- 1:19:09this is not the way we can send the Json
- 1:19:12to just show you that we can able to
- 1:19:14communicate producer and consumer I
- 1:19:17randomly copy paste the Json string and
- 1:19:19I'm sending it okay so now what I'll do
- 1:19:23total number of mess mes is N9 the
- 1:19:26offset count is reached up to eight okay
- 1:19:28and all messages went to the partition
- 1:19:31zero one and two is empty right so what
- 1:19:35I'll do I'll publish I mean we can give
- 1:19:37a try to publish a bulk messages to see
- 1:19:40whether in conflent capka the load is
- 1:19:43being distribut to multiple partition or
- 1:19:45not okay so I'll just copy
- 1:19:53this
- 1:19:58I need a new okay I can use the same
- 1:20:05producer top Pi one and then I need the
- 1:20:10location of my download right because
- 1:20:12that is where my file is present so CD
- 1:20:15come
- 1:20:20out this is
- 1:20:22right close it
- 1:20:29users. CSP just enter
- 1:20:32it can you see here it sent all the rows
- 1:20:38now we'll verify whether thousand row
- 1:20:39went to the topic one in conflent capka
- 1:20:43or
- 1:20:44not let me refresh this
- 1:20:47topic now if you'll run it you can see
- 1:20:51so many messages now did you observe one
- 1:20:54thing some messages went to the
- 1:20:56partition one some messages went to the
- 1:20:59partition
- 1:21:00two some messages went to the partition
- 1:21:02Z now if you see in the partition Z
- 1:21:05total nine messages received to the
- 1:21:07partition Z One received
- 1:21:10925 now rest 75 will be there in the
- 1:21:13partition two all good
- 1:21:1676 okay the messages will or the events
- 1:21:21will be distribute to the number of
- 1:21:23partitions okay okay again we can't
- 1:21:26guarantee who will get the more messages
- 1:21:28or who will get the less messages that
- 1:21:30will be again managed by the coordinator
- 1:21:32okay and this behavior is also same in
- 1:21:34the open source aaji kapka only few of
- 1:21:38the utility you will find in the
- 1:21:39conflent capka okay now in my all the
- 1:21:43upcoming videos I will use the open
- 1:21:45source Apachi kapka but if you are
- 1:21:47familiar now about the confluent capka
- 1:21:50of Community Edition I mean I already
- 1:21:52show you right how we can play with each
- 1:21:55component in open source Apachi capka
- 1:21:58and conflent capka now it is up to you
- 1:22:00when we're doing any program using the
- 1:22:03open source Apachi capka if you want you
- 1:22:05can give a try with this confident capka
- 1:22:07as well
- 1:22:09[Music]
- 1:22:15okay you can install capka on any
- 1:22:18operating system like Windows Mac or
- 1:22:20Linux but the installation process is
- 1:22:23somewhat different different for every
- 1:22:25operating system so instead you will use
- 1:22:28Docker I think it's an even better
- 1:22:30option since you don't have to install
- 1:22:32the tools manually instead you just need
- 1:22:35to write one simple Docker compos file
- 1:22:37which will take care of everything and
- 1:22:40the best part is it will work
- 1:22:42irrespective of any operating system so
- 1:22:44it does not matter what operating system
- 1:22:47you using everything will still work or
- 1:22:50in other word you can say this will act
- 1:22:52as a platform independent okay and I
- 1:22:54would say this will be a quick way to
- 1:22:57get started with capka so without any
- 1:22:59further delay let's move into the
- 1:23:01installation process so in order to
- 1:23:03install capka on the docker container we
- 1:23:06need two things we need an instance of a
- 1:23:09jeeper and an instance of capka
- 1:23:12therefore we are going to create a
- 1:23:14Docker ipen compass. yml file and we are
- 1:23:17going to Define two different service
- 1:23:19one is Juke keeper service and other one
- 1:23:21is the capka service and we'll ask
- 1:23:23Docker to create create these two
- 1:23:25instance for us okay so let's have a
- 1:23:27look how we can Define the docker hpen
- 1:23:29compost. yl file and how we can Define
- 1:23:32the services inside that particular yml
- 1:23:35file okay let's go to the intellig idea
- 1:23:38or you can create in the any editor I'll
- 1:23:42go to the intellig idea I have just
- 1:23:44created one simple folder capka hyen
- 1:23:47installation I'll create a file called
- 1:23:49new click on file and I will give the
- 1:23:52file name as docker hyphen compose do
- 1:23:57yml okay you can keep the extension
- 1:24:00either yml or Y ml it is up to you okay
- 1:24:06I'll just define y ml just click on
- 1:24:11enter now in this Docker compos file you
- 1:24:14can Define what all services you need as
- 1:24:18part of our scenario we need two
- 1:24:20Services zuke keeper and capka okay and
- 1:24:23to start any capka server first we need
- 1:24:26to start the zeper that is what we learn
- 1:24:28before right we have tried multiple
- 1:24:31example using the common line interface
- 1:24:33so we need to first start the Juke
- 1:24:34keeper then we need to start the capka
- 1:24:37correct so I need to define those two
- 1:24:40services so first I will Define the
- 1:24:43version which is
- 1:24:453.1 then I will Define the services you
- 1:24:48can Define n number of services here
- 1:24:51okay so that Docker will create the
- 1:24:53instance for for these Services whatever
- 1:24:55you will Define in this compos file in
- 1:24:58our case we need the services called
- 1:25:00Juke
- 1:25:01keeper okay now how can I get these
- 1:25:04Services simply you need to specify the
- 1:25:07image name of zeper but I don't know
- 1:25:09what is the image name simply go to the
- 1:25:12Chrome browser and you can type
- 1:25:16Ducker
- 1:25:19zeper you can find there will be so many
- 1:25:22Ducker images from the Ducker Hub but
- 1:25:25I'm using the stable one this one
- 1:25:30okay if you'll Zoom this this is the
- 1:25:34image name copy this name just paste it
- 1:25:39Define the image I want this image as a
- 1:25:42zookeeper so that Docker will start this
- 1:25:46zookeeper instances by pulling the image
- 1:25:48from this okay now once you have defined
- 1:25:51the images you can Define the name of
- 1:25:54the container so I'll give the same name
- 1:25:56as a Juke keeper okay now once you have
- 1:26:00defined the image name and container
- 1:26:02name next you can Define the ports where
- 1:26:05this zeper will run hostport and your
- 1:26:09machine Port so I Define 2181 this is
- 1:26:13the default one so I'll specify the
- 1:26:15same fine we have defined one Services
- 1:26:19similarly we need to Define another
- 1:26:20services that is capka so let me copy
- 1:26:23this
- 1:26:24and I will simply paste
- 1:26:28it second service name is kapka and
- 1:26:32images name you can search the
- 1:26:36image just simply change the name of
- 1:26:39your service what you want to use this
- 1:26:41particular tag having this zuke keeper
- 1:26:44and capka so if you want I can show you
- 1:26:49that in this repository if you will
- 1:26:52filter see here we have the capka okay
- 1:26:55I'm using the same now just change the
- 1:26:58container name as a capka and the
- 1:27:01default Port of the capka
- 1:27:049092
- 1:27:079092 once you have defined the service
- 1:27:09zeper and capka then you just need to
- 1:27:12define the environment by telling to the
- 1:27:15capka where is your Juke keeper up and
- 1:27:17running okay so simply just Define the
- 1:27:22environments or environment
- 1:27:24attribute just Define here capka
- 1:27:28advertised host name and the host is
- 1:27:31Local
- 1:27:33Host okay now we just need to Define
- 1:27:36capka Juke keeper connect where is your
- 1:27:39Juke keeper up and running you just need
- 1:27:41to Define that okay so let me quickly
- 1:27:43Define
- 1:27:49it now just
- 1:27:52Define Juke keeper
- 1:27:55it is on Port
- 1:27:572181 okay all looks
- 1:27:59good next once we have created this
- 1:28:03Docker compose file we can execute it
- 1:28:06using the docker compose command Okay
- 1:28:08Docker compose up but before that I just
- 1:28:11want to show you the feature in ID if
- 1:28:14you want to run The Individual Services
- 1:28:16go to that particular service name and
- 1:28:19here you can see the option called
- 1:28:20Docker compos of zeper okay similarly if
- 1:28:23you want to start only capka just go to
- 1:28:26this service I mean this icon here
- 1:28:28you'll find the option to Docker compose
- 1:28:30of capka if you want to run all the
- 1:28:33services what you have defined in your
- 1:28:35file just run on the I mean just click
- 1:28:38on this particular icon okay it will
- 1:28:40execute all the services defined inside
- 1:28:43this compos file okay but I just want to
- 1:28:46run it explicitly in the command prop so
- 1:28:50that you will be aware about the command
- 1:28:52so what I'll do I'll go to this
- 1:28:54particular directory
- 1:28:56first CD capka
- 1:29:00installation
- 1:29:01okay now if I check the ls I have this
- 1:29:04Docker compos yl file so what I want to
- 1:29:08do here I'll simply run the command
- 1:29:12Docker compose then give the file name
- 1:29:16which is Docker compose yml and up that
- 1:29:20particular compos file in the background
- 1:29:22so that is the reason I have defined
- 1:29:24Das D okay so the main command you just
- 1:29:28need to run Docker compose give the file
- 1:29:30name and define up okay just run
- 1:29:34it can you see here it's pulling the
- 1:29:38capka pulling the Juke keeper it will
- 1:29:41take few second to download those two
- 1:29:53images
- 1:29:55okay if You observe here container Juke
- 1:29:58keeper started and container capka
- 1:30:01started also Network UPA installation
- 1:30:03default created it created the default
- 1:30:05Network and it started these two
- 1:30:08container so we didn't ran this Juke
- 1:30:11keeper and capka from our local binary
- 1:30:13distribution rather we asked to the
- 1:30:15docker to create these two instance for
- 1:30:17us okay now to verify whether these two
- 1:30:21container is running or not what I can
- 1:30:24do first let me check the
- 1:30:28images can you see here it created the
- 1:30:31image called capka and you'll also find
- 1:30:33the image called Juke keeper okay now
- 1:30:37we'll verify whether these two instance
- 1:30:38is running or not so simply you can run
- 1:30:41Docker
- 1:30:43PS can you see here let me Zoom
- 1:30:51this so container ID is this and the
- 1:30:54name of the container we have
- 1:30:56defined Juke keeper and capka okay these
- 1:31:01two instances up and running now what
- 1:31:04we'll do let's move ahead one step into
- 1:31:08the Ducker container or capka container
- 1:31:11and we'll execute the command to see
- 1:31:13whether the capka is up or not okay so
- 1:31:17for that what I can do you need to fire
- 1:31:19the
- 1:31:21command let me clear it
- 1:31:26Docker
- 1:31:28execute integrated terminal and name of
- 1:31:32your container ID in our case the
- 1:31:35container ID is capka right that is what
- 1:31:38we have defined our container name next
- 1:31:40just Define slash
- 1:31:43bin
- 1:31:45slsh before that just add a dash here
- 1:31:50okay now if you'll enter it now if
- 1:31:53you'll do the ls you can find these many
- 1:31:56folder okay just go inside CD
- 1:32:00opt now if you'll do the ls you can find
- 1:32:04your capka binary distribution installed
- 1:32:07by the docker itself okay now go inside
- 1:32:10this particular folder I'll copy
- 1:32:15this now we are inside the if you'll
- 1:32:18check the directory we inside SL opt
- 1:32:22capka and the binary version okay now if
- 1:32:26you'll do the ls you'll find license
- 1:32:29notice B everything now if you go inside
- 1:32:32the
- 1:32:35bin you'll find all the script file sell
- 1:32:38script file or batch file for Windows
- 1:32:41can you see here all the Sal script file
- 1:32:45to start the cap card to create the
- 1:32:46topic to produce the messages to consume
- 1:32:50the messages everything you will find
- 1:32:52now what we'll do we'll quickly end of
- 1:32:55this we'll just create a topic and we'll
- 1:32:58publish a message and we'll just consume
- 1:33:00it or we'll just publish a message and
- 1:33:02we'll verify that in upset capka
- 1:33:04Explorer okay so to create a topic you
- 1:33:07know the command so let me enter
- 1:33:11it simply just run this capka
- 1:33:15topic.sh we want to create the topic and
- 1:33:18give the host and the service name of
- 1:33:20Juke keeper replication factor and we
- 1:33:22need partition one topic name is Quick
- 1:33:25Start okay let me Zoom
- 1:33:27this quick start now just enter
- 1:33:33it it created the topic called quick
- 1:33:36start so we'll verify whether this topic
- 1:33:38is created or not so what I can do I'll
- 1:33:41open the upset
- 1:33:46Explorer then just click on this I mean
- 1:33:50you just need to give the see here if
- 1:33:52you'll see the connection property I'm
- 1:33:54giving where is my juke keeper up and
- 1:33:55running now if you open the topic we can
- 1:33:58see the topic got created name quick
- 1:34:01start so there will be no data since we
- 1:34:03didn't publish anything to this topic so
- 1:34:06to publish to the topic what I can do
- 1:34:09I'll go to the
- 1:34:13terminal then I'll simply run the
- 1:34:16producer command capka consult
- 1:34:19producer. to which topic you want to
- 1:34:21produce the messages and what is the
- 1:34:25host and put where your bootstrap server
- 1:34:27or capka server is up and running right
- 1:34:30this is the basic command we already
- 1:34:31tried this many times in our capka
- 1:34:33Series so I'll just enter
- 1:34:36it then I'll send few message hi
- 1:34:41hello Java
- 1:34:45TIY then install kupka using Docker and
- 1:34:52Docker comp
- 1:34:56okay or anything some
- 1:34:59number some random
- 1:35:01messages now how can we verify whether
- 1:35:04the messages what I produce is there in
- 1:35:08the capka server or not or it is there
- 1:35:10in the topic or not simply go to the
- 1:35:13upset
- 1:35:16Explorer now just go to the data I want
- 1:35:21to see all the latest data just execute
- 1:35:24this can you see here all the messages
- 1:35:27what we have just produces is there in
- 1:35:29the topic now what I want to verify that
- 1:35:32is fine producer is able to publish the
- 1:35:35messages to the topic now I want to
- 1:35:38verify whether consumer is able to
- 1:35:41consume this messages or not okay so
- 1:35:44that is again Quick Step what I'll do
- 1:35:46I'll just open a new
- 1:35:48terminal then I will just um it's fine I
- 1:35:52can run from this command also so I'll
- 1:35:55just execute the command
- 1:35:57this then I'll go to the
- 1:36:00directory CD
- 1:36:03opt
- 1:36:04LS I want to go to this particular
- 1:36:07directory
- 1:36:11CD go inside the bin LS that's it right
- 1:36:16now we just need to start a consumer app
- 1:36:19so that it will read whatever the
- 1:36:21messages we have published to this
- 1:36:22particular topic
- 1:36:24so simply go to the
- 1:36:26terminal what I'll do I'll open the
- 1:36:29background here
- 1:36:32okay now let me clear
- 1:36:34everything I'll just run this command
- 1:36:37capka console consuma Dosh read the
- 1:36:41messages from this topic from the
- 1:36:43beginning read it and where is my capka
- 1:36:45server okay just enter
- 1:36:49it can you see here it take few second
- 1:36:52and we able to to consume all the
- 1:36:55messages okay so in real time you can
- 1:36:58follow this approach where you can
- 1:37:00install the Capa Juke keeper MySQL
- 1:37:03redish whatever the services you don't
- 1:37:06want to install on your machine and you
- 1:37:08want to play with them or you are hosted
- 1:37:11to the server where you don't want to
- 1:37:13install them physically but rather you
- 1:37:15are you want to get it from the docker
- 1:37:17so you can go with this particular
- 1:37:19approach all the services you can
- 1:37:21install through this particular Docker
- 1:37:23compose approach okay and this is really
- 1:37:26needed if your application hosted to the
- 1:37:29cloud infrastructure and if you are not
- 1:37:31using the managed service given by the
- 1:37:33cloud provider then you must need to
- 1:37:35play with this container platform using
- 1:37:37this Docker compose so do let me know in
- 1:37:40a comment section if you guys have any
- 1:37:42[Music]
- 1:37:49doubts so in this tutorial we'll
- 1:37:52understand how to produ produce or
- 1:37:54publish a message to the capka topic
- 1:37:55using spring Boot and also we'll Deep
- 1:37:58dive further to understand how the
- 1:38:00messages are getting distributed to
- 1:38:02multiple partitions okay all right so
- 1:38:05let's quickly create a new project to
- 1:38:07demonstrate capka using springboard
- 1:38:09first let me create a new project from
- 1:38:11the intellig
- 1:38:13itself
- 1:38:16project Define the group ID artifact ID
- 1:38:19version all the
- 1:38:22things
- 1:38:25change it to
- 1:38:26Maven let it be jdk 17 is required since
- 1:38:29we using the spring boot
- 1:38:323.0 then specify this project
- 1:38:38name next here you need to choose the
- 1:38:40correct dependency so I'm going to use
- 1:38:43the web dependency because I want to
- 1:38:45expose on endpoint if I'll trigger that
- 1:38:47endpoint it will publish message to the
- 1:38:48capka okay that is what I want to design
- 1:38:51I need web and the main dependency you
- 1:38:54need to add is capka dependency now if
- 1:38:57You observe here there are two
- 1:38:58dependency you are getting for capka one
- 1:39:00is for capka stream one one is for
- 1:39:03spring for aaji kapka okay so we are
- 1:39:05going to use this since we are not going
- 1:39:07to use the capka stream for now we'll go
- 1:39:09with the spring capka dependency okay I
- 1:39:13believe that's enough now click on
- 1:39:18next so next Once the project is
- 1:39:21imported successfully since we are
- 1:39:23playing with the capka first we need to
- 1:39:25start the capka server and already I
- 1:39:28explained the order to start the capka
- 1:39:30first you need to start the Juke keeper
- 1:39:32then you need to start the capka server
- 1:39:34then as for your need whether you want
- 1:39:35to create a topic and you want to
- 1:39:37configure partition count replication
- 1:39:39Factor all these steps you can do but to
- 1:39:41start the capka server first you need to
- 1:39:43start the Juke keeper okay so what I'll
- 1:39:46do I'll go to the directory where I
- 1:39:48install my capka this is what the
- 1:39:50directory then I'll open the terminal
- 1:39:53okay so I'll start the Juke keeper first
- 1:39:56so I already shared the GitHub link
- 1:39:58where I have defined all the command for
- 1:40:00the capka so the first command to start
- 1:40:02the Juke keeper I will copy this then
- 1:40:05I'll go to the
- 1:40:06terminal fine then I will just start the
- 1:40:09Juke keeper
- 1:40:12okay once I start the Juke keeper then
- 1:40:15immediately I'll start my capka server
- 1:40:18so for that I will go to the terminal
- 1:40:20and I will take the command to start the
- 1:40:23C
- 1:40:25server then I'll open a new
- 1:40:30terminal okay I can start the capka
- 1:40:33server
- 1:40:36here fine so if You observe my capka
- 1:40:40server is started on Port 9092 that is
- 1:40:44the default Port okay I just highlight
- 1:40:47here if you can see this let me Zoom
- 1:40:49this for you 9092 is the default Port of
- 1:40:52my kapka server okay and this is my juke
- 1:40:56keeper fine now all good all our kapka
- 1:41:01ecosystem is ready to start the
- 1:41:02development go to your intellig idea and
- 1:41:05first I will go to application.
- 1:41:07properties file or you can create the
- 1:41:09application. yml file to tell to this
- 1:41:12application where is your capka server
- 1:41:14is up and running if you'll not specify
- 1:41:16that how your application is know that
- 1:41:19where is your capka is running so that
- 1:41:20you will publish the message right so
- 1:41:22that's Theon
- 1:41:23you need to specify in your application
- 1:41:26that where is your bootstrap server of
- 1:41:28Capa producer is running okay so I'll
- 1:41:30create a file I'll name it application.
- 1:41:35yml then I can Define here first I will
- 1:41:38Define my server Port let's say I just
- 1:41:40want to define something 91
- 1:41:4291 then I will Define spring capka
- 1:41:47producer then bootstrap server okay and
- 1:41:51give the URL you can Define multiple
- 1:41:55server port and host but since we have
- 1:41:57only the default one which is our local
- 1:41:59host I just specify Here Local Host and
- 1:42:03the default Port is
- 1:42:079092 fine next let me create couple of
- 1:42:11package where we can Define our logic to
- 1:42:14publish the message to the capka
- 1:42:20okay next let me create the service
- 1:42:23class where actually we will write the
- 1:42:25code to publish the messages so I will
- 1:42:28create a class called let's say I'll
- 1:42:29will give a name something like capka
- 1:42:31message publisher something like that
- 1:42:36okay now this is the class so in this
- 1:42:39class to publish a message to the capka
- 1:42:41or from our application if you want to
- 1:42:44talk to the capka server then you need
- 1:42:46to use some class given by the spring
- 1:42:48framework that is capka template okay
- 1:42:54capka
- 1:42:56template Define the key and value for
- 1:42:58now let me Define the key as a string
- 1:43:00and value as a
- 1:43:04object let me annotate this class as a
- 1:43:06atate
- 1:43:07service and I need to Auto add
- 1:43:11this fine now here I'll write a simple
- 1:43:15method public will simply return the
- 1:43:18message nothing else I mean we'll simply
- 1:43:20return the message to the topic so so
- 1:43:23I'll name it send message to
- 1:43:28topic what message you want to send I'll
- 1:43:31pass it from the
- 1:43:34controller since we have injected the
- 1:43:37template now what we can do using this
- 1:43:39template we can talk to the capka server
- 1:43:41so I'll just use this instance template
- 1:43:44do send method If You observe send is
- 1:43:47the overloaded method it contains
- 1:43:49multiple argument you can use as for
- 1:43:51your need since I just want to send a
- 1:43:53simple message to the topic I'll use
- 1:43:56this second method overloaded method
- 1:43:58I'll specify the topic name let's say JT
- 1:44:01or Java
- 1:44:04Tey fine and then I need to specify what
- 1:44:08data I want to send to the topic this
- 1:44:11one now one thing you need to observe
- 1:44:14here we have used this template. send
- 1:44:16method but if you'll take a look into
- 1:44:19its return return type it is completable
- 1:44:22Future
- 1:44:23okay so this is the completable future
- 1:44:26so here if you want to block the sending
- 1:44:28thread and get the result about the sent
- 1:44:30messages then what you can do simply you
- 1:44:33can do future do get API call from the
- 1:44:38completable future then the thread will
- 1:44:40wait for the results to come but it will
- 1:44:42definitely slow down the producer as you
- 1:44:45know capka is a fast in processing
- 1:44:47platform therefore it's better to handle
- 1:44:50the results asynchronously so that the
- 1:44:52subsequent message do not wait for the
- 1:44:54result of the previous messages okay now
- 1:44:57how we can do that we can do this using
- 1:44:59a call back implementation I'll show you
- 1:45:02that how we can do that let's remove
- 1:45:04this piece of code we'll simply tell to
- 1:45:06this future object
- 1:45:09future when complete when the future
- 1:45:12will be complete okay then I will use
- 1:45:15the Lambda it will give me the result
- 1:45:17and also it will give me the parameter
- 1:45:19of exception okay now I can simply print
- 1:45:23whatever the things I
- 1:45:25need so just add a comma here and I'll
- 1:45:28just add some print statement okay so
- 1:45:31here what I'm doing I'm just checking if
- 1:45:34xal equal to null if there is no
- 1:45:36exception it means message is sent and I
- 1:45:38just want to tag what messages we have
- 1:45:41sent and what is the upset number of it
- 1:45:44also if you want you can also get the
- 1:45:46partition number where the message went
- 1:45:48okay so if you see here there is a
- 1:45:51method result Dot get record metadata do
- 1:45:56get not record
- 1:45:58metadata uh okay
- 1:46:02dot you can see here the partition okay
- 1:46:05you also if you want to print the
- 1:46:06partition you can do that but let's make
- 1:46:08it simple okay I am trying to make it
- 1:46:10asynchronous so that it don't wait for
- 1:46:13the result to come and I am just simply
- 1:46:15printing the message this is for success
- 1:46:17and this is for failure okay that's it
- 1:46:21now this is the me method who will send
- 1:46:23message to the topic now I will call
- 1:46:26this method from my controller class so
- 1:46:28I'll create a controller class
- 1:46:31quickly I'll name it
- 1:46:35event now I need to un your atate R
- 1:46:39controller I'll also Define request
- 1:46:42mapping and I'll give some URL
- 1:46:46here then I need to inject the capka
- 1:46:49message publisher the service class okay
- 1:46:52so I'll just inject private let me Zoom
- 1:46:59this so once you inject the publisher
- 1:47:02now I can simply Define the endpoint who
- 1:47:04will publish the message okay so I'll
- 1:47:07write public response
- 1:47:09entity up type
- 1:47:12generic let's say publish
- 1:47:16message okay I'll pass the message in
- 1:47:19the request
- 1:47:21URL and here at path
- 1:47:25variable also I will Define this as a
- 1:47:28gate
- 1:47:31mapping Define the URL I'll just Define
- 1:47:34publish and what message I want to
- 1:47:38publish simple right now what I can do
- 1:47:43I'll simply call
- 1:47:45Publisher dot send message to the topic
- 1:47:49and I will simply give the message
- 1:47:51name that's it
- 1:47:53but for safer side I'll keep it inside
- 1:47:56the try
- 1:48:00catch fine now here I'll just return
- 1:48:04response
- 1:48:06entity do okay and some
- 1:48:10messages I'll give something like
- 1:48:17message similarly in the catch also I'll
- 1:48:20just Define return respon
- 1:48:23entity dot I can specify the
- 1:48:27status which will be HTTP
- 1:48:31status do internal server error anything
- 1:48:34whatever you want to specify okay then
- 1:48:37I'll simply build
- 1:48:39it fine I have just handled the TR cage
- 1:48:42here to just track the response or to
- 1:48:46just track the uh event whether it is
- 1:48:48successfully published or not fine let's
- 1:48:51say my capka server is down and I'm
- 1:48:53sending the messages then definitely I
- 1:48:54will get the internal server error you
- 1:48:57can handle it in better way but this is
- 1:48:59the simple exception handling I have did
- 1:49:01here okay so now we are all set to start
- 1:49:05the server okay I mean application now
- 1:49:09if You observe one thing we have not
- 1:49:11create topic manually okay if You
- 1:49:14observe here while sending the messages
- 1:49:18to using this Capa template I'm just
- 1:49:20giving the topic name but I have not
- 1:49:23created it manually through the command
- 1:49:25prompt so whether the spring boot will
- 1:49:29create the topic for us or not that we
- 1:49:31are going to verify if spring boot will
- 1:49:33create the topic on behalf of us then
- 1:49:36what is the default configuration he
- 1:49:37will add that we are going to verify now
- 1:49:40okay so all good let me go to the main
- 1:49:45class just start the
- 1:49:50application so you can observe here it
- 1:49:52started on Port 9191 now let's go to the
- 1:49:56postman and try to hit this particular
- 1:49:58end point / producer hyen app SLP
- 1:50:01publish on the messages okay now let me
- 1:50:04clear
- 1:50:07this let's say I'll just pass a message
- 1:50:13welcome can you see here message
- 1:50:16published successfully and now if you go
- 1:50:19and check in your console you can find
- 1:50:21that is what we have just added in our
- 1:50:23completable future when complete this
- 1:50:25messages we are going to print the
- 1:50:27message we have sent is welcome let me
- 1:50:29Zoom this for you and the upset count is
- 1:50:32zero offset is nothing the position of
- 1:50:34your message inside the partition okay
- 1:50:37so that already I explain in the capka
- 1:50:39component and its architecture so I hope
- 1:50:42that makes sense to understand this uh
- 1:50:44term okay now I'll add some other
- 1:50:47messages let's say
- 1:50:50ABC
- 1:50:51DF
- 1:50:55Ram
- 1:50:57bimol okay so I have published near to
- 1:50:59four or five messages now let's see go
- 1:51:02and check in your console see all the
- 1:51:04messages is assigning to the upset but
- 1:51:07we are not sure whether it is going to
- 1:51:09which partition because even we don't
- 1:51:11know the topic configuration created by
- 1:51:13Spring boot so let's verify the topic
- 1:51:16configuration then I will show you how
- 1:51:18you can create manually and how you can
- 1:51:20create programmatically okay so now to
- 1:51:24describe the topic what I can do I'll
- 1:51:26will go to the uh
- 1:51:29terminal then simply fire this command
- 1:51:32give your topic name which is Java Tei
- 1:51:36hyphen demo one so already all the
- 1:51:40required command I have shared in this G
- 1:51:42repo okay so I will also share this link
- 1:51:45in video description so you can use that
- 1:51:48now let me trigger
- 1:51:51this so so if You observe here the
- 1:51:54partition count equal to 1 and
- 1:51:56replication Factor equal to 1 so
- 1:51:58whenever we are allowing spring boot to
- 1:52:01create the topic on behalf of us it will
- 1:52:03go with the default configuration with
- 1:52:05the partition count one it won't give
- 1:52:08you the correct or proper throw put if
- 1:52:10your application is being used by
- 1:52:12multiple consumer they cannot able to
- 1:52:15handle your request concurrently because
- 1:52:17of this partition count now how you can
- 1:52:19create topic with your own configuration
- 1:52:22let's say I want to go with Partition
- 1:52:23count 5 or 10 15 as for my needed then
- 1:52:27you need to create it manually how we
- 1:52:29can create that you can go to this and I
- 1:52:32have already shared this command copy
- 1:52:35this and just go to the terminal
- 1:52:40again just trigger this command and you
- 1:52:43just need to change the topic
- 1:52:47name so I will change it
- 1:52:49to Java Tey
- 1:52:54okay let me create
- 1:52:59this it created the topic for me but
- 1:53:03before we use this javat demo to topic
- 1:53:05first let's verify javat demo and topic
- 1:53:08whatever the message we have sent
- 1:53:10whether it is there in the topic or not
- 1:53:12if it is there what is the partition and
- 1:53:14what is the upset we can visualize okay
- 1:53:17so to visualize the capka ecosystem I
- 1:53:19already explained you can use the offset
- 1:53:22EXP Explorer okay just open this then
- 1:53:26simply take the
- 1:53:28connection now if you'll observe here
- 1:53:31there are two topic right javate demo
- 1:53:34one javate demo hyphen 2 okay so the
- 1:53:39first topic we use this demo one so if
- 1:53:41I'll open this topic if I'll click on
- 1:53:44the data and if I'll change it to the
- 1:53:47new and if I'll send the request we can
- 1:53:50see all the five messages we have
- 1:53:51published right welcome ABC d a ram Bal
- 1:53:56okay and corresponding upset and all
- 1:53:58went to the partition zero because we
- 1:54:00have only one partition now since we
- 1:54:02have created the other topic with the
- 1:54:05three partition 0 and one two you can
- 1:54:07see here now let's see how we can send
- 1:54:11multiple or bulk messages to the topic
- 1:54:13so that message will be distribute to
- 1:54:16multiple partition so what we'll do
- 1:54:19we'll use this topic in our code now
- 1:54:21rather than using the default topy
- 1:54:22created by springboard to verify the
- 1:54:25messages bulk messages is getting
- 1:54:27distributed to multiple partition that
- 1:54:30is have the increasing the throw put
- 1:54:32from the middleman of kfka okay so what
- 1:54:35I'll do I'll just go to the code first
- 1:54:38I'll change the topic
- 1:54:40name Java demo hyphen 2 this topic have
- 1:54:45three partition okay so you f send bulk
- 1:54:48messages let me write some logic to send
- 1:54:50some bulk messages so go to the
- 1:54:53controller what I'll do I'll just Define
- 1:54:56a for Loop
- 1:54:58here int I = to z i less than equal to
- 1:55:04let's say I just want to send 10,000
- 1:55:06messages okay that is how I can show you
- 1:55:09how I can send the bulk message
- 1:55:14i++ so here also I will just do some
- 1:55:16changes with the message I'll just upend
- 1:55:19the
- 1:55:21number
- 1:55:23whatever the message I will send let's
- 1:55:24say I send user then user one 2 like
- 1:55:2710,000 user string message will publish
- 1:55:30to the Capa topic then we'll verify
- 1:55:32since we have the three partition how
- 1:55:34many messages is went to each partition
- 1:55:36and how the messages is getting
- 1:55:38distributed okay this is the simple way
- 1:55:41to demonstrate this particular concept
- 1:55:44all good let me restart
- 1:55:48it so you can see here application
- 1:55:50started on Port 9 91 go to the postman
- 1:55:54then I'll simply trigger here let's say
- 1:55:56user so that it will append user 1 2 3
- 1:55:59like that okay now I will send the
- 1:56:01request which will trigger 10,000
- 1:56:04message to the capka topic send the
- 1:56:07request it took 1 1915 millisecond okay
- 1:56:12now go to
- 1:56:13the this particular upside Explorer now
- 1:56:16let me refresh
- 1:56:18this okay I'll go to this particular
- 1:56:21topic
- 1:56:22and if I'll just verify the data there
- 1:56:25are so many messages okay now if I'll
- 1:56:29verify each message count from each
- 1:56:32partition let's say I went to partition
- 1:56:34zero it received 1682 messages if you
- 1:56:38want to see the data of those particular
- 1:56:41thing you can hit this all the messages
- 1:56:43you will get from the partition zero and
- 1:56:47you can find this upset here right so
- 1:56:49the count we want to see it receive one
- 1:56:521682 out of 10,000 now we verify
- 1:56:55partition 1 it receed 5120
- 1:56:595,120 now we'll see the partition two it
- 1:57:02received
- 1:57:043,19 all the three partition received
- 1:57:07the message from the producer what we
- 1:57:10just published now okay so it clearly
- 1:57:13says that whenever you have multiple
- 1:57:15partitions capka scheduler or Juke
- 1:57:18keeper will take care to handle the
- 1:57:20message load and splited into to the or
- 1:57:22distributed into the multiple partition
- 1:57:25based on their availability okay so you
- 1:57:28can see here
- 1:57:29again partition 0 1 and two they both
- 1:57:33have some information or some
- 1:57:35messages so this is how capka distribute
- 1:57:39the load to multiple partition now the
- 1:57:42next thing in a real time do we need to
- 1:57:45create the topic from this command the
- 1:57:47way I have created here javat demo 2 it
- 1:57:51is up to you
- 1:57:52in some case in some industry they have
- 1:57:55their own portal to create the topic
- 1:57:57from the automated way but if you don't
- 1:58:00want to do or if you don't want to
- 1:58:02create it from the command promt then
- 1:58:04there is a way given by the spring
- 1:58:06framework you can programmatically
- 1:58:08create the topic now I will show you how
- 1:58:11you can create the topic
- 1:58:12programmatically okay just go to your
- 1:58:15code so what I'll do I'll create another
- 1:58:19package config now I'll create a
- 1:58:27class then simply uned this atate
- 1:58:31configuration
- 1:58:32fine then next there is a class given by
- 1:58:36Spring framework that is new topic using
- 1:58:38that you can create a new topic so let
- 1:58:41me do that public new topic make sure to
- 1:58:45import it from ca. clients.
- 1:58:50admin now you can just return new new
- 1:58:55topic this is also overloaded method and
- 1:58:58you can see here new topic you can
- 1:59:00specify the topic name you want to
- 1:59:02create you can specify the number of
- 1:59:04partition you want to use and
- 1:59:06replication Factor see all the required
- 1:59:09information you can specify here let me
- 1:59:11specify the name javat Tei
- 1:59:16demo three okay and number of partition
- 1:59:20let's say I just want to Define five
- 1:59:23and replication Factor let's say I just
- 1:59:26want to Define one so it will be in the
- 1:59:31S specify the count
- 1:59:35one okay let me Zoom
- 1:59:38this now just Define this at theate
- 1:59:42B fine so we are not creating the topic
- 1:59:46using the command or we are not allowing
- 1:59:48our spring boot to create the default
- 1:59:51topic rather
- 1:59:52we tell to the spring whatever the topic
- 1:59:54configuration we need programmatically
- 1:59:57we are telling him this is the topic
- 1:59:59name I want to create and this is the
- 2:00:00number of partition I want to keep as
- 2:00:02part of my topic and number of
- 2:00:04replication Factor okay now let me use
- 2:00:06the same topic in my
- 2:00:12services fine now let me rerun our
- 2:00:17app so again it started now go to the
- 2:00:20postman let me clear this
- 2:00:23now we have five topic okay so we'll
- 2:00:25verify first whether the topic is
- 2:00:27created or not can you see here topic is
- 2:00:31created with five partition see the
- 2:00:34partition count 0 1 2 3 4 we don't have
- 2:00:37any data because we have not published
- 2:00:39any message to this particular topic now
- 2:00:41same 10,000 messages again I want to
- 2:00:43publish to the topic Java demo three and
- 2:00:46we have five partition now now let's see
- 2:00:49how much record how many record is going
- 2:00:52to each partition there there could be
- 2:00:54possible that some partition will not
- 2:00:56get any single message also okay because
- 2:00:59that is not in our hand that will be
- 2:01:01handled by the capka schuer okay or Juke
- 2:01:04keeper now I'll just send the
- 2:01:08request all the message is published
- 2:01:11let's verify
- 2:01:13it just refresh
- 2:01:17it go to the data trigger it
- 2:01:22now let's see whether all the partition
- 2:01:25is getting the messages or not partition
- 2:01:28zero received
- 2:01:303,717 partition one received
- 2:01:331,954 partition 2 received
- 2:01:37496 partition three didn't received
- 2:01:39anything partition four received
- 2:01:41something
- 2:01:423834 out of five partition parti 3
- 2:01:47didn't received any single message okay
- 2:01:49so that is the reason it's not enough
- 2:01:51our hand it will be decided by Juke
- 2:01:53keeper to whom he need to send the
- 2:01:55messages since we are only sending
- 2:01:5710,000 messages Juke keeper assume that
- 2:02:00okay these four partition is capable
- 2:02:02enough to handle those kind of messages
- 2:02:05to increase the throw putut okay so that
- 2:02:07is how they will decide so this is how
- 2:02:09you can develop your producer
- 2:02:11application and you can also customize
- 2:02:13your topic as for your required
- 2:02:15configuration to increase the
- 2:02:17application performance okay so in my
- 2:02:20upcoming tutorial I will explain about
- 2:02:22the consumer part so we understand
- 2:02:24producer and we send message to the
- 2:02:26topic and topic is being splitted to the
- 2:02:29multiple partition so all the message
- 2:02:31went to multiple partition okay now in
- 2:02:33my next tutorial I will show you how you
- 2:02:36can develop the consumer and how
- 2:02:38concurrently you can handle the events
- 2:02:40coming from the cap cut topic to this
- 2:02:43consumer
- 2:02:50okay so who when your producer
- 2:02:52application will produce the messages it
- 2:02:55needs to consume by the consumer
- 2:02:57otherwise by keeping the messages in a
- 2:02:59topic without using it does not make any
- 2:03:02sense so you need to write a consumer
- 2:03:04app who will consume the messages from
- 2:03:06your topic now if You observe in the
- 2:03:09flow one consumer is listening to three
- 2:03:12different partition which is not good I
- 2:03:14mean it don't give you the better
- 2:03:16throughput then how can we design our
- 2:03:18consumer app to overcome this issue it's
- 2:03:21very simple just share the workload to
- 2:03:24multiple consumer instance so what you
- 2:03:26can do we can simply create multiple
- 2:03:29consumer instance and then we can group
- 2:03:31them with an unique ID so that this ID
- 2:03:34will indicate the task of a specific
- 2:03:36consumer group so now if You observe we
- 2:03:40have three consumer instance right and
- 2:03:42how many partition do we have three so
- 2:03:45each partition can consumed by
- 2:03:48individual consumer instance which will
- 2:03:50definitely increase the through put
- 2:03:52earlier a single consumer instance is
- 2:03:54pointing to multiple partition now we
- 2:03:57have a individual consumer instance for
- 2:03:59each partition of our topic definitely
- 2:04:02this will give you the better
- 2:04:04performance but what if I have another
- 2:04:06consumer instance then what you'll do
- 2:04:09because already three partition is being
- 2:04:10assigned to three consumer instance then
- 2:04:13what the fourth consumer instance will
- 2:04:15do nothing he will simply stay on bench
- 2:04:18the way it company hires the employee
- 2:04:21and keeping the them in a bench for
- 2:04:22future project and once the project came
- 2:04:25they assign the employee right similarly
- 2:04:28if any consumer instance dies then this
- 2:04:31consumer 4 will replace that and start
- 2:04:34consuming the messages and this concept
- 2:04:37is called consumer rebalancing okay so
- 2:04:40this is what all about the theory of
- 2:04:43capka consumer on the basic level so we
- 2:04:47understand the theory now let's start
- 2:04:49implementing our consumer app to
- 2:04:51understand this PR ially okay so let's
- 2:04:53go to the intelligent idea so this is
- 2:04:56the cup producer example we have tried
- 2:04:58last time now let's quickly create a
- 2:05:00Capa consumer application click on file
- 2:05:03new click on
- 2:05:05Project click
- 2:05:07next then specify all the required
- 2:05:14field now let's add all the required
- 2:05:16dependency I'll just add
- 2:05:19capka spring for Apachi kapka
- 2:05:22and then I will just add web
- 2:05:24dependency spring web that's it click on
- 2:05:28next click on
- 2:05:30finish so now we have created the
- 2:05:33consumer application now if You observe
- 2:05:36here consumer application is the one who
- 2:05:39connect to the capka topic to receive
- 2:05:41the messages so you need to tell to your
- 2:05:43consumer application where is your capka
- 2:05:45server up and running then only he can
- 2:05:47connect to the topic and he can retrive
- 2:05:49the messages right so first
- 2:05:52you need to tell to your consumer where
- 2:05:54is your capka server is up and running
- 2:05:56so go to the resource or application or
- 2:05:58properties it is up to you I'm using the
- 2:06:00application. IML
- 2:06:04file now here you will Define spring
- 2:06:08capka then consumer and where is your
- 2:06:11bootstrap server is up and running in
- 2:06:13our case it is running on Port Local
- 2:06:16Host 9092 make sure before you start
- 2:06:19your producer and consumer you should
- 2:06:21start your Juke keeper then your capka
- 2:06:23server okay and also I'll just change
- 2:06:26the server Port
- 2:06:28here server
- 2:06:31Port
- 2:06:329292 because I believe my producer is up
- 2:06:36and running on Port 91 okay 91 91 so I
- 2:06:40just change it to the 9292 that's it now
- 2:06:43let me quickly create a package so that
- 2:06:46we can start creating a class who will
- 2:06:47do the actual logic or who will try to
- 2:06:51face the messages from the topic okay
- 2:06:53let me create a
- 2:06:56package then we'll create a
- 2:07:00class now just unnoted this class with
- 2:07:02adate
- 2:07:04service fine now we'll write a method
- 2:07:08who will do the actual logic to retrieve
- 2:07:10the message from the
- 2:07:11topic public void so what I'll do I need
- 2:07:15to an your SL forj okay I didn't added
- 2:07:18the lbo fine I'll create the logger
- 2:07:20object
- 2:07:27then complete the
- 2:07:28method consume that's
- 2:07:32it now which type of data we are sending
- 2:07:35from the producer if you go and check in
- 2:07:38your
- 2:07:39producer go to the class we are sending
- 2:07:43message of type string right so you need
- 2:07:47to specify the type what messages this
- 2:07:49method will consume so I'll will specify
- 2:07:53string
- 2:07:55message fine now I'll just add the log
- 2:07:58statement to print that
- 2:08:00messages logger or just change it to the
- 2:08:07log just add some messages
- 2:08:10okay
- 2:08:12consumer consume the
- 2:08:17message and the messages we are
- 2:08:20receiving
- 2:08:22fine we are not doing anything we are
- 2:08:24simply write a method to consume the
- 2:08:26message and we are just simply printing
- 2:08:28them now how this consumer will know
- 2:08:32from which topic he need to read the
- 2:08:34messages that you need to tell to the
- 2:08:37you need to tell to the consumer okay
- 2:08:40how you can tell that simple The
- 2:08:41annotation is capka listener now in this
- 2:08:46Capa listener annotation you can tell to
- 2:08:48the application or you can tell to this
- 2:08:51consumer who is the topic from where you
- 2:08:53want to read the messages okay so in our
- 2:08:57case what is the topic we created in our
- 2:08:58producer
- 2:09:00application Java demo right we are
- 2:09:04creating the topic in our config class
- 2:09:07this is where we are creating the topic
- 2:09:10with three partition so you need to
- 2:09:13specify the same topic name in consumer
- 2:09:16from where you want to listen the
- 2:09:18messages now all good let me start the
- 2:09:21the producer then we'll start the
- 2:09:23consumer app so go to the
- 2:09:26producer start it once it will up we'll
- 2:09:30start our
- 2:09:31consumer so producer is up on put 9191
- 2:09:35now let me start the consumer app go to
- 2:09:38the main class and simply start
- 2:09:41it meanwhile I will open the
- 2:09:44postman let's check the producer code in
- 2:09:48our
- 2:09:49application just go to the public
- 2:09:52not here go to the controller just
- 2:09:54wanted to check how many messages we are
- 2:09:56publishing okay 10,000 right that's fine
- 2:09:59we can send the 10,000 and we'll see
- 2:10:02whether it is getting logged in our
- 2:10:03consumer or not then we'll change
- 2:10:06it so we'll go and check our consumer
- 2:10:09app we are getting the error now if You
- 2:10:13observe carefully the error clearly says
- 2:10:16that no group ID found in consumer
- 2:10:19config as I already explained even
- 2:10:22though you have a single consumer
- 2:10:24instance still you need to map it to a
- 2:10:27consumer ID but in our code we have not
- 2:10:30specify any consumer ID okay that is the
- 2:10:33reason kapka is giving error hey just
- 2:10:36map a group ID to your consumer instance
- 2:10:40so if you'll
- 2:10:42Define group ID here group ID and let me
- 2:10:47give something like JT hyphen group
- 2:10:52hen one something like that okay and
- 2:10:54also I need to configure this particular
- 2:10:57group ID in yml file or application or
- 2:11:00properties file as well okay because
- 2:11:02that is the central place where we are
- 2:11:04specifying our configuration so
- 2:11:09spring okay so already we have defined
- 2:11:11this so in the consumer section only you
- 2:11:14can Define group ID what is the group ID
- 2:11:17JT group one this is the group ID we
- 2:11:20have created you can any unique name
- 2:11:22okay it will just help you to find the
- 2:11:25role of your consumer group or to find
- 2:11:27the role of your this particular
- 2:11:30consumer Group which task it is doing
- 2:11:32that is the reason you can segregate
- 2:11:33using this consumer or group ID okay
- 2:11:37that's fine now go to The Listener we
- 2:11:41have defined the same and we have
- 2:11:42defined same in our config now let me
- 2:11:44restart it
- 2:11:47again so consumer app is up and running
- 2:11:50now if You observe
- 2:11:52carefully let me Zoom this for you JT
- 2:11:56hyen group hyen one is my consumer group
- 2:12:00ID right this particular consumer group
- 2:12:04is assigned to all the three partition
- 2:12:08can you see
- 2:12:10here Java key demo 0o demo 1 demo 2 the
- 2:12:14Java H demo is my topic name and
- 2:12:17partition name happened with the number
- 2:12:190 1 and two and all the three partition
- 2:12:22of my topic is being assigned to single
- 2:12:26consumer instance by specifying this
- 2:12:29particular group ID so that is what the
- 2:12:32first step we have discussed here a
- 2:12:34single consumer instance is trying to
- 2:12:36read from multiple partition this is the
- 2:12:39happy scenario what we are trying now
- 2:12:41lat will increase the consumer instance
- 2:12:44and we'll see whether the each consumer
- 2:12:46instance is pointing to each partition
- 2:12:48or not that is the next step okay now
- 2:12:50for now
- 2:12:51you can read this particular statement
- 2:12:54to understand this is the consumer
- 2:12:56instance and this is the group ID is
- 2:12:58being assigned to all three partition
- 2:13:00that's fine let me clear this go to the
- 2:13:03postman let me send the message it will
- 2:13:06send 10,000 messages right it will
- 2:13:09definitely take few
- 2:13:11second message published successfully
- 2:13:14and it take
- 2:13:15157
- 2:13:17millisecond and now let's check in our
- 2:13:19consumer app can you see
- 2:13:22here consumer consume the message user
- 2:13:26this one I mean this is the
- 2:13:28message this is the message and this is
- 2:13:30my log statement consumer consume the
- 2:13:33message this all the 10,000 messages is
- 2:13:37being consumed okay now how you can see
- 2:13:40that all the 10,000 messages what you
- 2:13:42have published from producer is being
- 2:13:44consumed by your consumer simply go to
- 2:13:47your upset Explorer then refresh the
- 2:13:50topic
- 2:13:51you'll find the Java TI demo as a topic
- 2:13:55now let me refresh
- 2:13:56this let me yeah you can see here right
- 2:14:00this is what our topic and this is what
- 2:14:02the consumer group we have created now
- 2:14:04first let's verify the data in topic so
- 2:14:08just run
- 2:14:10it now if you'll go and check in the
- 2:14:13Java TI demo go inside the partition you
- 2:14:17can verify each partition I mean getting
- 2:14:20the messages
- 2:14:21partition 0 received
- 2:14:235,150 partition 1 received
- 2:14:272346 partition 2 also received something
- 2:14:302505 I mean all the three partition
- 2:14:33received the messages okay now if you'll
- 2:14:36go and check in the consumer group if I
- 2:14:40will open the
- 2:14:41offset can you see here the topic name
- 2:14:44is Java demo and message received to all
- 2:14:49the three partition 102 two okay and
- 2:14:53count of each partition messages each
- 2:14:56messages in Partition you can see here
- 2:14:58partition one is this much zero is this
- 2:15:01much and two is this much now upset is
- 2:15:03the number of your
- 2:15:05sequence now if You observe here lag is
- 2:15:08zero now all the partition whatever the
- 2:15:12messages they have received they are
- 2:15:14able to deliver that messages to the
- 2:15:16consumer I mean consumer successfully
- 2:15:19read 23 6 messages from partition 1
- 2:15:235,150 messages from partition 0
- 2:15:27255 messages from partition two that is
- 2:15:30the reason you can see the lag count is
- 2:15:32zero there is no message is pending from
- 2:15:35your producer side to consumed by your
- 2:15:38consumer all the messages whatever your
- 2:15:40producer is published is being consumed
- 2:15:43by your consumer okay so by seeing this
- 2:15:45lag number you can understand how many
- 2:15:48messages your consumer is not received
- 2:15:50okay this is what I always prefer to use
- 2:15:53this offset Explorer because this will
- 2:15:55give you the lag information okay that's
- 2:15:57fine now what do you understand
- 2:16:00partition 012 received all the messages
- 2:16:03but we have only one consumer instance
- 2:16:05who is listening to these three
- 2:16:08partition now can we try creating three
- 2:16:12consumer instance so that we'll verify
- 2:16:15whether all the three instance is
- 2:16:17pointing to individual partition or not
- 2:16:19we'll do that right away let's go to the
- 2:16:22code you can create a separate class and
- 2:16:25Define the this particular Capa listener
- 2:16:28and you can specify from which topic you
- 2:16:31want it to read but we'll simply create
- 2:16:35duplicate method to consume the messages
- 2:16:37okay rather than create duplicate class
- 2:16:40so I'll just copy the messages I'll will
- 2:16:43create three consumer instance okay let
- 2:16:45me Zoom this for you this is the second
- 2:16:49one so I'll just Define 2 I'll Define
- 2:16:54one now I created another one this this
- 2:16:57is 1 2 3 right we have created three
- 2:17:00consumer instance also as for the flow I
- 2:17:03explained right if there is a fourth
- 2:17:05consumer instance then what you'll do
- 2:17:08because we have three partition all the
- 2:17:09three partition will assign to each and
- 2:17:11every consumer instance then what is the
- 2:17:14role of consumer 4 we'll see that as
- 2:17:16well okay so we'll create another
- 2:17:19consumer instance for
- 2:17:21backup so I'll name
- 2:17:23it anything you can give consumer
- 2:17:26consume four that's it now I'll just
- 2:17:29change the log statement otherwise it is
- 2:17:31difficult for us to filter out who
- 2:17:33received the messages okay so I'll just
- 2:17:36Define this is consumer 1 consumer 2
- 2:17:40consumer
- 2:17:423 and four that's it okay group ID
- 2:17:47everything is same we have just created
- 2:17:49multiple consumer instance
- 2:17:51so we'll see messages will come to if
- 2:17:54messages will come to two partition then
- 2:17:57two consumer instance should read those
- 2:18:00messages from the partition it can be
- 2:18:03consumer one and consumer 2 or three and
- 2:18:05four any order that is depends on the
- 2:18:07Joe keeper coordinator this is not in
- 2:18:09our hand that's fine so for safer side
- 2:18:13what we'll do we'll just change the
- 2:18:15topic name and group ID okay I'll change
- 2:18:18it to the group J
- 2:18:21group same I'll replace in all the
- 2:18:27method also I'll change the topic
- 2:18:30name same I need to do the change in my
- 2:18:33producer as well right so producer will
- 2:18:37publish the message
- 2:18:38here and also in config right that's it
- 2:18:43so we have created the same three
- 2:18:44partition let me restart the
- 2:18:48producer producer is up now let's go and
- 2:18:51start our
- 2:18:53consumer now if you see the statement
- 2:18:56consumer is up and running three
- 2:18:58individual consumer instance is being
- 2:19:01assigned to individual partition you can
- 2:19:04see here right our topic name Java demo
- 2:19:07one hyphen 1 hyphen 2 Zer is nothing
- 2:19:10your partition number and it is just
- 2:19:12giving the group name because we don't
- 2:19:14specify our consumer instance name all
- 2:19:16the consumer instance we are wrapped
- 2:19:18inside a group ID so you can see
- 2:19:21instance one instance two and instance
- 2:19:23three is pointing to three different
- 2:19:26partition from the topic now let's
- 2:19:28verify that right away okay clear
- 2:19:31it I'll send the request now my producer
- 2:19:35will publish 10,000 message to the topic
- 2:19:37and the messages will distribute to
- 2:19:40three different topic okay now then we
- 2:19:43have three consumer instance we'll see
- 2:19:45whether all the messages is being
- 2:19:48distributed to individual uh partition
- 2:19:50and and consumer instance combination or
- 2:19:52not so go to the postman simply hit the
- 2:19:57request it publish the 10,000 messages
- 2:20:01first let's verify how many partition
- 2:20:04received the messages okay then based on
- 2:20:07that we'll validate our consumer
- 2:20:08instance go and check here let me
- 2:20:11refresh the topic also Let me refresh my
- 2:20:14consumer we have the new consumer JT
- 2:20:17group we'll verify here only okay
- 2:20:21partition 1 2 0 all three partition
- 2:20:24receive the messages there is no lag it
- 2:20:26means there should be three consumer
- 2:20:29instance log statement we can see okay
- 2:20:33if you are not seeing three consumer
- 2:20:34instance log then our concept is wrong
- 2:20:37what we discuss is wrong okay so let's
- 2:20:40verify it right away so how we can
- 2:20:42filter it simple right that is the
- 2:20:44reason I have changed the first
- 2:20:46statement so I'll just filter here
- 2:20:49consumer one
- 2:20:51received
- 2:20:532,445 messages now we'll see who else
- 2:20:57received the messages does Consumer 2
- 2:20:59received no consumer 3 yes consumer 3
- 2:21:03received
- 2:21:052371 okay now we'll see consumer 4 also
- 2:21:08received something yeah
- 2:21:112176 so what we cleared here we have
- 2:21:14four instance consumer 1 2 3 4 consumer
- 2:21:181 3 and four received the messages you
- 2:21:20can see in four we can see some messages
- 2:21:23right and the count you can see here now
- 2:21:25if you change it to the
- 2:21:27three consumer three three also receive
- 2:21:32the messages right now if you change it
- 2:21:34to the one one also receive the messages
- 2:21:37but two consumer is not receive the
- 2:21:40messages okay consumer 1 three and four
- 2:21:43only these three consumer instance
- 2:21:45receive the messages not this
- 2:21:48guy why simple statement right we have
- 2:21:52three partition that is the reason it
- 2:21:53assigned to the three consumer instance
- 2:21:56one consumer instance two consumer
- 2:21:59instance three consumer instance if
- 2:22:02message is being pushed to only two
- 2:22:05partition let's say uh partition one and
- 2:22:08partition two then only two consumer
- 2:22:11instance will print the statement that
- 2:22:13is how the capka maintain the balance
- 2:22:16between the partition and consumer
- 2:22:18instance okay now what I'll do the order
- 2:22:21is not same right now this time consumer
- 2:22:23two didn't receive the messages but
- 2:22:24there is no guarantee that when the
- 2:22:26message will come next time consumer 2 3
- 2:22:294 or 1 2 3 like that anyone can receiv
- 2:22:32it okay we'll just verify it right
- 2:22:35away now I'll again send 10,000
- 2:22:40messages fine now we'll see who all
- 2:22:43received the messages consumer 3
- 2:22:46received the
- 2:22:48messages fine consumer 2 no it didn't
- 2:22:52receive consumer one receive consumer 4
- 2:22:54receive okay if you'll try many times
- 2:22:57there might be you can see the change in
- 2:23:00order okay it can be 1 2 3 or 3 to 1 any
- 2:23:03any order we can't predict it it will be
- 2:23:05decided by
- 2:23:07coordinator that's fine so I believe
- 2:23:10you're pretty clear with this concept
- 2:23:12okay number of partition we have number
- 2:23:14of consumer instance make the mapping
- 2:23:17and get the messages that is the simple
- 2:23:20statement
- 2:23:21okay but just keep a note in real time
- 2:23:24you should not write multiple consumer
- 2:23:27instance like this okay so this is just
- 2:23:30I give you the demo to give the clear
- 2:23:34picture about this partition and
- 2:23:36consumer instance mapping in real time
- 2:23:39you can make it in a better way by
- 2:23:40implementing the concurrency so you can
- 2:23:43design or you can increase the
- 2:23:44throughput of your application okay just
- 2:23:47keep a note again this is just for the
- 2:23:49demo purpose so you should not write the
- 2:23:51code like this in your consumer that's
- 2:23:54fine now let's understand if I have
- 2:23:58published 10,000 message to the topic
- 2:24:01whether all the messages is being
- 2:24:03consumed by my consumer or not let's say
- 2:24:05producer publish the 10,000 messages but
- 2:24:08in between my capka consumer is shut
- 2:24:11down then that time how can I verify
- 2:24:14that how many messages my consumer
- 2:24:17received or how many messages he didn't
- 2:24:19received that is what the lag we just
- 2:24:21discuss now right but let's prove it I
- 2:24:24mean we'll forcefully stop our capka
- 2:24:26consumer and we'll validate how the lag
- 2:24:29show the exact Behavior or not okay or
- 2:24:32whether it shows the exact number or not
- 2:24:34fine so what I'll do I'll just open my
- 2:24:38Postman so I'll I'll keep this in two
- 2:24:44different uh section so that I can
- 2:24:47quickly stop my
- 2:24:49consumer
- 2:24:51Okay cool so let me clear
- 2:24:54everything I'll send the request then
- 2:24:58immediately I will stop
- 2:25:00it I don't know how many messages is
- 2:25:02being published because from the
- 2:25:04producer we have published the 10,000
- 2:25:06messages okay so the topic name is also
- 2:25:09same let's verify it right away go to
- 2:25:11the upset
- 2:25:14Explorer just reconnect
- 2:25:17it open the topic just Java tiet demo
- 2:25:22one go into the partition but it won't
- 2:25:24give you any clue here if you look into
- 2:25:26the topic better you go to the
- 2:25:29consumer there is no lag because we not
- 2:25:33able to shut down immediately okay so
- 2:25:37what I'll do I'll publish more messages
- 2:25:40so that I will get enough time to shut
- 2:25:42down my consumer so I'll go to my
- 2:25:44producer
- 2:25:46application then this is how many 10,000
- 2:25:49right I will send
- 2:25:50one L messages and also I'll just change
- 2:25:54my topic
- 2:25:56name uh go to the config here okay we
- 2:26:00did change in config change this topic
- 2:26:02name now go to your consumer
- 2:26:06code I mean meanwhile let me start my
- 2:26:09producer right we have did the changes
- 2:26:12just restart it I'll go to the
- 2:26:16consumer in consumer also you need to
- 2:26:19change the topic name
- 2:26:22uh let's say two only right and I'll
- 2:26:25change the group J group hyphen
- 2:26:30new let me specify in all the
- 2:26:35instance fine we have changed the
- 2:26:38consumer group and we have changed the
- 2:26:40topic name in both producer and consumer
- 2:26:42and we start our producer now let me
- 2:26:44start my
- 2:26:48consumer okay consumer is up and running
- 2:26:51let me clear
- 2:26:54everything let me send the request now
- 2:26:57now it will send one lakh record to the
- 2:26:59capka okay now send the request and
- 2:27:03meanwhile I will stop
- 2:27:06it I'm able to stop or not I'm not sure
- 2:27:09but I've have tried it but let's
- 2:27:11see I will just reconnect it okay so
- 2:27:15that it will refresh
- 2:27:17everything go to the topic this is the
- 2:27:20topic newly we created right go to the
- 2:27:23consumer J group new is there yeah we
- 2:27:26are able to deliver it Java demo
- 2:27:302 can you see
- 2:27:32here not don't don't verify the demo one
- 2:27:35this is the recent topic we created in
- 2:27:37demo 2 partition one it received
- 2:27:4132358 messages and only 976 message is
- 2:27:47being consumed by your consumer that is
- 2:27:49the reason you are able to to see the
- 2:27:50upset count rest
- 2:27:5331382 message is not being consumed by
- 2:27:56your consumer okay you won't find those
- 2:27:59messages in your consumer section okay
- 2:28:02now if you'll see the partition zero it
- 2:28:06received
- 2:28:073276 messages it was able to deliver
- 2:28:103,217 messages 29
- 2:28:149,549 messages is not being consumed
- 2:28:17producer is able to publish the messages
- 2:28:19but but somehow your consumer is shut
- 2:28:21down that time so he is not able to
- 2:28:23receive those messages okay now you can
- 2:28:26easily identify right which partition is
- 2:28:29don't I mean which partition messages is
- 2:28:32not being consumed by your consumer so
- 2:28:35in real time in production issue you
- 2:28:37might need to republish your messages or
- 2:28:39you need to publish your messages again
- 2:28:41to resolve such kind of issue so if you
- 2:28:43have this tool at least you can figure
- 2:28:45out right the number of messages in
- 2:28:47which partition okay that will give you
- 2:28:50the clue and you can publish the
- 2:28:52messages from that specific partition
- 2:28:54okay and you can see here partition
- 2:28:57two 34,000 messages he received but he
- 2:29:01only published
- 2:29:033,257 because these are the messages is
- 2:29:05being consumed by consumer 31620 is not
- 2:29:09consumed okay so this is how you can
- 2:29:12figure out the lag in your producer and
- 2:29:14consumer so it will make your job easier
- 2:29:17or it will give you the better
- 2:29:18productivity while analyzing the issue
- 2:29:21on capka okay don't worry I will also
- 2:29:24tell you in my upcoming tutorial how you
- 2:29:26can retrive a specific messages from the
- 2:29:29partition we'll go in deeper okay with
- 2:29:32the code in my upcoming session so I
- 2:29:34hope you are clear about the consumer
- 2:29:37and consumer instance and partition
- 2:29:39mapping and how to verify the lag in the
- 2:29:43upset Explorer
- 2:29:49two
- 2:29:52[Music]
- 2:29:53in the previous capka tutorial we
- 2:29:55understood the complete capka pops up
- 2:29:57mechanism right where we publish the
- 2:29:59plan string message to the topic and
- 2:30:01from the topic our consumer application
- 2:30:03consumed it even we Deep dive more to
- 2:30:06understand partition and consumer group
- 2:30:08mapping with examples but the pop sub
- 2:30:11system which we designed will it accept
- 2:30:14any data type to publish and consume the
- 2:30:16event the answer is big no let's prove
- 2:30:19it right away then we'll find out a
- 2:30:21solution for it okay so to verify this
- 2:30:24Behavior instead of sending the string
- 2:30:26message to the topic and consumed it
- 2:30:28let's send some object or let's create
- 2:30:31some pujo class which we want to send to
- 2:30:33the capka topic and we want that PUO
- 2:30:36need to be consumed by my consumer so
- 2:30:38for that what I'll do I'll simply create
- 2:30:40a new
- 2:30:42package dto now I'll create some class
- 2:30:45here let's say customer
- 2:30:47class now let's define couple of
- 2:30:57fi so we have defined ID name email and
- 2:31:00contact number so we added the lbook
- 2:31:03right so I can use theate data
- 2:31:05annotation to avoid okay we have not
- 2:31:08added the lumbo no is go to your pom.xml
- 2:31:12and then you can add the lbook
- 2:31:15dependency just update it that's it now
- 2:31:19we can input this statement so we don't
- 2:31:21need to write the getter and Setter
- 2:31:22because we already added the at theate
- 2:31:24data annotation that's fine now I want
- 2:31:27to send this customer object to the
- 2:31:29capka topic from the producer and I want
- 2:31:32this customer object need to be consumed
- 2:31:34by my consumer application so let's do
- 2:31:37that go to the publisher where exactly
- 2:31:40we are publishing the message so let me
- 2:31:43go to that
- 2:31:45class we have defined the publisher
- 2:31:47right so I don't want to touch the
- 2:31:49existing in code let it be I'll create a
- 2:31:52new
- 2:31:54method so I'll name
- 2:31:56it let's say send events to
- 2:32:01topic and here I don't want to send the
- 2:32:04message I want to send the object which
- 2:32:07is
- 2:32:09customer fine now everywhere I need to
- 2:32:13change
- 2:32:16it fine so I will also change thep toic
- 2:32:20name that will do later now let me send
- 2:32:23the object so this is the exact place
- 2:32:27where we are sending the object okay
- 2:32:30earlier we are sending the string
- 2:32:31message now we are sending one custom
- 2:32:33object which is customer that's fine now
- 2:32:36we'll call this method from our
- 2:32:38controller go to the controller
- 2:32:41class then simply write another endpoint
- 2:32:45so I'll write something
- 2:32:46like
- 2:32:48public
- 2:32:50void send events okay and what type of
- 2:32:55event I want to send customer
- 2:32:58object now simply I can call Publisher
- 2:33:01dot send events to the topic and
- 2:33:05customer and I need to Define this as a
- 2:33:07endpoint so what I'll do I'll just
- 2:33:10Define post mapping because I want to
- 2:33:13send the object from the postman and I
- 2:33:15need to annotate this at theate request
- 2:33:18body that's it I need to define the URL
- 2:33:22here
- 2:33:24publish okay so what we are doing here
- 2:33:27from the postman we'll send the customer
- 2:33:29object in the form of Json and then we
- 2:33:32are giving that to the capka template to
- 2:33:34send this particular customer object to
- 2:33:36this given topic okay so what we'll do
- 2:33:39we'll also change the topic name so go
- 2:33:42to the config this is the place where we
- 2:33:44are creating the topic I'll name it Java
- 2:33:47demo and go to the controller class not
- 2:33:50controller go to the publisher
- 2:33:53class there also we need to change the
- 2:33:56topic name so you want to send the
- 2:33:58customer object to this particular Java
- 2:34:01demo topic okay now once we send it to
- 2:34:05the topic who will consume it our
- 2:34:07consumer now also we need to tell to the
- 2:34:10consumer hey consumer please consume
- 2:34:13customer object instead of string right
- 2:34:16so I also need to do the code change in
- 2:34:18my consumer so go to the consumer
- 2:34:20application first let me copy this dto
- 2:34:25class in consumer also you need to tell
- 2:34:28right what need to be consumed so what
- 2:34:31I'll do I will also comment all the
- 2:34:33consumer group related Stu this is my
- 2:34:37consumer he will point to the Java key
- 2:34:39demo what what I have defined just now
- 2:34:41in my producer and instead of string
- 2:34:44message we we need to tell to this
- 2:34:46consumer please consume the object
- 2:34:49object which is customer what data type
- 2:34:52you need to consume you can specify here
- 2:34:55okay so I'll change the method name to
- 2:34:57consume so it is crying because we don't
- 2:34:59have the customer class in this
- 2:35:01particular project so rather than keep
- 2:35:03the same class in both producer and
- 2:35:05consumer application you can create a
- 2:35:07multimodule project and you can create a
- 2:35:09common module and you can keep all the
- 2:35:12common dto or PUO class okay but for
- 2:35:15this demo purpose I'm just adding
- 2:35:18it so you have added the customer again
- 2:35:21we don't have the lbook dependency here
- 2:35:23so just go to the
- 2:35:24producer and copy the lbook
- 2:35:29dependency then simply go to
- 2:35:32the consumer pom.xml and paste
- 2:35:35it all good go to The
- 2:35:38Listener and just input the
- 2:35:41statement so I will just print your
- 2:35:44customer dot two string I just want to
- 2:35:47print it the plain string okay
- 2:35:50okay consumer consume the
- 2:35:55events fine we are also specifying the
- 2:35:58topic name and group ID so group ID you
- 2:36:01can give any
- 2:36:02name that's
- 2:36:04it just go to the producer one
- 2:36:08again let me verify we have defined this
- 2:36:12right also what I want to do I'll keep
- 2:36:14this in TR catch so that if there is any
- 2:36:16exception while publishing the event we
- 2:36:19can easily catch
- 2:36:22it just add another bracket that's it so
- 2:36:25here what I do I just print it see out
- 2:36:29let's say
- 2:36:30error and what is the error message we
- 2:36:33got if there is
- 2:36:35any get
- 2:36:37message that's it right so I have
- 2:36:39started my juke keeper and capka server
- 2:36:42that is the first step you need to do
- 2:36:44before you start your application now we
- 2:36:46are all good let's run our application
- 2:36:48so let me start the producer
- 2:36:52app once it will up then I will quickly
- 2:36:54start my consumer
- 2:36:57app so producer is up on put 9191 go to
- 2:37:01the consumer and just start
- 2:37:07it so it is up and running also it
- 2:37:10clearly says that the group ID what you
- 2:37:12have created is being assigned to the
- 2:37:14topic of three partition 012 because we
- 2:37:18are creating Java demo
- 2:37:20topic with three partition that is what
- 2:37:21we have written in our cavka config
- 2:37:23that's fine let me clear this this is my
- 2:37:27consumer right and go to the producer
- 2:37:30let me clear the console fine now what
- 2:37:33Endo we need to trigger this particular
- 2:37:36endpoint SLP publish and you need to
- 2:37:38give the customer as a Json so I'll go
- 2:37:41to the
- 2:37:42postman let me check the
- 2:37:46endpoint producer of this is are the
- 2:37:50producer app and publish so I will send
- 2:37:53the customer object ID name email and
- 2:37:56contact number now let me send this
- 2:37:59object to the from the producer app to
- 2:38:02the topic let's see what is the result
- 2:38:03we are getting status is
- 2:38:0620 let's see what if there is any
- 2:38:10error can you see here we are getting
- 2:38:14the error here can't convert value of
- 2:38:18class this D customer to Apachi kapka
- 2:38:21common serial serialization string
- 2:38:24serializer specified in value
- 2:38:26serializer what it says usually when we
- 2:38:29are sending the messages to the topic it
- 2:38:33will always accept bite AR as a input
- 2:38:36because we are serializing the object we
- 2:38:38are sending the object over the network
- 2:38:40to whom topic so if you look into the
- 2:38:44presentation from the producer to the
- 2:38:46topic we are serializing the data okay
- 2:38:49in the form of object or in the form of
- 2:38:51Json but in the consumer perspective we
- 2:38:55are deserializing the data from producer
- 2:38:58which we are sending the bite array over
- 2:39:00the network to the topic and in consumer
- 2:39:03it consumes the bite array from the
- 2:39:05topic so this is the part of
- 2:39:07serialization and deserialization
- 2:39:09concept basic core Java if some object
- 2:39:12you'll send over the network that must
- 2:39:14need to be serialized and if you'll read
- 2:39:17that particular object then that must
- 2:39:19need to be that that concept is called
- 2:39:21der serialize right that is what exactly
- 2:39:23you are doing producer is serializing
- 2:39:26the object to the topic and consumer is
- 2:39:28dis realizing the object from the topic
- 2:39:31so that is the reason we are getting the
- 2:39:33error here because Capa looking for bite
- 2:39:36array to be input so that is the reason
- 2:39:38we can happily play with the string
- 2:39:40without adding any configuration if
- 2:39:42you'll send any string that will be
- 2:39:44converted to the bite array and will
- 2:39:46publish to the topic and consumer can
- 2:39:48consume that bite array can convert it
- 2:39:50to the string but if you want to play
- 2:39:52with the different type of object then
- 2:39:54you need to tell to the capka hey capka
- 2:39:58while doing the serialization please use
- 2:40:00this serializer or in the consumer you
- 2:40:03can tell him hey consumer while
- 2:40:06deserializing that particular object
- 2:40:08please use the deserializer what I am
- 2:40:10giving to you tell that I'm sending the
- 2:40:13customer object please serialize it and
- 2:40:15consumer need to tell to the capka hey I
- 2:40:18am expecting consumer object please
- 2:40:20derealize it now how you can tell to the
- 2:40:22capka very simple just go to the
- 2:40:25application. IML
- 2:40:26file you need to mention in both
- 2:40:29producer and consumer because from the
- 2:40:31producer you are serializing the data
- 2:40:34and from the consumer you are
- 2:40:35deserializing the data okay so go to the
- 2:40:38resource application. IML let me Zoom
- 2:40:41this for
- 2:40:42you now here you can tell him hey capka
- 2:40:48please use this is key serializer and
- 2:40:51what is the key serializer I'm sending
- 2:40:53the key in the form of string so use
- 2:40:56this string serializer to just convert
- 2:40:59the key okay and for Value serializer
- 2:41:03use the Json
- 2:41:05serializer just specify that very simple
- 2:41:09so I'm sending the key in the form of
- 2:41:11string you just play with the string
- 2:41:13serializer but the value what I'm
- 2:41:16sending the customer object for that
- 2:41:18please you use this Json serializer
- 2:41:21because I'm sending in the form of Json
- 2:41:23you please take that Json and send it to
- 2:41:25the topic okay and now same thing you
- 2:41:28need to tell in your consumer also so go
- 2:41:31to the
- 2:41:33consumer go to the application.
- 2:41:37yl now here in the consumer also you
- 2:41:40need to tell him what is your key
- 2:41:42deserializer see the word here see the
- 2:41:44key here in consumer we are
- 2:41:47deserializing the data right so that is
- 2:41:49the region key dis realizer what is the
- 2:41:51key dis realizer that is the same string
- 2:41:54string dis realizer right so just Define
- 2:41:58it now what is the value D realizer same
- 2:42:02Json D realizer just Define that can you
- 2:42:06see here key serializer key der
- 2:42:08serializer is the string deserializer
- 2:42:10and value der serializer is the Json D
- 2:42:13realizer because we are not giving any
- 2:42:15key as per our example okay we are
- 2:42:17playing with the value so it's fine you
- 2:42:19can keep string serializer and der
- 2:42:21serializer as a key for now but you're
- 2:42:24are sending the object so serializer and
- 2:42:26deserializer must be Json one because we
- 2:42:29sending the object in the form of Json
- 2:42:32fine all good now let's run our producer
- 2:42:35and consumer to verify whether the way
- 2:42:38we have configured for the object it is
- 2:42:40working or not my capka is able to
- 2:42:43publish the object to the topic and my
- 2:42:45consumer is consume that object from the
- 2:42:47topic or not okay so very simple just
- 2:42:50restart go to the producer first restart
- 2:42:54it then go to the
- 2:42:57consumer here also you can simply
- 2:42:59restart
- 2:43:01it both producer and consumer is up and
- 2:43:04running you can see here this is the
- 2:43:06consumer let me clear the console go to
- 2:43:09the producer let me clear the
- 2:43:11console that's fine now let me send the
- 2:43:14request
- 2:43:17okay they're getting status code 20 and
- 2:43:21then let's see whether is there any
- 2:43:24error there is no error from the
- 2:43:26producer producer successfully send this
- 2:43:29particular object you can see here okay
- 2:43:31because we are just printing the two
- 2:43:33string now let's verify in the
- 2:43:36consumer there is a exception in
- 2:43:38consumer let me stop it
- 2:43:41now let's see what is the error we are
- 2:43:44getting from the consumer okay the class
- 2:43:47com. java. d.com customer is not the
- 2:43:51trusted package Java util Java Lang if
- 2:43:54you believe this class is safe to dis
- 2:43:56realize please provide its name if the
- 2:44:00if the serialization is only done by a
- 2:44:02trusted trusted Source you can also
- 2:44:04enable trust all okay so very simple
- 2:44:07okay the exception it clearly says that
- 2:44:10the D what you are sending or what your
- 2:44:14consumer is trying to consume is not a
- 2:44:16trusted package so you need to tell to
- 2:44:18the
- 2:44:19hey if you're getting any object from
- 2:44:22this particular package please accept it
- 2:44:25please listen to it okay that is what
- 2:44:28you can Define so just mention that in
- 2:44:31our yl file so what I'll do I will tell
- 2:44:34here
- 2:44:36properties
- 2:44:38spring
- 2:44:41Json fine
- 2:44:44trusted and specify the packages which
- 2:44:48packages you trusted to use or you are
- 2:44:50allowing your consumer to read that so
- 2:44:53what you can do for now our package is
- 2:44:56come.
- 2:44:58java. dto if you want to include any
- 2:45:02kind of package just Define the star
- 2:45:05instead of specify the single package
- 2:45:07name okay com. java. dto now let me
- 2:45:11restart the consumer producer is up and
- 2:45:13running he is able to publish the
- 2:45:15message there is no error but consumer
- 2:45:18was crying so if fixed it let's see
- 2:45:20whether it is able to consume or not
- 2:45:21yeah can you see
- 2:45:23here now my consumer consume the
- 2:45:27previous event which is the object this
- 2:45:30is what we have sent right now to cross
- 2:45:33verify it let me clear this let me clear
- 2:45:37the producer as well and then what you
- 2:45:40can do let me publish some different
- 2:45:42event
- 2:45:46okay fine let me send the request
- 2:45:50go to the
- 2:45:52producer producer send the messages to
- 2:45:56the topic now go to the
- 2:45:59consumer can you see here consumer
- 2:46:02consume the event this is what just now
- 2:46:04we sent
- 2:46:05right we're able to consume it fine all
- 2:46:10good we are able to play with the object
- 2:46:12from the producer to the consumer but is
- 2:46:16there the only one way to play with the
- 2:46:18object object by defining just the
- 2:46:20serializer and deserializer in
- 2:46:22application. IML file no you can
- 2:46:26customize that by defining your Java Bas
- 2:46:28config class okay rather than add the
- 2:46:31key and value or whatever the properties
- 2:46:34in the application.yml file still you
- 2:46:36have option to go with the Java Bas
- 2:46:38configuration if you want to do any
- 2:46:40customization okay I will also show you
- 2:46:42that part how you can go with the Java
- 2:46:45base configuration first let me comment
- 2:46:48this piece of code because we are not
- 2:46:49going to use the application. IML
- 2:46:52configuration we want to write our own
- 2:46:54Java base configuration go to the
- 2:46:56producer as
- 2:46:58well let me not all actually these
- 2:47:02things fine now go to the config class
- 2:47:07which is Capa producer config and then
- 2:47:09here you can Define your own
- 2:47:11configuration when I say own
- 2:47:13configuration whatever the things key
- 2:47:16and value you have defined in this
- 2:47:18application IML same things you need to
- 2:47:21configure in the Java code I mean you
- 2:47:23need to create a bin of it okay so I'll
- 2:47:26show you how we can do
- 2:47:28that so first create a producer config
- 2:47:31map so
- 2:47:34public map of
- 2:47:37string
- 2:47:40object now create a object of
- 2:47:44map now all the key whatever we have
- 2:47:47defined in application. yml
- 2:47:49that needs to be configured here so you
- 2:47:51can Define like this okay producer
- 2:47:54config what is the bootstrap server
- 2:47:56Local Host 9092 what is your key
- 2:47:58serializer what is your value serializer
- 2:48:01so you just need to Define this at the
- 2:48:05r bin fine now by giving this producer
- 2:48:10config you can create producer Factory
- 2:48:12object so
- 2:48:17public now simply you can return return
- 2:48:19here new default Capa producer Factory
- 2:48:23and you need to give this config
- 2:48:25whatever you have Define about your um
- 2:48:28capka server serializer all the things
- 2:48:31okay so just Define that
- 2:48:33config and annotate this at theate
- 2:48:36bin now again you need to create capka
- 2:48:39template object by giving this producer
- 2:48:41Factory now your capka template will
- 2:48:43have all the information okay because
- 2:48:45he's not going to read from yml he is
- 2:48:48going to read from this config map so
- 2:48:50let me create
- 2:48:55it just Define this at theate bin that's
- 2:48:59it okay so we Define the config map
- 2:49:02where you have defined all the K and
- 2:49:03value then we created the producer
- 2:49:05Factory and then we give that producer
- 2:49:08Factory to the capka template now same
- 2:49:11things you need to write in your
- 2:49:13consumer as well in consumer also you
- 2:49:15can specify the consumer config map by
- 2:49:18defining all the D realizer property and
- 2:49:21all the bootstrap server related
- 2:49:23properties okay so that is also
- 2:49:25straightforward what is the error
- 2:49:28here missing return statement okay I
- 2:49:31need to return that map
- 2:49:35object fine now go to the consumer and
- 2:49:38do the same kind of configuration so go
- 2:49:41to the consumer
- 2:49:43app we already commented all the key and
- 2:49:46value from the application.
- 2:49:49now what I do I will create a new
- 2:49:51package then quickly create a
- 2:49:56class fine just annotate this at theate
- 2:50:01configuration now again here also I need
- 2:50:03to define a consumer config map by
- 2:50:06defining all the key and value key key
- 2:50:09dis realizer value dis realizer trusted
- 2:50:11packages bootstrap server group ID
- 2:50:14everything you can configure here okay
- 2:50:16so just go to the config
- 2:50:19class just Define a
- 2:50:25bin now you can create the object of a
- 2:50:29map now you can Define all the key and
- 2:50:31value so name it
- 2:50:34props then just
- 2:50:38Define all the value
- 2:50:42okay so we have defined Bop server key D
- 2:50:46serializer value d serializer and
- 2:50:48trusted package so IED this atate
- 2:50:52bin next you need to define a consumer
- 2:50:55Factory by giving this consumer config
- 2:50:58so quickly create that class
- 2:51:02public now return new
- 2:51:06default capka consumer Factory and
- 2:51:09Define the config which is nothing
- 2:51:11consumer
- 2:51:13config fine Define this at the
- 2:51:17bin now you just need to define the
- 2:51:20container Factory Capal listener
- 2:51:22container Factory okay so just Define
- 2:51:26that Capa listener container Factory and
- 2:51:30we are creating the object of it and we
- 2:51:31are giving the consumer Factory here
- 2:51:34okay this is the simple way we are
- 2:51:36creating the consumer config consumer
- 2:51:38Factory and consumer listener okay or
- 2:51:41Capa listener so the advantages of using
- 2:51:44this Java base config if your
- 2:51:46application is straightforward you don't
- 2:51:48want to do any customization then you
- 2:51:51can happily go with the application. yl
- 2:51:53file okay but for example let's say you
- 2:51:55want to access a secured Capa cluster
- 2:51:58which will be https okay then in that
- 2:52:01case you need to Define all the
- 2:52:02certificate and all in your producer
- 2:52:05config place okay you need to define the
- 2:52:08uh all the SSL related information in
- 2:52:10the config class you you need to
- 2:52:13manually configure them but that you
- 2:52:14cannot do using the application. yl file
- 2:52:17okay but if you check in the consumer
- 2:52:19there also if you want to set the
- 2:52:21concurrency label for your consumer that
- 2:52:24you can configure here in this listener
- 2:52:26container Factory that things cannot be
- 2:52:28done using the application. yml file so
- 2:52:31it's pretty simple if your approach is
- 2:52:33straightforward just to push the event
- 2:52:36to the topic and consume it without any
- 2:52:38customization you can go with the
- 2:52:40application. IML file but if you want
- 2:52:42any customization then you can go with
- 2:52:45the Java base config approach okay so
- 2:52:47all good now what I'll do I'll just
- 2:52:50restart both producer and
- 2:52:54consumer so it seems application is up
- 2:52:57and running you can see here the
- 2:52:58consumer this is the group assign to the
- 2:53:01three different partition let's see the
- 2:53:03producer yeah producer is up and running
- 2:53:06let me clear
- 2:53:08everything from the consumer clear
- 2:53:10everything now we comment out all the
- 2:53:14application. imlb configuration and
- 2:53:16rather we created our own Java config
- 2:53:19right so I'll go to the postman I'll
- 2:53:22change something 1 2 3 name
- 2:53:27John some number okay now let me send
- 2:53:30the
- 2:53:31request status code 20 go to the
- 2:53:36producer can you see here producer is
- 2:53:39able to publish the message this is what
- 2:53:41the messages now let's verify whether
- 2:53:44consumer is consumed it or
- 2:53:46not yeah can you see here consumer
- 2:53:50consume the event customer ID 1 2 3 name
- 2:53:52John email is this contact number is
- 2:53:55this okay so this is how you can play
- 2:53:58with the object in capka I mean you can
- 2:54:01serialize and deserialize any type of
- 2:54:03object in the form of Json using the
- 2:54:06springboard capka okay so I already
- 2:54:08explained both the approach application.
- 2:54:10IML and Java base config so just choose
- 2:54:13them based on your
- 2:54:17need
- 2:54:21so we'll start with brushing of capka
- 2:54:23internal workflow and then we'll move
- 2:54:25into the capka partitioning
- 2:54:27demonstration okay so in capka ecosystem
- 2:54:30once producer send bulk messages it will
- 2:54:33split into different partition okay so
- 2:54:36let's say I have sent thousand of
- 2:54:37messages then each and every message
- 2:54:40can't be guaranteed that it will go to
- 2:54:42the same partition it will go to the
- 2:54:44multiple partition the number of
- 2:54:46partition we have in our broker
- 2:54:48similarly consumer will consume from all
- 2:54:50the partition so this is the typical C
- 2:54:53cup of sub mechanism right now how you
- 2:54:56can ensure your message go exactly where
- 2:54:58you want them to go maybe you want to
- 2:55:01optimize the data processing or improve
- 2:55:03load balancing then how you can control
- 2:55:05that I want my producer will send all
- 2:55:09the messages to a single partition and
- 2:55:11say I want my consumer to read from a
- 2:55:14single partition how you can make the
- 2:55:17control on this partition
- 2:55:19that is what I'm going to demonstrate in
- 2:55:20this video so this is the capka producer
- 2:55:22example and here we have the capka
- 2:55:24consumer example this is what the
- 2:55:26example we are starting learning from
- 2:55:28the capka series beginning right I'm
- 2:55:31taking the same example to just
- 2:55:33demonstrate on this Capa partitioning
- 2:55:36okay so what we'll do first I will
- 2:55:38create a topic with couple of partition
- 2:55:41then first we'll see the general
- 2:55:43Behavior or the default Behavior how the
- 2:55:45messages are getting splitted to the
- 2:55:48different partition then we'll
- 2:55:50understand how we can make a control on
- 2:55:52that how I can send message to a
- 2:55:55specific partition rather than sending
- 2:55:56it to all the partition then we'll move
- 2:55:59to the consumer and we'll read from a
- 2:56:01specific partition okay so first let me
- 2:56:04create a
- 2:56:06topic so you know the command to create
- 2:56:08the topic right capka topic s then
- 2:56:11provide the topic name and then provide
- 2:56:13the number of partition and replication
- 2:56:16Factor so the topic name I'm creating
- 2:56:19here javat hyen topic just enter it so
- 2:56:23now this topic is created with five
- 2:56:26partition count okay can you see here
- 2:56:28the partition count is five it means now
- 2:56:31we can distribute our messages to
- 2:56:33different five partition fine so I will
- 2:56:36copy this topic name then I will go to
- 2:56:40the consumer I'll just change here then
- 2:56:44I'll go to the
- 2:56:45producer I will just change here
- 2:56:49here now simply start your
- 2:56:53producer and start your
- 2:56:56consumer but make sure before you play
- 2:56:59with your producer and consumer you
- 2:57:01should start your Juke keeper and capka
- 2:57:04server okay so I have already started
- 2:57:06both the server now I have started my
- 2:57:08producer and consumer then we'll publish
- 2:57:10the message and we'll verify meanwhile I
- 2:57:13will open the offset Explorer then I
- 2:57:15will refresh the topic so we created
- 2:57:18this topic right javat hyen topic I'll
- 2:57:21just go inside this
- 2:57:24topic there is no data because you have
- 2:57:27not published anything right so I cannot
- 2:57:29see anything here let's see if the
- 2:57:31producer and consumer is up so producer
- 2:57:34is up and
- 2:57:35running also consumer is up and running
- 2:57:38and the consumer group is assigned to
- 2:57:41the three partition of my topic sorry
- 2:57:44five partition of my topic can you see
- 2:57:46here 0 1 2 3 4
- 2:57:48that's fine just go to the producer and
- 2:57:52we'll just publish couple of messages so
- 2:57:55this is what the messages I mean this is
- 2:57:57what the method I want to execute who
- 2:57:59will use the template to send the
- 2:58:01messages now if I see here this is what
- 2:58:04the controller now here I'm sending
- 2:58:0610,000 messages so if I'll send bulk
- 2:58:09messages then only I can differentiate
- 2:58:11okay all the messages is going to
- 2:58:13different partition if I'll send one
- 2:58:16messages then simply it can go to the
- 2:58:18any single petition right so that is the
- 2:58:20reason I'm sending the bulk messages
- 2:58:22here and yeah I need to hit this
- 2:58:24particular Endo so let me go to the
- 2:58:27browser or I will go to the postman this
- 2:58:30is what the endpoint I'll will send some
- 2:58:33messages like
- 2:58:34welcome okay now let me send
- 2:58:38this so all the message has been
- 2:58:40published here you can see the messages
- 2:58:42now to verify that okay all the messages
- 2:58:45is goes to the top PE I will go to this
- 2:58:48offset Explorer and I'll simply refresh
- 2:58:52this then if I execute this I can see
- 2:58:56bunch of messages right and now let's
- 2:58:58see how many messages went to each
- 2:59:01partition so the number of partition we
- 2:59:03have is five can you see here 0 1 2 3 4
- 2:59:07I mean begin from the zero it is five
- 2:59:09now let's see the number of messages in
- 2:59:12each partition partition
- 2:59:14zero contains 2242 messages
- 2:59:20partition 1 contains
- 2:59:233,206 partition 2
- 2:59:261843 partition 3
- 2:59:28709 partition 4 2011 okay I mean we have
- 2:59:33sent 10,000 messages now different
- 2:59:36messages went to different partition and
- 2:59:39we can see the number of messages hold
- 2:59:41in each partition right so now that that
- 2:59:45is what the typical uh producer flow I
- 2:59:48mean uh whenever we are publish the
- 2:59:50messages messages will distributed to
- 2:59:52different partition but I want to make
- 2:59:54the control here I want to send messages
- 2:59:57to a single partition from this topic I
- 3:00:00might need to perform the data
- 3:00:03processing or optimization so I might
- 3:00:05need to send the messages to a single
- 3:00:07partition let's say I want to send to
- 3:00:09the partition three how I can do that
- 3:00:12that that is what we are going to learn
- 3:00:14here right so simple thing you need to
- 3:00:17tell to the capka template while sending
- 3:00:20the messages to which partition you want
- 3:00:23to put this messages that thing you need
- 3:00:25to tell to the capka template now how
- 3:00:27you can tell that simple go to your code
- 3:00:30go to the capka template where where
- 3:00:33exactly you are sending the messages now
- 3:00:35here if you'll open this send method
- 3:00:37this is the overloaded method in the
- 3:00:39template class can you see here topic
- 3:00:42and data topic key and data and here is
- 3:00:46the method top top now you can tell here
- 3:00:49which partition you want to send your
- 3:00:52messages you can specify the partition
- 3:00:54along with the key key is something
- 3:00:56which will help you to avoid duplicate
- 3:00:58message in capka that will cover in the
- 3:01:00separate session but this is what the
- 3:01:02method you can use to specify to which
- 3:01:05partition you want to send the messages
- 3:01:08so just do the simple code change go to
- 3:01:10the publisher and tell that okay I want
- 3:01:12to send to partition three and just pass
- 3:01:16the key I mean for now I will pass the
- 3:01:18the null I don't want to play with the
- 3:01:20key at this moment so now all the
- 3:01:22messages will send to this partition so
- 3:01:25shall we verify that simple thing just
- 3:01:27go to the controller I'll just reduce
- 3:01:29the message count I want to send 100
- 3:01:32messages okay now it will go to the
- 3:01:34partition three so just verify we want
- 3:01:38to send into this partition the current
- 3:01:40message count is 79 and we will push the
- 3:01:43100 more messages okay so simply I'll
- 3:01:46will just restart my product
- 3:01:48producer so producer is up and running
- 3:01:51now go to the postman I'll just send
- 3:01:54some different message let's say user
- 3:01:57and it will trigger 100 user messages
- 3:02:00that is what we are just looping in the
- 3:02:01code right we are sending 100 with the
- 3:02:04user count user 1 2 3 like 100 now let
- 3:02:08me send the
- 3:02:09request all the messages has been sent
- 3:02:13can you see here send message user 93
- 3:02:15and this are the partition and also so
- 3:02:18if you'll observe in the
- 3:02:19consumer okay what is the problem here
- 3:02:22okay so that's fine we'll we'll come to
- 3:02:24the consumer because we have Define your
- 3:02:26object and we are sending the string
- 3:02:28that is the reason this consumer is
- 3:02:30crying but we'll do the code change once
- 3:02:31we'll demonstrate the consumer part okay
- 3:02:34that's fine now just go to the offset
- 3:02:37Explorer see the count is 709 right and
- 3:02:40we have pushed 100 more messages so just
- 3:02:43let me refresh this then simply come to
- 3:02:46the partition three
- 3:02:48can you see here the count is 810
- 3:02:52earlier it was 709 right now all the 100
- 3:02:56messages is come to partition 3 only now
- 3:02:59similarly this is how you can manage in
- 3:03:01the producer part and you can specify to
- 3:03:04which partition you want to send the
- 3:03:06messages like this you just need to
- 3:03:10specify the partition number that's it
- 3:03:13now this is what we understand from the
- 3:03:15producer perspective now whatever all
- 3:03:18the messages we are sending to a
- 3:03:19partition or the broker is being
- 3:03:22consumed by my consumer right now I
- 3:03:25don't want my consumer to look into all
- 3:03:28the partition I want to tell to my
- 3:03:30consumer hey consumer can you please
- 3:03:33read into a specific partition for
- 3:03:35example from partition 3 or partition
- 3:03:37two not from all how I can do that we
- 3:03:40understand from the producer perspective
- 3:03:42how to push the message to a specific
- 3:03:44partition now we we just need to start
- 3:03:47implementing from the consumer
- 3:03:49perspective how can my consumer read
- 3:03:51from a specific partition okay so it's
- 3:03:55very simple you just need to play with
- 3:03:57The annotation just go to the
- 3:04:00consumer okay let me stop this first
- 3:04:04then I will just do one thing I'll just
- 3:04:06copy this or I can create different
- 3:04:09instance okay I'll take the input as a
- 3:04:12string because that is what I'm sending
- 3:04:14from my producer so it found the object
- 3:04:17which is is not able to dis realize so
- 3:04:19that is the reason we are getting the
- 3:04:20error in the consumer that is fine now
- 3:04:23here in The
- 3:04:25Listener here exactly in the capka
- 3:04:27listener you need to specify from which
- 3:04:32partition you want to execute this
- 3:04:34consumer to consume the messages that is
- 3:04:37what you can just simply Define by just
- 3:04:39writing the topic
- 3:04:43partition okay now in this topic
- 3:04:45partition you can simply just Define
- 3:04:48here at
- 3:04:50theate topic
- 3:04:52partition now you just need to define
- 3:04:54the what is the topic name and what is
- 3:04:56the partition count you want to read the
- 3:04:59topic name you can Define the same topic
- 3:05:01name from where you want to
- 3:05:03consume and then you can just Define the
- 3:05:07partition count the count I want to read
- 3:05:10it from let's say from partition two
- 3:05:13okay I want to read it from partition
- 3:05:16two that's it so this is how you can
- 3:05:20control from your consumer perspective
- 3:05:23to read from a specific partition okay
- 3:05:26that's fine now let's verify that so now
- 3:05:29I'll just use the different topic javate
- 3:05:31techy topic one fine and then okay let
- 3:05:36it be I mean better I will just comment
- 3:05:38this okay now also go to the
- 3:05:42producer and I'll will just change your
- 3:05:44Java topic one rather than doing like
- 3:05:47this what what I will do I just want to
- 3:05:49send messages to the different partition
- 3:05:53then I can prove that okay my consumer
- 3:05:55is reading from the specific partition
- 3:05:57what I have mentioned there so for that
- 3:05:59what I can do capka template or what I
- 3:06:02have different template right template
- 3:06:04do
- 3:06:06send messages to I mean I'll just simply
- 3:06:09copy this
- 3:06:11right so what is the
- 3:06:13messages I'll just Define let's say hi
- 3:06:17then something some random string
- 3:06:20okay I'm sending multiple messages to
- 3:06:22the same topic to different partition
- 3:06:25one
- 3:06:27two okay first let me see what number we
- 3:06:30have specified two right so I'll send
- 3:06:33more messages to the partition
- 3:06:352 so I'm sending to all the partition
- 3:06:39but I want my consumer to read from only
- 3:06:42these two partition okay I mean this is
- 3:06:45the same partition my consumer should
- 3:06:47receive only these two messages because
- 3:06:49these two messages went to partition two
- 3:06:52so just change the
- 3:06:55messages okay so welcome and YouTube
- 3:06:58should my consumer consumed I mean that
- 3:07:01is what we can find in my consumer
- 3:07:03console let's verify that so looks good
- 3:07:08we have specified the topic name
- 3:07:09correctly now we just need to create
- 3:07:11this topic with five partition okay so
- 3:07:15go to the terminal just
- 3:07:18run this the topic name the same Java
- 3:07:21topic one fine the topic is created
- 3:07:24let's verify that topic is
- 3:07:27created yeah topic is here we don't have
- 3:07:30any messages it contains five partition
- 3:07:33fine so what I'll do I'll just start my
- 3:07:38producer and my consumer as
- 3:07:42well so both are up and running let's
- 3:07:45verify the producer all good let's go to
- 3:07:48the
- 3:07:49consumer in consumer If You observe the
- 3:07:52console statement consumer client ID
- 3:07:55this is are the consumer group group ID
- 3:07:58is this resetting offset for partition
- 3:08:01javat topic one 2 can you see here this
- 3:08:06is the topic name Java hen topic one and
- 3:08:09it is only listening to partition two
- 3:08:12that is what you can see in the
- 3:08:14statement itself phas position offset
- 3:08:17from the zero I mean it's just telling
- 3:08:19that okay read it from the partition two
- 3:08:21that is what I can understand from this
- 3:08:23console now let's hit the endpoint and
- 3:08:26we'll
- 3:08:27verify I'll send the some some messages
- 3:08:30okay because anyway it is not going to
- 3:08:32print we have hardcoded the value while
- 3:08:34sending the messages send the
- 3:08:37request message has been sent now if you
- 3:08:39will go and check in the
- 3:08:43topic just open the
- 3:08:45data white send lot of
- 3:08:48messages okay oh it's sending on the
- 3:08:51loop yeah that's fine but let's see the
- 3:08:54number of messages in the partition 2 20
- 3:08:58202 now let's see the consumer go to the
- 3:09:02consumer see consumer only consume the
- 3:09:05method I mean messages welcome and
- 3:09:07YouTube because it's in the loop so it's
- 3:09:10send multiple I mean 100 messages so can
- 3:09:13you see here it only receive welcome and
- 3:09:16YouTube
- 3:09:18because only YouTube and welcome these
- 3:09:20two messages we are sending to partition
- 3:09:242 in our producer can you see here
- 3:09:26welcome and YouTube so it is clearly
- 3:09:29state that my consumer is listening to
- 3:09:32only partition 2 because I cannot see
- 3:09:35this hi hello Java TI key these messages
- 3:09:37in my consumer console okay so you
- 3:09:41cannot let me filter it okay we have not
- 3:09:44see this no
- 3:09:48no right but only welcome and
- 3:09:52YouTube can you see here the count is
- 3:09:54itself is 101 and
- 3:09:57YouTube it should be same one1 okay
- 3:10:01that's it I mean this is how you can
- 3:10:03play with the Capa partition to specify
- 3:10:06while publishing the messages and
- 3:10:08specify the partition while consuming
- 3:10:11the messages this is sometimes required
- 3:10:13Whenever there is some uh message event
- 3:10:16failed in the production you might need
- 3:10:18to republish the event okay and that
- 3:10:21time you might need to play with this
- 3:10:22partition kind of concept so just give a
- 3:10:25try and let me know in a comment section
- 3:10:27if you have any
- 3:10:33doubts if you remember in the last capka
- 3:10:36session we have designed our capka
- 3:10:38producer and consumer application right
- 3:10:41so in this tutorial I'll guide you how
- 3:10:43to write capka integration test with
- 3:10:46test container
- 3:10:48okay all right being a developer always
- 3:10:50writing a capka integration test can be
- 3:10:53frustrating because that's mainly due to
- 3:10:55the complex test configuration that
- 3:10:57involves for registering the consumer
- 3:10:59and producer or to read and write
- 3:11:02messages isn't it and also without
- 3:11:04proper integration testing you cannot be
- 3:11:06confident about the stability of your
- 3:11:08production environment So to liage that
- 3:11:11we can use the test containers which can
- 3:11:13reduce test configuration to almost zero
- 3:11:16and not worry about how to read and
- 3:11:18write messages from and to capka when
- 3:11:21writing a test case instead a developer
- 3:11:24can entirely focus on only testing
- 3:11:26needed functionality okay so if you're
- 3:11:29not aware about what is test containers
- 3:11:31and how to use it then I would strongly
- 3:11:33suggest you to check out this particular
- 3:11:35video in my YouTube channel spring boot
- 3:11:373 integration testing with test
- 3:11:40container where I have explain in detail
- 3:11:42by taking the example of my SQL using
- 3:11:44the test containers I will also share
- 3:11:46the link in the video description for
- 3:11:48your reference okay so without any
- 3:11:50further delay let's Circle back to the
- 3:11:52demo so let's get started so before we
- 3:11:55start writing the test case first let's
- 3:11:58verify whether our producer and consumer
- 3:12:00example application is working or not so
- 3:12:02I have started my producer application
- 3:12:05and also I have started my consumer
- 3:12:07application here and if You observe here
- 3:12:10in the terminal I have started my juke
- 3:12:12keeper and capka server okay that's
- 3:12:15enough now if you will go to the
- 3:12:17producer app in the producer application
- 3:12:20we have defined the controller from this
- 3:12:22controller method we are sending the
- 3:12:24customer object to the topic fine so the
- 3:12:27endpoint is SL publish let's go to the
- 3:12:30postman and I have this particular
- 3:12:32endpoint with this particular payload
- 3:12:34now just trigger the
- 3:12:36request status code is 20 now let's
- 3:12:39verify in the consumer whether that
- 3:12:41particular object is received or not
- 3:12:43yeah can you see here consumer consume
- 3:12:46the event this this is what we are
- 3:12:48sending from our producer one1 name is
- 3:12:51basant email ID and contact number so
- 3:12:55all good now let's start writing the
- 3:12:57test case for that I don't want to use
- 3:13:00my capka which I installed on my machine
- 3:13:03so I will stop
- 3:13:04everything I will stop my
- 3:13:07consumer and producer as
- 3:13:10well and from the terminal I'll close my
- 3:13:13juke keeper and capka server because I'm
- 3:13:16not going to use this particular capka
- 3:13:18instance I want to take the help from
- 3:13:20capka taste containers to spin up a new
- 3:13:23capka instance for me so I'll show you
- 3:13:26that how you can do that so let me close
- 3:13:28everything now the local Capa is shut
- 3:13:30down for me so what I'll do very simple
- 3:13:34just go to the pal.
- 3:13:36XML in pal. XML we just need to add the
- 3:13:40tast container specific dependency if
- 3:13:42I'll open it so I have added the tast
- 3:13:45container dependency then taste
- 3:13:47container for capka and I have added the
- 3:13:50Jupiter and also I have added this
- 3:13:52utility to just pause for some time I
- 3:13:56mean whenever I'm running the capka
- 3:13:58producer and consumer I want my thread
- 3:14:01will pause for some specific second or
- 3:14:04for some time interval that is the
- 3:14:05reason I I have just added this utility
- 3:14:08okay I I'll tell you in the code
- 3:14:10whenever we will implement this that's
- 3:14:12fine then just update your
- 3:14:15project then what I want to do I want to
- 3:14:18disable this config class because I want
- 3:14:20to load it from the application. yml
- 3:14:23file so I'll just uncomment this that's
- 3:14:27it just go to the test class and start
- 3:14:30writing your test case now here okay let
- 3:14:33me close everything so first step we can
- 3:14:37just run this particular test case in
- 3:14:40random Port okay so I'll Define web
- 3:14:44environment there is something called
- 3:14:47random
- 3:14:48Port then I just let me Zoom this for
- 3:14:51you then I just need to annotate here at
- 3:14:54theate test
- 3:14:56containers by adding this annotation I'm
- 3:14:59telling to the spring boot test that
- 3:15:01this particular class is using some
- 3:15:04container fine now which container I
- 3:15:08want to use I want to use the capka
- 3:15:10container so for capka container how
- 3:15:12from where you'll get this capka
- 3:15:14container just go to this documentation
- 3:15:17you can copy this and then then simply
- 3:15:21just paste it here we are using the
- 3:15:23capka container by saying to The capka
- 3:15:25Container hey capka container can you
- 3:15:28please bring this specific capka version
- 3:15:30for me
- 3:15:316.2.1 if you don't want to specify any
- 3:15:33capka version just tell that hey
- 3:15:36container can you give me the latest
- 3:15:38capka version which you have with you
- 3:15:40that's fine you can Define the version
- 3:15:42whatever you want here okay now annotate
- 3:15:45this atate container
- 3:15:48next since we have created the Capa
- 3:15:50container we need to configure the
- 3:15:52bootstrap server so for that I will
- 3:15:54simply write a method public
- 3:15:57viid let's say init cap Properties or
- 3:16:00something like
- 3:16:02that then you need to just pass your
- 3:16:04Dynamic Property
- 3:16:07registry and also here you need to
- 3:16:09annotate Dynamic Property Source these
- 3:16:13annotation came from the this specific
- 3:16:16test containers okay so you don't need
- 3:16:18to remember this annotation name but
- 3:16:20whenever you are using the test
- 3:16:21container you need to play with these
- 3:16:24annotation fine now
- 3:16:26registry do
- 3:16:29add key name is spring. ca. bootstrap
- 3:16:33hypen server I just need to copy the
- 3:16:35same name what I have defined here so
- 3:16:38this is
- 3:16:39right so you can Define spring.
- 3:16:44capka then bootstrap servers
- 3:16:47and what is the bootstrap server you
- 3:16:49want to specify I don't know exactly the
- 3:16:51bootstrap server name because this
- 3:16:54particular container will get the server
- 3:16:56for me every time you will run you'll
- 3:16:58get different server okay so what we can
- 3:17:01do we can tell here hey get it from the
- 3:17:04Capa container
- 3:17:06itself this is how we can
- 3:17:09specify fine so we are done with the
- 3:17:13infrastructure setup now let's start
- 3:17:15writing our actual test case so let's go
- 3:17:19to the project we are going to write the
- 3:17:21test case for this guy okay this capka
- 3:17:24message publisher we want to send the
- 3:17:26custom object which is customer this
- 3:17:29particular method so let me copy this
- 3:17:31method
- 3:17:33name then I'll simply write the test
- 3:17:35method
- 3:17:37public
- 3:17:38viid
- 3:17:41test fine so to call this specific
- 3:17:44method from the capka message publish
- 3:17:47sir first I need to inject this so just
- 3:17:50go here and I just need to Auto add that
- 3:17:56private then just do the auto
- 3:17:59add that's it now I just need to call
- 3:18:04Publisher dot send events to the topic
- 3:18:09there are two method we have defined in
- 3:18:10that publisher one will send the string
- 3:18:13message one will send the custom object
- 3:18:15where we have just find the customer as
- 3:18:17a d and we are sending it so I'm using
- 3:18:20the second method I mean wherever I am
- 3:18:22sending the custom object so I will
- 3:18:24create a new customer new customer what
- 3:18:28all field I need to pass ID name email
- 3:18:31and contact number so let's say ID
- 3:18:34something some random
- 3:18:37number name let's say test
- 3:18:40user email I guess email is test at
- 3:18:45gmail.com
- 3:18:47then what else contact number some
- 3:18:50random number
- 3:18:52okay next just annotate this at theate
- 3:18:55test make sure to import the test from
- 3:18:59Jupiter not from the junit once you have
- 3:19:02written this this sent event to the
- 3:19:04topic will publish the message to the
- 3:19:06topic and after that I want to pause for
- 3:19:09some time interval so that is the reason
- 3:19:11we have used this utility class so let
- 3:19:13me add the input
- 3:19:15statement this should be UT from org do
- 3:19:18this particular class okay and I want to
- 3:19:21wait the pool interval is 3 second and
- 3:19:25minimum I want to wait for 10
- 3:19:27second and next to that if you have any
- 3:19:30asset statement you can add it inside
- 3:19:33this particular
- 3:19:35block so just UT
- 3:19:38this so we are specifying the pool
- 3:19:41interval of 3 second and at most I mean
- 3:19:44the minimum amount I want to wait for 10
- 3:19:46second once this particular message will
- 3:19:48push to the topic so after that if you
- 3:19:51have any asset statement you can write
- 3:19:53it here since my method is B I don't
- 3:19:56have any Asser statement I cannot write
- 3:19:58anything here even you can give a try by
- 3:20:01persisting that particular customer
- 3:20:03object to the DV and you can just write
- 3:20:05the assert statement here by fetching it
- 3:20:08from the
- 3:20:09DV so this is the simple example guys
- 3:20:11okay nothing to worry now I believe all
- 3:20:14good let's run our code
- 3:20:17so let me run the test
- 3:20:20case okay it seems there is error this
- 3:20:24particular method must be static okay
- 3:20:26let's change it to the
- 3:20:28static so if you'll make this static
- 3:20:30better let's change this particular
- 3:20:32field also
- 3:20:34static now looks good let me run this
- 3:20:38we'll find out some error then I'll tell
- 3:20:40you how we can resolve
- 3:20:42that so if You observe here when the
- 3:20:45particular test will execute first it
- 3:20:48will get that particular capka image can
- 3:20:51you see here starting to pull the image
- 3:20:54what is the version latest now it will
- 3:20:56take few second to pull that Docker
- 3:20:59images so you can see here it
- 3:21:02successfully pull the image of this
- 3:21:04specific version I mean the latest
- 3:21:06version then it's trying to start that
- 3:21:08particular container this is what right
- 3:21:11now to verify whether it pull the
- 3:21:13specific Docker image or not what you
- 3:21:15can do go to the terminal just type
- 3:21:19Ducker
- 3:21:21images If You observe here can you see
- 3:21:23here C capka and the tag is latest now
- 3:21:27if you want to verify whether that is
- 3:21:29the container itself is running or not
- 3:21:31just run Docker PS you can see here
- 3:21:35right the capka container is running
- 3:21:37here let me Zoom this for you the
- 3:21:41container is running here and this is
- 3:21:43the Ducker image it pull
- 3:21:47kapka fine so go back here let's see
- 3:21:51what is the result so if You observe
- 3:21:54here it's giving us the
- 3:21:57error because this particular test case
- 3:22:01will look into your application. yml
- 3:22:03file which you have defined in
- 3:22:06your Source folder I mean in the main
- 3:22:08folder not in the test folder there you
- 3:22:11have specified the bootstrap server as a
- 3:22:149092 local host but our bootstrap server
- 3:22:18is running dynamically fine so for that
- 3:22:21reason you need to override this
- 3:22:23application. yml in your test
- 3:22:27folder so just create a directory
- 3:22:32resource just paste
- 3:22:35this and remove the bootstrap server so
- 3:22:38that it will pick from the container and
- 3:22:40will run it Dynamic bootstrap I mean the
- 3:22:43bootstrap server will be dynamic here so
- 3:22:46this is not not required for test you
- 3:22:47can remove
- 3:22:49this fine now go to the test class and
- 3:22:53run it
- 3:22:55again now second time when I'm running
- 3:22:58it it won't pull the image again because
- 3:23:01it is already pulled that specific image
- 3:23:03at first time fine so let's wait it to
- 3:23:08complete see directly pull the image
- 3:23:11here and starting the
- 3:23:15container okay can you see here we got
- 3:23:18the message here sent message this is
- 3:23:21what the object we are sending right
- 3:23:23with upset zero now if you go and check
- 3:23:25in the publisher class after
- 3:23:28successfully send the messages we are
- 3:23:31just printing the statement
- 3:23:33system.out.print and send message this
- 3:23:35is what right if there is a exception
- 3:23:38then it will print this particular
- 3:23:40statement since in our case we are able
- 3:23:42to successfully send the messages we are
- 3:23:44able to get the first statement that
- 3:23:46sent message to the particular message
- 3:23:49and to which upset that is what we can
- 3:23:52see here
- 3:23:53right if you go here to the test
- 3:23:57case get the result go to this specific
- 3:24:00test case now if you'll scroll down you
- 3:24:02can see the log messages here right this
- 3:24:06is what we have written for producer now
- 3:24:08since I am a producer I might have 100
- 3:24:11of consumers right so now consumers
- 3:24:14people need to write their own test case
- 3:24:16for Consumer class okay so let's start
- 3:24:19begin writing the test case for Consumer
- 3:24:21we are good with the producer now if You
- 3:24:23observe in the code I have not worried
- 3:24:26more about the infrastructure I have
- 3:24:28just defined these two line everything
- 3:24:30will be Take Care by this capka
- 3:24:31container itself what I focused I only
- 3:24:34focused on writing my functionality test
- 3:24:37whether this particular method where I'm
- 3:24:39exactly publishing the event is working
- 3:24:42or not okay now let's start writing the
- 3:24:44test case for consumer application so I
- 3:24:47need this information so first of all
- 3:24:49first let me copy all the dependency
- 3:24:52what I have added here because same
- 3:24:54dependency we need to add in our
- 3:24:56consumer as
- 3:24:57well copy this go to the consumer
- 3:25:01application go to the pom.xml
- 3:25:05then just scroll down and paste it
- 3:25:10here just update
- 3:25:12it fine now go to the test case here
- 3:25:16this are the test case now let me copy
- 3:25:19all the field from the producer I mean
- 3:25:21we are going to reuse the same
- 3:25:23annotation and um capka container so
- 3:25:25I'll just copy
- 3:25:27paste just go to the
- 3:25:30consumer paste it
- 3:25:33here go to the producer get the
- 3:25:37container and Dynamic Property
- 3:25:39Source copy
- 3:25:41this then paste it here I mean whoever
- 3:25:45have the consumer for your application
- 3:25:48they need to test their application that
- 3:25:50whether they are able to successfully
- 3:25:52face the events from that specific topic
- 3:25:54or not okay so that is the reason they
- 3:25:57need to write the test case so for that
- 3:25:59how they can validate that whether the
- 3:26:01message is coming to the topic or not so
- 3:26:03forcefully you need to send some event
- 3:26:06or messages to the topic then in
- 3:26:08consumer itself you need to write the
- 3:26:09logic to the or you need to validate
- 3:26:12that whether it is consuming or not okay
- 3:26:14that is simple statement guys so let me
- 3:26:17write the test case
- 3:26:19public first let me go to The Listener
- 3:26:22class right this is the class where my
- 3:26:24consumer is written here I mean this is
- 3:26:27the place I'm consuming the event from
- 3:26:29this specific topic so that is the
- 3:26:31reason I have defined the adate Capal
- 3:26:34listener so I'll just take the method
- 3:26:36name go to the test class public
- 3:26:40void
- 3:26:43test so to validate whether this part
- 3:26:46particular method is being consumed or
- 3:26:48not so what I'll do forcefully I need to
- 3:26:52send some message to this topic so for
- 3:26:54that what I need I need a capka template
- 3:26:57so I just need to inject here Capa
- 3:27:02template up type string as a key and
- 3:27:06value as a
- 3:27:08object just inject using Auto
- 3:27:13add so here what I will do I'll just
- 3:27:16send Capa template do
- 3:27:19send object to which topic the same
- 3:27:23topic what you have defined here because
- 3:27:25I only forcefully send some message to
- 3:27:26this topic to validate okay this
- 3:27:29consumer is working
- 3:27:30correctly just specify the topic and
- 3:27:33just Define the customer object so I
- 3:27:35will better copy the customer object
- 3:27:37from here
- 3:27:40itself just paste it here or I'll just
- 3:27:44Define above okay
- 3:27:47then just pass this customer object to
- 3:27:49this Capa template so that it will
- 3:27:52publish the message to this topic and
- 3:27:54this particular listener need to execute
- 3:27:56and it will print this log statement
- 3:27:58that is what we are expecting fine so
- 3:28:01what I do go to the test first I'll do
- 3:28:04one thing I'll just add some log
- 3:28:06statement here SL
- 3:28:084J then I'll just
- 3:28:11add
- 3:28:12log
- 3:28:15doino let's say test consume
- 3:28:19events method
- 3:28:22execution started something like that
- 3:28:25and also I will just add execution
- 3:28:29ended so execution will be start then we
- 3:28:32are building the customer object and
- 3:28:34through capka template we are sending it
- 3:28:36to the this specific topic fine now once
- 3:28:40it will publish to to this specific
- 3:28:42topic this listener will read that
- 3:28:44messages and it will print the log stat
- 3:28:46once everything done we are just
- 3:28:48printing that okay the execution is
- 3:28:50ended and you can also annotate your
- 3:28:54test and make sure to input it from the
- 3:28:57Jupiter okay this is as much as
- 3:29:01simple fine now here also what I want to
- 3:29:04do I want to enable the application.
- 3:29:08yml I mean I just want to uncomment I
- 3:29:11want to use this application. IML and I
- 3:29:13want to comment this config class again
- 3:29:17you can try in both the way but let's go
- 3:29:19with the application. IML so that if any
- 3:29:22test configuration required I can add
- 3:29:24directly the yml file here rather than
- 3:29:26Define the config class now here also I
- 3:29:29need to define the resource folder new
- 3:29:33directory resources and just add the
- 3:29:36application. IML
- 3:29:38file fine so make sure you need to
- 3:29:42define the same group in your test case
- 3:29:45so so let's go here I mean in not in the
- 3:29:48test case this particular capka template
- 3:29:51will publish the message to this topic
- 3:29:53now if you go and check in The Listener
- 3:29:55the topic name is correct and the group
- 3:29:57name need to be this okay I mean it is
- 3:30:01same only that is how I have configured
- 3:30:02in yml file so go to the test case all
- 3:30:07good now let me run
- 3:30:13this so it started the container here
- 3:30:15can you see see here now it will start
- 3:30:17executing our test
- 3:30:20case so this test is executed
- 3:30:24successfully now if you filter
- 3:30:27this we have not received the I mean the
- 3:30:29request is not goes to my consumer fine
- 3:30:33if you'll see here not here to go to the
- 3:30:36main one you'll find some error right or
- 3:30:40ca.on producer is closed forcefully so
- 3:30:43for that what you can do you need to
- 3:30:46wait for some time I mean the way we
- 3:30:49have done in the producer similarly you
- 3:30:51need to wait for some second because it
- 3:30:54might take some time to establish the
- 3:30:55connection right or it might take some
- 3:30:58time to connect to the container so it's
- 3:31:00good to keep some buffer uh time
- 3:31:02interval for
- 3:31:04sleep fine now let's run it just run
- 3:31:09from
- 3:31:11here restarting the container
- 3:31:14now
- 3:31:16so it started our
- 3:31:18application now can you see here let's
- 3:31:21wait it to complete yeah within a second
- 3:31:23it will be complete now if you'll scroll
- 3:31:26down the JT group partition assigned to
- 3:31:29this particular topic partition I mean
- 3:31:31the partition count is zero and this is
- 3:31:33your topic name and this is the group
- 3:31:36name and can you see here the messages
- 3:31:39let me Zoom this consumer consume the
- 3:31:42event this is what the 263 test user t
- 3:31:46start the email and contact number this
- 3:31:48is the object we are publishing through
- 3:31:50the capka template and we have defined
- 3:31:52the consumer who listen to that
- 3:31:54particular topic and my consumer is able
- 3:31:57to consume it so if You observe here go
- 3:32:00to the consumer the statement we have
- 3:32:02written is same right consumer consume
- 3:32:05the event then what is the event it is
- 3:32:07getting we're just converting it to the
- 3:32:09string the same statement we are able to
- 3:32:12see it I mean we are able to
- 3:32:15successfully consume the event from our
- 3:32:17test case and using the test containers
- 3:32:21so if You observe here again I didn't
- 3:32:23take much heck to Define my
- 3:32:25infrastructure I just Define these two
- 3:32:27statement and rest everything will be
- 3:32:29Take Care by this capka container itself
- 3:32:32I mean the tast container itself okay so
- 3:32:35the another advantages here in this
- 3:32:37particular approach the Capa Ducker
- 3:32:39container can be kept alive which can
- 3:32:42help with the test performance and
- 3:32:44results in a testing environment closer
- 3:32:46to the production one so this particular
- 3:32:49tutorial will give you the fully context
- 3:32:51how you can write the integration test
- 3:32:54for capka using the tast containers
- 3:32:56whenever you'll try to apply this
- 3:32:58specific use case to your real life
- 3:33:00situation you can take it further by
- 3:33:02adding more test cases you could also
- 3:33:04start refactoring your app or adding new
- 3:33:07features more reliable with the test
- 3:33:09that you have written
- 3:33:14okay
- 3:33:18as you know Apachi kapka run in a
- 3:33:21distributed manner across multiple
- 3:33:23containers or machines it is also
- 3:33:26crucial to address the errors and
- 3:33:28recover the data in the event of failure
- 3:33:32sometimes what happen when a producer
- 3:33:34send a messages and we process that
- 3:33:37message from the capka topic error can
- 3:33:39occur for instance consumer services or
- 3:33:42related infrastructure such as database
- 3:33:45connection
- 3:33:46might be unavailable during that period
- 3:33:49so in that case there is a risk of
- 3:33:52losing any events send or received due
- 3:33:55to the failure so here what do we want
- 3:33:57to ensure we don't lose any data and try
- 3:34:01to handle the failed messages that's the
- 3:34:04challenging part isn't it but don't
- 3:34:06worry in this tutorial I'll be
- 3:34:09explaining different way to handle
- 3:34:11errors in capka okay all right let me
- 3:34:15walk you through an use case consider a
- 3:34:17scenario where you are processing
- 3:34:19financial transaction if a transaction
- 3:34:22fails to be processed due to some
- 3:34:24temporary issue then how will you handle
- 3:34:27it because here we want to ensure
- 3:34:29reliable message processing so what
- 3:34:32we'll do here we'll instruct capka to
- 3:34:36retry the failed event now what capka
- 3:34:39will does here he will simply reattempt
- 3:34:42process the messages multiple times as
- 3:34:46per your configuration for example if
- 3:34:48you set the retry count to four capka
- 3:34:51will make three retry attempts in a
- 3:34:54sequential order because it follows n
- 3:34:57minus one order okay so if you specify
- 3:35:00four times to retry capka will do three
- 3:35:03times so that is how you need to decide
- 3:35:05what number you need to set for retry
- 3:35:08attempt next if the process succeed
- 3:35:11Within These retrace that's excellent
- 3:35:14however if the retray account is
- 3:35:16exceeded then the message will be
- 3:35:18directed to the dead letter topic okay
- 3:35:21or the messages will be route to the DLT
- 3:35:24now you might have a question hey what
- 3:35:27is DT DT stands for dead letter topic
- 3:35:31which is nothing but creating a new
- 3:35:33failure topic to store all the failed
- 3:35:36messages or all the failed events if a
- 3:35:38consumer is not able to process messages
- 3:35:41even after retry then those unprocessed
- 3:35:44messages will p push to this particular
- 3:35:47dead letter topic okay now since you
- 3:35:49have failed records with you in a topic
- 3:35:52subsequently you can monitor them to
- 3:35:54identify and investigate further so this
- 3:35:58mechanism ensures that no data is lost
- 3:36:01in the event of failure we still rain
- 3:36:03the ability to handle and manage those
- 3:36:06failed messages right so this is the
- 3:36:09simple way to ensure that there is no
- 3:36:12data is lost and reliable message
- 3:36:15process procing so let's quickly see
- 3:36:17this in a action let's go to our
- 3:36:19intelligent idea since already you
- 3:36:21discuss how to produce and consume
- 3:36:24messages using capka so I have not
- 3:36:26wasting the time and I have created a
- 3:36:28project already so if you check here I
- 3:36:31have a publisher and I have a consumer
- 3:36:33if you're not aware about how to play
- 3:36:35with the popup using capka you can check
- 3:36:38out my playlist capka for beginers from
- 3:36:41my YouTube I will also share the link in
- 3:36:43video description for your refer friend
- 3:36:46okay now let's jump into the code and
- 3:36:48try to understand what pops up we are
- 3:36:50doing in this particular example let's
- 3:36:53go to the
- 3:36:55publisher publisher is the simple thing
- 3:36:57what we are doing here we are taking the
- 3:36:59user object as a input and we are
- 3:37:02publishing it to the
- 3:37:03capka fine now if you'll open this
- 3:37:06particular user object we have couple of
- 3:37:08field called ID first name last name
- 3:37:11email gender and IP address now let's go
- 3:37:15to the consumer what exact action we are
- 3:37:18doing in the
- 3:37:19consumer if I open the
- 3:37:23consumer we are simply taking that input
- 3:37:26user input what we received from our
- 3:37:29topic then I'm just adding one
- 3:37:31validation here to simulate the error
- 3:37:34scenario I have just added this
- 3:37:36restriction okay so as part of the user
- 3:37:40you can see there is a IP address so
- 3:37:42when I will receive any user request or
- 3:37:44when I will consum
- 3:37:46any event specific to the user I will
- 3:37:48validate that okay if user is giving any
- 3:37:52of the IP specified in this particular
- 3:37:54list then I will not allow that user to
- 3:37:57process further okay so I have just
- 3:38:00added a simple if Clause to check if the
- 3:38:04address if the list contains the address
- 3:38:07then throw some simple exception to just
- 3:38:09simulate the error scenario okay now
- 3:38:13what I'll do let's simply wrun on this
- 3:38:15particular use case then we'll find out
- 3:38:18what is the problem with this particular
- 3:38:20implementation then we'll go with the
- 3:38:22solution so first I need to start the
- 3:38:25jke
- 3:38:27keeper start
- 3:38:30it then let me simply start the capka
- 3:38:34server just enter
- 3:38:37it so here my Joe keeper and capka
- 3:38:41server both are up and running to
- 3:38:43validate that what I'll do I'll go to
- 3:38:45this
- 3:38:47upset yeah if I'll open this I can able
- 3:38:50to make the
- 3:38:53connection okay great so this is the
- 3:38:56tool to visualize your cup cup payload
- 3:38:59from the producer end and you can find
- 3:39:01the consumer related stop as well okay
- 3:39:04everything I have covered in my Capa
- 3:39:05playlist you can take a look to that
- 3:39:07particular
- 3:39:09videos all good now let's start our
- 3:39:13application so in config we have defined
- 3:39:16to create the topic can you see here we
- 3:39:20have the topic name and here we are
- 3:39:22saying to create the topic with three
- 3:39:25partition and the topic name is defined
- 3:39:29in this application. properties
- 3:39:33file
- 3:39:34fine so let's see the
- 3:39:38log can you see here Java group that is
- 3:39:42what my consumer group partition
- 3:39:44assigned two since we have created topic
- 3:39:47with the three partition you can see
- 3:39:49each and every partition here right 1
- 3:39:53two so now just let's test our pops up
- 3:39:57flow so we have defined a controller
- 3:39:59class here Here If You observe we have a
- 3:40:02endpoint called publish new where we are
- 3:40:05just giving a user object and it will
- 3:40:08call our publisher to publish the event
- 3:40:11to the topic okay then our consumer will
- 3:40:14listen to that part particular topic and
- 3:40:16we'll do the steps mention here that's
- 3:40:20fine just go to the
- 3:40:22postman I'll give some valid input
- 3:40:25first so let me send the
- 3:40:29request message published successfully
- 3:40:32now if you validate in your
- 3:40:34console Can You observe here this is the
- 3:40:38log from producer okay the messages is
- 3:40:42pushed to upset zero and this is the log
- 3:40:46from the
- 3:40:47consumer can you see
- 3:40:51here 2
- 3:40:5435 let's change the name to basant then
- 3:40:58last name something
- 3:41:00Hut and give some random email
- 3:41:04id id will be two now let me send the
- 3:41:09request the message published
- 3:41:11successfully again if you check in your
- 3:41:14log
- 3:41:15the second record with id2 from the
- 3:41:18producer log and this is what the
- 3:41:21messages log from the
- 3:41:23consumer okay it read from the upset one
- 3:41:28and also the upset one in the producer
- 3:41:31this is the happy scenario right it is
- 3:41:33working as expected now let me clear
- 3:41:35this let me give some IP address which
- 3:41:39is specified
- 3:41:41here let me
- 3:41:43just change change in
- 3:41:46the request
- 3:41:48body now if I'll send the request let's
- 3:41:51see what will be the
- 3:41:54result message published successfully
- 3:41:56producer is successfully able to publish
- 3:41:59the message to the topic now the
- 3:42:01validation is there in our consumer
- 3:42:03right so what is the result here invalid
- 3:42:07IP address received that is what we are
- 3:42:10expecting here right we are expecting to
- 3:42:13get the exception and we are getting it
- 3:42:15now what is the problem
- 3:42:18here there is nothing right what is the
- 3:42:20problem here the simple thing is that
- 3:42:23the record which we are not processed if
- 3:42:27we not take a look into it there could
- 3:42:29be a chance for data lost right the
- 3:42:31record which we just received and what
- 3:42:34just we validated that we need to take
- 3:42:36care of it we need to reattempt because
- 3:42:39see here I am just added the if
- 3:42:42statement and throwing the exception but
- 3:42:44let's assume from this piece of code it
- 3:42:46is connecting to the DV or it is
- 3:42:49connecting to the AWS S3 okay in that
- 3:42:52case if that particular server is down
- 3:42:54and your request is not able to process
- 3:42:57then that is the wrong practice right
- 3:42:59because you completely lost that input
- 3:43:02what you received from your producer
- 3:43:04because of because of your AWS failure
- 3:43:06or because of your DV failure since this
- 3:43:09is the simple scenario it's fine we are
- 3:43:11just throwing the exception okay but
- 3:43:13still we are losing the that particular
- 3:43:16event or messages and we are not doing
- 3:43:18anything with that so what we need to do
- 3:43:21as for the presentation what we
- 3:43:23understand we need to keep a retry we
- 3:43:26need to tell to the capka hey can you
- 3:43:29please retry for two or three times that
- 3:43:31is what the number I will tell to you
- 3:43:33once kapka will do the retry if within
- 3:43:37that time period I mean let's say I have
- 3:43:38defined for three if after three times
- 3:43:42reattempt if it is able to resolve the
- 3:43:44issue
- 3:43:45then well and good if not if the ret
- 3:43:48account exceed then it will simply push
- 3:43:51that failed event or failed messages to
- 3:43:55De letterer topic okay that is what we
- 3:43:58just need to tell to the capka so we
- 3:44:00need to tell to the capai about the
- 3:44:02retry and if retry count exceed then
- 3:44:05push the messages to the DT so that we
- 3:44:08have a copy of our failed messages or
- 3:44:10failed events so in future we can
- 3:44:13investigate it and we can reprocess it
- 3:44:15if required fine so just go back to the
- 3:44:20code now here let me stop the
- 3:44:24server so what you can do here simply
- 3:44:27you can tell to the capka by defining
- 3:44:30one
- 3:44:31annotation retriable topic okay here you
- 3:44:37can tell to the
- 3:44:38capka how many attempts you want to
- 3:44:41perform so you can Define I want to
- 3:44:44retry
- 3:44:45four times so by default attempt count
- 3:44:48is three if you will not specify any
- 3:44:50attempt kka will do three times attempt
- 3:44:54now since I have defined the four capka
- 3:44:57will do four times attempt I mean to the
- 3:45:00same consumer will be keep executing for
- 3:45:02four time okay so what will happen
- 3:45:06internally capka will create three
- 3:45:11topic just appending the suffix with
- 3:45:15retry okay each topic will retry once
- 3:45:19since we have defined there will be
- 3:45:21three times rise attempt okay because as
- 3:45:24I mentioned it will work on N minus one
- 3:45:28order okay we'll test in a moment now
- 3:45:31what next we are saying to the kapka by
- 3:45:33defining this particular annotation
- 3:45:35please enable retry mechanism for me and
- 3:45:38I want to perform four times retry okay
- 3:45:43now after four time four times retry if
- 3:45:47the issue is not resolved what do you
- 3:45:49want to perform we want to Simply push
- 3:45:51that messages to dead letter topic right
- 3:45:54so for that what you can do simply you
- 3:45:57can define a
- 3:46:00method
- 3:46:04public then just copy the same header
- 3:46:09because the event is
- 3:46:12same just Define here now simply just
- 3:46:16add a log
- 3:46:18statement do info and just Define some
- 3:46:25messages just add some meaningful
- 3:46:30messages now what we have defined here
- 3:46:33we are just saying what is the body and
- 3:46:38what is the topic from where we are
- 3:46:39listening and what is the offset count
- 3:46:42okay so I can capture the entire body
- 3:46:45but better to simplify the log statement
- 3:46:48I'll just capture the
- 3:46:51name
- 3:46:53then
- 3:46:55topic so this looks good now what next
- 3:46:59we need to annotate this method with DT
- 3:47:03Handler dead letter topic and Handler
- 3:47:07this is what The annotation if you open
- 3:47:09this annotation this is the empty
- 3:47:12one fine so using this annotation we can
- 3:47:16able to enable DT Logic for failure
- 3:47:20message based on our Capa topic okay so
- 3:47:24using this annotation this will listen
- 3:47:27all the failed messages so looks good
- 3:47:30now let me start the
- 3:47:34application so it started If You observe
- 3:47:37the console log carefully see here Java
- 3:47:41group DLT this is the topic name
- 3:47:45partition assign this and since as I
- 3:47:48already mentioned for if you specify the
- 3:47:51retry count for example go here ret
- 3:47:54count if you'll specify four it will
- 3:47:57create internally three topic to do the
- 3:48:00reattempt okay can you see here Java
- 3:48:04group retry zero this is one topic capka
- 3:48:07error handle retry zero retry
- 3:48:111 retry two can you see here now to just
- 3:48:17visualize it in a better way what I'll
- 3:48:19do I'll open the I'll just refresh this
- 3:48:23capka upset
- 3:48:26Explorer can you see here this is our
- 3:48:29main topic or Target topic where you
- 3:48:31want to push the messages and listen to
- 3:48:34it then it created DT can you see here I
- 3:48:38cannot maximize this let me try no I
- 3:48:41cannot do that but you can see here
- 3:48:43right then how many ret toic get created
- 3:48:46guys 0 1 and two so the total count is
- 3:48:51three here right that's super isn't it
- 3:48:55now let's do one thing let's quickly
- 3:48:57test this scenario let me rest okay we
- 3:49:01already restarted it right now I'll go
- 3:49:03to the
- 3:49:05postman first let me clear this
- 3:49:08console what I'll do I'll just use the
- 3:49:12same messages but I'll just give some
- 3:49:15valid IP so that I don't want to see any
- 3:49:17error first time so just send the
- 3:49:21request there is no error we got the
- 3:49:25success
- 3:49:26result this is the producer log this is
- 3:49:29the consumer log now let me clear this
- 3:49:33now let's try try with one failed
- 3:49:36scenario I'll give the IP address which
- 3:49:39is already restricted in my consumer
- 3:49:42okay let's change it to
- 3:49:45now let's observe carefully let me send
- 3:49:48the
- 3:49:51request see here it's doing the retry
- 3:49:53then it throw the exception we'll verify
- 3:49:55in the log itself
- 3:49:57okay see here this is the first time it
- 3:50:01received then it did the three times
- 3:50:04retry to verify okay B you just showing
- 3:50:07the log how can you get confirm that
- 3:50:10these three log statement being execute
- 3:50:13from the retri topic let me Zoom this
- 3:50:16first then what we'll
- 3:50:19do just go here because we are already
- 3:50:23logging the topic name and offset name
- 3:50:25right what is the topic name here guys
- 3:50:29capka error handle retry zero capka
- 3:50:32error handle retry one capka error
- 3:50:36handle retry two so you tried three
- 3:50:40times reattempt with three different
- 3:50:43topic now again if you want to check or
- 3:50:46if if you want to validate again go
- 3:50:49simply to your topic
- 3:50:52here just open the data you can see
- 3:50:56here each since in the log we found each
- 3:51:00topic retry topic receive the messages
- 3:51:03from the upset you can validate that
- 3:51:05directly in
- 3:51:06the topic
- 3:51:09itself can you see here fine now the
- 3:51:14retry is succeeded after three try three
- 3:51:17times it retried still we are not
- 3:51:20getting any result it is failed can you
- 3:51:24see here invalid IP address because
- 3:51:27since we are simulating the error
- 3:51:29scenario forcefully you'll get the error
- 3:51:32messages but between that retry if your
- 3:51:36DB connection or any aw specific infra
- 3:51:39isue is being resolved then you'll get
- 3:51:41the response I mean that is the happy
- 3:51:43scenario since we are getting the error
- 3:51:47after four times retry or after three
- 3:51:50times retry then where suppose to the
- 3:51:53message will
- 3:51:54go to DT topic right so if You observe
- 3:51:58here it created a topic capka error
- 3:52:02handle DLT now let's validate message
- 3:52:05went to this particular topic or not
- 3:52:08just run it can you see
- 3:52:12here this is the failed event went to
- 3:52:15the particular
- 3:52:18topic but if you'll send the not failed
- 3:52:22event I mean not failure scenario
- 3:52:25payload let's say
- 3:52:26239 I'll keep sending the
- 3:52:31request now if I
- 3:52:34validate in the upset
- 3:52:37Explorer no record right nothing will
- 3:52:40come to DT because only dead letter top
- 3:52:45will listen the failed events that is
- 3:52:48what we are saying in the code and there
- 3:52:50will be nothing in the
- 3:52:51retry only
- 3:52:54one good so it's working as expected
- 3:52:58right now there is also you can play
- 3:53:01with this retryable topic annotation now
- 3:53:05continuously our retry is happening
- 3:53:08right I mean there is no Gap but if you
- 3:53:10want to set the attempt in every or in
- 3:53:14some time interval you can do that you
- 3:53:18can simply Define
- 3:53:22backup okay then you can specify the
- 3:53:26what is the delay you expecting and what
- 3:53:28is the multiplier all these stops you
- 3:53:30can
- 3:53:33Define then what is the Marx
- 3:53:38delay you can define those value this is
- 3:53:41the number not integer for example I
- 3:53:43will just Define 3,000 I mean 3 second
- 3:53:48then multiply
- 3:53:511.5 and Max delay could be
- 3:53:5515,000 Max 15
- 3:53:59second now this is how you can Define
- 3:54:02the when you want to do the retry
- 3:54:05attempt not if you don't want to do it
- 3:54:07immediately if you want to specify your
- 3:54:09own time interval you can Define this
- 3:54:13now also you can tell to this retryable
- 3:54:15topic for what kind of exception you
- 3:54:19want to perform this retry now for
- 3:54:23runtime exception if you don't want to
- 3:54:25perform this retry what you can do
- 3:54:27simply you can Define here please
- 3:54:31exclude the retry for this
- 3:54:35exception what is the exception name
- 3:54:38null
- 3:54:40pointer do class or I don't run for this
- 3:54:44run exception right runtime
- 3:54:46exception do class you can also control
- 3:54:51for what exception you don't want to do
- 3:54:54the retry okay you can exclude them for
- 3:54:59this kind of exception there will be no
- 3:55:01retry so it is up to you based on your
- 3:55:04need you can configure
- 3:55:06it fine so this is how you can override
- 3:55:10the behavior of this retriable topic now
- 3:55:13let's do one thing let's try to produce
- 3:55:16bulk record and we'll see this
- 3:55:18particular Behavior okay let's go to the
- 3:55:21publisher I'll go to the controller
- 3:55:24class so I have a csb file list of user
- 3:55:28okay I mean I guess 100 is there there
- 3:55:31is total 100 user object I want to load
- 3:55:34this CSV file and we'll convert this
- 3:55:37particular to list of object this
- 3:55:40particular user in csb to list of user
- 3:55:43object then we'll push one by one to our
- 3:55:46publisher okay so it's a simple step I
- 3:55:49have just written one util class to load
- 3:55:53the CSV file and convert it to the list
- 3:55:55of user
- 3:55:56object then what I'll do I'll go to the
- 3:56:00controller
- 3:56:01class I'll will call this method to load
- 3:56:04the csb data and convert it to the list
- 3:56:08of user now simply what I'll do I'll
- 3:56:10just use users. for each user
- 3:56:16object then simply call
- 3:56:22this now it will trigger 100 event out
- 3:56:27of 100 event in my consumer I have
- 3:56:30restricted four user IP address okay it
- 3:56:36means out of 100 row 100 user object
- 3:56:40four user will be discarded and we go to
- 3:56:43the DLT okay when I say DT it is dead
- 3:56:48letter topic okay now if You observe we
- 3:56:51have only one record in our capka error
- 3:56:55handle DT if I run it again only one
- 3:56:58record now once I will execute this
- 3:57:02there should be four more record which
- 3:57:04need to add in this particular topic
- 3:57:07okay so what we'll do because see the
- 3:57:10simple thing guys whatever the IP I have
- 3:57:12specified here see I have taken it from
- 3:57:16this particular
- 3:57:19csb
- 3:57:20okay all the four are there I took from
- 3:57:23here only now let's do one thing let's
- 3:57:28rerun our
- 3:57:31app okay let me stop
- 3:57:34this because we need to it's fine we can
- 3:57:37run it I mean what I supposed to do in
- 3:57:40the controller class now I'm not giving
- 3:57:43the input even though I'm giving that
- 3:57:45input it doesn't make any sense anyway
- 3:57:48at the end it will iterate the csb and
- 3:57:51will push those messages what I'm giving
- 3:57:54here is not being pushed so I supposed
- 3:57:56to define a new endpoint saying get or
- 3:57:58something like that but it's fine we can
- 3:58:00play with the same because whatever the
- 3:58:02value will give will not be used so fine
- 3:58:05let me clear
- 3:58:09this just go to the postman this you can
- 3:58:13give any value okay this is not not
- 3:58:15being used send the
- 3:58:20request message publish
- 3:58:23successfully now if you'll go and check
- 3:58:25your
- 3:58:27console can you see
- 3:58:31here these are the messages from
- 3:58:35producer okay and if you scroll
- 3:58:39down these are the messages right 1 2 3
- 3:58:434 5
- 3:58:45is there any error let's
- 3:58:48see invalid IP address
- 3:58:51received it will do the retry can you
- 3:58:54see here DT also received some messages
- 3:58:58and if You observe here it did retry I
- 3:59:00mean since there are 100 record it's not
- 3:59:04easy to check the log to find out the
- 3:59:07retry account and its validation now to
- 3:59:10validate in a simple way see there are
- 3:59:13lot of error right
- 3:59:14DT received this just verify one thing
- 3:59:17how many dted messages you have there in
- 3:59:20your
- 3:59:21console four right because we have only
- 3:59:24four restricted IP address out of 100
- 3:59:27whatever the file we have in the
- 3:59:29csb we have here okay let me show you
- 3:59:33here only four user IP address are not
- 3:59:37correct so rest others should process
- 3:59:39the result now to validate that go and
- 3:59:43check in your up set Explorer just run
- 3:59:46this can you see here 1 2 3 4
- 3:59:495 this is the first record which we
- 3:59:52added before executing This bul Record
- 3:59:56then 1 2 3 4 with the same IP
- 4:00:01address take anyone's IP and just check
- 4:00:05in your
- 4:00:10csb okay also check in your consumer
- 4:00:14whether you have restricted for this
- 4:00:15user or not can you see
- 4:00:18here fine so whenever you are performing
- 4:00:21the bulk upload kind of scenario using
- 4:00:23capka you can do this kind of validation
- 4:00:27using this ret template and for failed
- 4:00:30event rather than losing those data
- 4:00:33better to keep it in some different
- 4:00:35topic and you don't need to create that
- 4:00:37topic manually that will be Take Care by
- 4:00:39capka by defining the simple annotation
- 4:00:42DT Handler that is called dead letter
- 4:00:45topic okay and everything will be Take
- 4:00:48Care by capka so that you have the filed
- 4:00:50messages with you you can play with them
- 4:00:52or you can investigate with them in
- 4:00:54future whenever it needed to process
- 4:00:57those details so this is how you can
- 4:01:00handle the error in capka using this DT
- 4:01:04and retry mechanism
- 4:01:12okay in in this tutorial we'll learn
- 4:01:15about how to use abro schema to produce
- 4:01:18and consume messages using spring boot
- 4:01:20additionally we'll explore the concept
- 4:01:22of schema registry what does it mean and
- 4:01:25what benefits does it provide okay all
- 4:01:28right let's assume I have a stable capka
- 4:01:31producer and consumer application where
- 4:01:33I'm sending employee details to a capka
- 4:01:35topic and the consumer read it from the
- 4:01:38capka topic and perform some further
- 4:01:41operation this is a happy scenario of
- 4:01:43pop up mechanism right but what if I
- 4:01:46change the producer data if I change the
- 4:01:49record to this new format where I'm
- 4:01:52adding a new field called middle name
- 4:01:54and I'm renaming one existing field
- 4:01:57called email ID and I'm removing one
- 4:01:59field called do now the problem starts
- 4:02:02from here will my consumer work with
- 4:02:05these new records the answer is no
- 4:02:08because the consumer is not aare of data
- 4:02:10changes and that's a big problem isn't
- 4:02:13it if there is any changes in my data it
- 4:02:16will directly impact all my consumer
- 4:02:19which is not acceptable in real-time
- 4:02:21implementation so how to overcome this
- 4:02:24The Simple Solution is to write a new
- 4:02:26producer and consumer application so
- 4:02:28that new data will be produced and
- 4:02:30consumed by the new instance without
- 4:02:33impacting the existing flow however that
- 4:02:36requires too much work and that is not a
- 4:02:38feasible solution then again how to
- 4:02:42solve this problem or how to handle data
- 4:02:45changes or schema evalution so that's
- 4:02:48where confident kapka introduced the
- 4:02:50concept of abro schema and to store the
- 4:02:53schema it has a component called schema
- 4:02:55registry okay now let's go one step
- 4:02:57ahead and understand what this abro
- 4:03:00schema is how to use it and how to
- 4:03:02handle schema Evolution then further
- 4:03:05we'll understand how it works internally
- 4:03:08so when I am saying Abu schema it's
- 4:03:10nothing but a contract between your
- 4:03:12producer and consumer it represents the
- 4:03:15data that you are going to serialize
- 4:03:17while producing and deserialize while
- 4:03:20consuming now if I try to represent the
- 4:03:22equivalent Abu schema of our record it
- 4:03:25will look something like this now if You
- 4:03:27observe carefully we are defining in
- 4:03:29abro schema saying that okay these are
- 4:03:32my fields which the producer will send
- 4:03:35and the consumer will consume in case
- 4:03:37the producer is not sending any of the
- 4:03:39field or adding any of the new field
- 4:03:42then there should not be any imp on the
- 4:03:44consumer side also if any changes are
- 4:03:47made to the schema it will be considered
- 4:03:49as a new version okay that is what
- 4:03:52called schema Evolution I have the
- 4:03:55contract or schema with me now how can I
- 4:03:58use it that's pretty simple once you
- 4:04:01prepare the schema you need to generate
- 4:04:03the corresponding object out from the
- 4:04:05schema so you need to Define all the
- 4:04:08field whether it is optional or
- 4:04:10mandatory in your schema file then you
- 4:04:13need to generate the abro object to
- 4:04:16convert abro schema to abro object you
- 4:04:19have multiple option you can directly
- 4:04:21use the abro tool or you can use Maven
- 4:04:24or gradel plugins to generate the object
- 4:04:26for you once you have generated the abro
- 4:04:29class then going forward you can use
- 4:04:31this employee. Java as a data object
- 4:04:34which you can produce and consume but
- 4:04:37while producing this object we need a
- 4:04:39serializer from the producer end we will
- 4:04:42serialize the encoded message and send
- 4:04:45them back to the capka topic so we can
- 4:04:47use capka abro serializer similarly we
- 4:04:51need a deserializer from the consumer
- 4:04:53end who will deserialize abro encoded
- 4:04:56message and convert them back to the
- 4:04:58object so we can use capka abro der
- 4:05:02serializer that's straightforward right
- 4:05:04here capka abro serializer and der
- 4:05:06serializer takes all the responsibility
- 4:05:09of converting your encoded abro message
- 4:05:12to the corresponding object and perform
- 4:05:14the serialize and deserialize operation
- 4:05:17that's interesting but where is the
- 4:05:19validation part the encoded abro message
- 4:05:21that my producer is sending whether it
- 4:05:23matches to the schema I have defined or
- 4:05:25not where are we validating it or how
- 4:05:28does my deserializer this realize the
- 4:05:30raw bites to an object without knowing
- 4:05:32the schema that's the major part we need
- 4:05:35to deal with to avoid failure on schema
- 4:05:38Evolution right so that's where schema
- 4:05:40registry comes into the picture what
- 4:05:43schema registry does it store your
- 4:05:45schema so when abro serializer serialize
- 4:05:48the record it first validate and store
- 4:05:51the schema into the registry then when
- 4:05:53abro deserializer deserialize it first
- 4:05:56takes the schema from the registry
- 4:05:58validate it with the messages and then
- 4:06:01deserialize it back to an object the
- 4:06:03primary purpose of schema registry is to
- 4:06:06provide a way to store retrive and
- 4:06:08evolve the schema in a consistent manner
- 4:06:11so in the future if I'm making any
- 4:06:13changes to my existing schema it will
- 4:06:16create a new version and store it into
- 4:06:18the schema registry since schema
- 4:06:21registry has the flexibility to support
- 4:06:23both backward and forward compatibility
- 4:06:26if any upgrade happens to the schema it
- 4:06:28will support the old schema as well as
- 4:06:30the new schema that is how it handle
- 4:06:33schema Evolution hope this makes sense
- 4:06:36to you no worries we'll do a handson
- 4:06:38example with all the theories we
- 4:06:41understand okay so without any further
- 4:06:43delay let's get
- 4:06:44[Music]
- 4:07:04started so let's create a new project to
- 4:07:07demonstrate this schema registry I'll
- 4:07:09create a new
- 4:07:11project then Define all the required
- 4:07:17field click next then just Define all
- 4:07:21required dependency I add web
- 4:07:24dependency then I need Lum book
- 4:07:28dependency then I need capka
- 4:07:33dependency click
- 4:07:36next now once it imported just go to the
- 4:07:39pom.xml
- 4:07:41and we'll validate that we have added
- 4:07:43the web dependency and we added spring
- 4:07:46capka and we added the lumo dependency
- 4:07:50right and anything else no but why there
- 4:07:54is
- 4:07:55error okay it is saying stter parent
- 4:07:593.2.2 not found okay let's downgrade it
- 4:08:02to the
- 4:08:043.2.1 let me update the
- 4:08:07project now let me clear this next we
- 4:08:11need to set up our capka environment so
- 4:08:13you want to use the conflent capka you
- 4:08:16can download the binary distribution and
- 4:08:18you can start the Juke keeper capka
- 4:08:20server service registry everything from
- 4:08:23the command prompt but I would prefer
- 4:08:24here to use the docker Compass file I
- 4:08:27will just Define all the container I
- 4:08:29need so that if I'll do the docker
- 4:08:31compose of it will start all the
- 4:08:33services for me okay so I'll just Define
- 4:08:36the docker compose file in this root
- 4:08:40directory then I can Define all the
- 4:08:43services I need so if You observe here
- 4:08:46let me Zoom this for you I need the Juke
- 4:08:49keeper then I need the broker which is
- 4:08:52nothing the capka server and then I need
- 4:08:55the Capa tools so this Capa tools is
- 4:08:59available for administrative task so you
- 4:09:01can ignore it this is optional and I
- 4:09:03need the schema registry so you
- 4:09:06understand what is the role of schema
- 4:09:07registry right so we need this schema
- 4:09:09registry component who will store the
- 4:09:11schema and validate with the encoded
- 4:09:13message and we need the control center
- 4:09:16so this is just a UI platform to monitor
- 4:09:19your capka activity okay so this control
- 4:09:23center depends on Juke keeper broker and
- 4:09:25schema registry it collect all the
- 4:09:28information from these three component
- 4:09:30and mapped into the UI okay and it will
- 4:09:32run on the port 9021 and these are the
- 4:09:35environment variable of each Services
- 4:09:37what I have defined okay no is I'll
- 4:09:40share this GitHub code in the video
- 4:09:41description you guys can to the same
- 4:09:44file so now to start these Services we
- 4:09:47need to start our Docker first right so
- 4:09:49I'll just open my Docker desktop in my
- 4:09:52local machine and once it will be start
- 4:09:55then I can just trigger Docker compose
- 4:09:57up so that all the services what I have
- 4:09:59defined in this compos file will be up
- 4:10:02and running for me and I can do the
- 4:10:04further step with this confident Capa so
- 4:10:07let's wait it to start so Docker is up
- 4:10:10and running now now what I'll do I'll go
- 4:10:12to the terminal Al so simply I will
- 4:10:14trigger
- 4:10:17Docker
- 4:10:19compose hyen
- 4:10:21D okay now if You observe it is just
- 4:10:25starting my all the
- 4:10:27services so if You observe it started
- 4:10:30the Juke keeper capka tools then capka
- 4:10:33server schema registry and control
- 4:10:36center so to verify that whether it
- 4:10:38started or not what I'll do I will just
- 4:10:41check the port of my control center that
- 4:10:43is
- 4:10:459021 I'll go to the
- 4:10:47browser and I'll will trigger
- 4:10:509021 it's loading right so it means it
- 4:10:53started so we are good here now what you
- 4:10:56need to do we'll follow the same step
- 4:10:59what we have discussed here first we'll
- 4:11:01create the schema then using the mavin
- 4:11:03plugin we'll generate the corresponding
- 4:11:05abro class then we'll Define the
- 4:11:07producer and consumer and we have
- 4:11:09defined the serializer and deserializer
- 4:11:11and schema registry okay okay so the
- 4:11:13first step let's create the schema so
- 4:11:17just go back to ID in the resource I
- 4:11:19will create a folder called
- 4:11:23abro then I need to define the schema
- 4:11:26here okay so let me create a file called
- 4:11:30employee dot
- 4:11:34avsc so this do absc is the extension of
- 4:11:38abro schema just create it then here you
- 4:11:43you need to Define what is the package
- 4:11:45where you want to keep your generated
- 4:11:47class so just Define the package name
- 4:11:50I'll define com dot
- 4:11:52java. dto let me Zoom this for
- 4:11:56you fine then next you need to Define
- 4:12:00what is the type so you want to define
- 4:12:03the abro record right so I'll Define the
- 4:12:06type as a record now what what is the
- 4:12:10class name you want to keep so from the
- 4:12:12schema the class what it will generate
- 4:12:15what what name you want to give okay so
- 4:12:17the name I'll Define
- 4:12:21employee next to that you need to Define
- 4:12:24all the field which you want to produce
- 4:12:26and consume so what you can do you can
- 4:12:29Define fields that need to be one array
- 4:12:32and you can Define all the field you
- 4:12:34need so if you'll verify the field we
- 4:12:37have defined ID first name last name
- 4:12:40email and do so let's go to the ID I'll
- 4:12:44Define the field name
- 4:12:46is the first field name
- 4:12:51is ID of
- 4:12:55type it will be
- 4:12:58string okay simple right now I will
- 4:13:02Define another field let me copy the
- 4:13:05same the field name is first name okay
- 4:13:12and okay there is a spelling
- 4:13:14mistake rst name of type string then the
- 4:13:19next last
- 4:13:22name type is
- 4:13:25string then what
- 4:13:27else first name last name email and date
- 4:13:30of Worth right so Define the email
- 4:13:35here up type string then D
- 4:13:40OB up type string and also let me add
- 4:13:43one additional field called age of type
- 4:13:48int okay you can Define any data type
- 4:13:52here string int double as for your
- 4:13:55requirement okay also you want to keep
- 4:13:58let's say I just want this email to be
- 4:14:00optional so what I can do I can Define
- 4:14:03here default value should be you can
- 4:14:06define null or you can Define
- 4:14:10empty okay so this is what the schema I
- 4:14:13have defined this particular field my
- 4:14:16producer will publish and my consumer
- 4:14:18will consume it that's fine now next If
- 4:14:21You observe the flow we have the abro
- 4:14:23schema with us next we need to generate
- 4:14:26the abro object by using the mavin
- 4:14:28plugin okay so let's go to the ID and
- 4:14:31we'll add all the required dependency
- 4:14:33and
- 4:14:34plugin so just go to the pal.
- 4:14:37XML and then I will add all the required
- 4:14:42dependency so if You observe here we
- 4:14:44need the abro serializer let me Zoom
- 4:14:46this for you and then we need the schema
- 4:14:49registry and we need this abro object
- 4:14:53okay then we just need to add the
- 4:14:56plugin just add it
- 4:14:59here just update
- 4:15:01it so we are getting the error but it is
- 4:15:04saying cannot resolve this dependency
- 4:15:07okay so it is looking from the mavin
- 4:15:10central repo to load these two depend
- 4:15:13schema registry and abro serializer but
- 4:15:16in Maven it is not available okay so
- 4:15:19what I'll do I'll tell in the pum to
- 4:15:22load it from the conflent repository
- 4:15:26okay now if I'll just update
- 4:15:31it there should not be any
- 4:15:34error okay there is no error now if
- 4:15:37you'll go to the plugin
- 4:15:39section If You observe here in the
- 4:15:42configuration
- 4:15:43I am telling here this is what my source
- 4:15:45directory where where you can find my
- 4:15:48schema and this is what the output
- 4:15:50directory where I want you to generate
- 4:15:53the object for me okay and the version
- 4:15:56of abro Maven plugin we are using
- 4:16:001.8.2 so that's fine now let's run the
- 4:16:03maven goal to generate the object for us
- 4:16:07just go to the mavin run the
- 4:16:11install so here here simply we are
- 4:16:13telling to take the abro schema and
- 4:16:16generate the abro object for us so the
- 4:16:19build is succeed now to validate that
- 4:16:21what we can do go to the main Java now
- 4:16:25you can see here right there is a
- 4:16:27package called D and this is what the
- 4:16:31package name we have defined in our
- 4:16:33schema can you see here the name space
- 4:16:35or the package we want to generate the
- 4:16:37abro object is this now if I'll open
- 4:16:41this I can see the abro class now if I
- 4:16:44open the employee. class let me Zoom
- 4:16:46this for
- 4:16:47you this class extends
- 4:16:50from specific record base specific
- 4:16:54record and this is are the schema what
- 4:16:57we have
- 4:16:58defined fine now if you'll scroll down
- 4:17:01we will find all the field can you see
- 4:17:04here ID first name last name email and
- 4:17:08date of birth and age also you can find
- 4:17:11the Constructor here this is your the
- 4:17:13default Constructor and you can find the
- 4:17:16getter and Setter method as well I mean
- 4:17:18this is just a class generated by the
- 4:17:21abro now we can use this particular
- 4:17:24employee as a data object to perform the
- 4:17:27producer I mean we can publish this
- 4:17:29particular object and we can consume
- 4:17:31this particular object so to do that
- 4:17:34let's create one producer and consumer
- 4:17:36application okay so what I'll do I'll
- 4:17:39create a separate package here
- 4:17:41itself and I'll create another for
- 4:17:44Consumer so I'll just Define
- 4:17:48package fine so let's create one
- 4:17:51producer and consumer class but before
- 4:17:54that let's verify the flow what we have
- 4:17:56discussed we have we have the schema
- 4:17:58with us and we created the employee.
- 4:18:00Java class which is nothing the abro
- 4:18:02record now we'll create the producer and
- 4:18:05consumer then we'll Define the abro
- 4:18:07serializer and deserializer and we'll
- 4:18:09validate the schema registry I mean we
- 4:18:12also need to give the configuration of
- 4:18:13schema registry to store all the schema
- 4:18:17okay so we are in the third step to
- 4:18:19create the producer and consumer so just
- 4:18:22go to the idea and then I'll just create
- 4:18:24a class
- 4:18:30here then just annotate here at theate
- 4:18:34service then I just need to define the
- 4:18:37capka template right
- 4:18:40private capka template
- 4:18:44of type you can Define string and of
- 4:18:48object employee right that is what our
- 4:18:50object abro object so I can define
- 4:18:53template just do the auto
- 4:18:55add then just Define the method public
- 4:19:00board send something like that okay
- 4:19:03Define the object what you want to send
- 4:19:05employee this employe is not our object
- 4:19:07right it's the abro object I mean you
- 4:19:10can consider it as a abro record so to
- 4:19:13do that I will just do Capa
- 4:19:15template do
- 4:19:18send what is the topic name we have not
- 4:19:21defined the topic and we want this topic
- 4:19:23to be autocreate so for that what I'll
- 4:19:26do I'll just create a config
- 4:19:28class or I'll just create a package
- 4:19:30called
- 4:19:31config then I'll Define class name capka
- 4:19:36config so I just need to annotate this
- 4:19:39with at theate
- 4:19:41configuration then I'll just create the
- 4:19:44topic I just need to create the B of
- 4:19:47topic
- 4:19:48public new topic this is the class name
- 4:19:52okay now create topic I'll just Define
- 4:19:56the method
- 4:19:58name then
- 4:20:01return see if You observe here you just
- 4:20:04need to define the name of your topic
- 4:20:06number of partition and replication
- 4:20:08Factor so I'll will just Define the name
- 4:20:11let's say Java
- 4:20:14hyphen
- 4:20:15abro and number of partition I want
- 4:20:18three and replication factor I want one
- 4:20:22so I can
- 4:20:24Define
- 4:20:26one just need to type c to
- 4:20:30sort just annotate this at theate B let
- 4:20:33me maximize this yeah Define at theate
- 4:20:37ban now I don't want to hard code this
- 4:20:39topic name here so what I want to do I
- 4:20:43have the application. properties file or
- 4:20:45I will just create a file called yml
- 4:20:49file and here I can Define my topic name
- 4:20:53okay so I'll Define
- 4:20:55topic name equal to this so I want to
- 4:20:59make it generic so I have defined the
- 4:21:01topic name in my configuration I mean
- 4:21:03the yml configuration and that I will
- 4:21:05load in my config so you know how to
- 4:21:08load the value from the properties file
- 4:21:10right I'll just Define private
- 4:21:14string topic
- 4:21:16name I just need to annotate your at
- 4:21:19theate
- 4:21:20Value and I just need to get that value
- 4:21:23from my yl file so just Define the
- 4:21:25dollar and give the key what is the key
- 4:21:28name topic dot
- 4:21:31name that's it right so I no need to
- 4:21:34hard code here I can use the topic
- 4:21:38name same topic name I can use while
- 4:21:41producing the messages
- 4:21:43so here also I can use the topic from
- 4:21:45the properties file so let me copy this
- 4:21:51syntax I need this topic to publish the
- 4:21:54messages and what data you want to
- 4:21:57publish
- 4:21:59employee
- 4:22:01fine next I'll just capture the return
- 4:22:04statement so that we can just add the
- 4:22:07log statement what data you are sending
- 4:22:09what is the offset count and everything
- 4:22:10okay so this will give give you the
- 4:22:13future
- 4:22:14object can you see here this will return
- 4:22:16the completable future now I'll take
- 4:22:19that future object and I'll will just
- 4:22:21get the result okay so this will be
- 4:22:25employee nothing I'm just checking if
- 4:22:27there is no exception print that okay
- 4:22:30message has been sent to this this this
- 4:22:32is what the message we have sent and
- 4:22:33this is what the upset count but if
- 4:22:35there is a error just print that unable
- 4:22:37to send the messages okay that that is
- 4:22:40the simple things we have done here now
- 4:22:42let me Define the consumer let me copy
- 4:22:45the name I'll go to the consumer
- 4:22:48package create a new
- 4:22:51class fine consumer is the simple one
- 4:22:54what you need to do just Define The
- 4:22:56annotation first to make it a component
- 4:22:59Define service then you just need to
- 4:23:02write a method to listen to the topic so
- 4:23:05just Define
- 4:23:06public
- 4:23:08boid read
- 4:23:11messages which type of messages it will
- 4:23:13read let's say Define
- 4:23:17consumer consumer
- 4:23:19record of type key as a string value as
- 4:23:22a
- 4:23:25employee fine now I just need to tell to
- 4:23:29this consumer from where he need to read
- 4:23:31that messages okay so just annoted Capa
- 4:23:34listener and just Define what is the
- 4:23:38topics just tell that okay read it from
- 4:23:42the topic what I have defined in my
- 4:23:45properties file that is the reason I
- 4:23:47move it to the properties file to make
- 4:23:48it generic rather than hard code in each
- 4:23:50and every place I have just defined it
- 4:23:53in a single place okay so topic dot
- 4:23:57name so we telling to this particular
- 4:24:00listener to read the messages from topic
- 4:24:03what I have defined and what what
- 4:24:05messages it will read of type employee
- 4:24:08so what we can do just add the log
- 4:24:10statement okay to verify that SL
- 4:24:13forj this we added from the lambo now
- 4:24:16here I will simply first get the key are
- 4:24:19you sending any key no we are not
- 4:24:20sending any key right but still let's
- 4:24:23capture it or we can pass it from the um
- 4:24:26producer
- 4:24:27itself consumer record do get
- 4:24:32key what is the value we are
- 4:24:35sending consumer record do
- 4:24:39value so now we can just add a log
- 4:24:41statement
- 4:24:43just add the log statement abro message
- 4:24:46received for the key is this and value
- 4:24:48is employee right okay I have done the
- 4:24:52mistake this this value should be of
- 4:24:54type employee that is what we have
- 4:24:56defined so this need to be
- 4:24:58employee because this is what the value
- 4:25:00you are getting right key is a string it
- 4:25:03can be serialized and deserialized by a
- 4:25:05string serializer or der serializer but
- 4:25:08value is something object which is abro
- 4:25:11object and it it needs to be serialized
- 4:25:14and deserialized by the ab serializer
- 4:25:16and der serializer so that we are going
- 4:25:18to Define in a moment so this is what we
- 4:25:21are
- 4:25:22getting fine we have defined the
- 4:25:25producer and we have defined the
- 4:25:26consumer now let me Define a controller
- 4:25:29so that I can trigger an event so just
- 4:25:32create another
- 4:25:35class just Define the package called
- 4:25:39controller just annoted here at theate r
- 4:25:44controller then you can define a I mean
- 4:25:47you need to inject the
- 4:25:50producer then just Define an end point
- 4:25:53to trigger the I mean event to the cap
- 4:25:56copy okay so I'll just Define public
- 4:25:59[Music]
- 4:26:01string what type of messages we want to
- 4:26:04send
- 4:26:05employee then you can simply Define your
- 4:26:09producer dot send what is that employee
- 4:26:15object and I can return some dummy
- 4:26:18string
- 4:26:21message and this need to be annoted at
- 4:26:23theate request
- 4:26:25body and I need to annoted here at theed
- 4:26:29post mapping since I have the request
- 4:26:30body and I'll Define the endpoint
- 4:26:34events and this particular application I
- 4:26:38want to run it in a different port so
- 4:26:40I'll just Define
- 4:26:47serverport let's say 9292 or let's take
- 4:26:51some different number okay
- 4:26:548 1 81 something like
- 4:26:58that we are good now okay so if you
- 4:27:02cross verify the flow we have created
- 4:27:04the producer and consumer and now the
- 4:27:07next step we need to Define what is the
- 4:27:10serializer of my value what is the
- 4:27:12deserializer of my value and then where
- 4:27:15is my bootstrap server is up and running
- 4:27:18everything I need to Define as a
- 4:27:20configuration so what I can do I'll
- 4:27:23prefer to Define in a application. yml
- 4:27:25file you can also create a javab base
- 4:27:27config class but I can Define here in
- 4:27:30the yl file that is the easy approach
- 4:27:33okay because I don't want more
- 4:27:35customization so I just need to use the
- 4:27:38producer and consumer config in the yml
- 4:27:40file but if you need more configuration
- 4:27:42I mean customization of more
- 4:27:44configuration you need to go with the
- 4:27:46javab base config approach so I need to
- 4:27:48Define here
- 4:27:49spring
- 4:27:51capka producer first okay in the
- 4:27:55producer first I need to Define where is
- 4:27:58my bootstrap server I mean where is my
- 4:28:00capka is running so where exactly it is
- 4:28:03running you can
- 4:28:06Define Local Host 9092 right that is
- 4:28:10what usually we Define Local Host
- 4:28:139092 so if You observe your Docker
- 4:28:15compose
- 4:28:17file go to the docker compose file the
- 4:28:21bootstrap
- 4:28:23server if you scroll yeah can you see
- 4:28:26here this is what
- 4:28:27your kka server right broker and it is
- 4:28:31running on Port 9092 that is what we are
- 4:28:33just defining here so just go to the ml
- 4:28:39file 9092 but we are not running it in
- 4:28:43our local we are running it in a Docker
- 4:28:46so you need to give the default IP so my
- 4:28:51default IP is 127 do 0 do 0 do
- 4:28:581 colum
- 4:29:009092 fine
- 4:29:03now now I can Define the
- 4:29:08serializer key serializer okay so the
- 4:29:11key seral ier is the plain string I can
- 4:29:14use this string serializer which is the
- 4:29:16default given by the capka itself okay
- 4:29:19but if I Define the value serializer I
- 4:29:22cannot Define the Json D serializer here
- 4:29:25I need to define the abro capka abizer
- 4:29:29okay coming from the confluent that is
- 4:29:32what we have discussed in our flow right
- 4:29:34we need the capka abro serializer to
- 4:29:38serialize my value to the topic so we
- 4:29:41need to Define that so just Define the
- 4:29:45class io. conf. kapka serializers can
- 4:29:49you see the class name capka abro
- 4:29:52serializer we are using this abro
- 4:29:54serializer to take the encoded abro
- 4:29:58messages and send it to the Capa opy so
- 4:30:02we are good now also we need to define
- 4:30:04the schema registry right so I'll just
- 4:30:07define
- 4:30:08properties
- 4:30:10schema registry
- 4:30:13now where is your schema registry is up
- 4:30:15and running if you go to your Docker
- 4:30:17compos file the registry is running on
- 4:30:21Port
- 4:30:238081 right so I need to Define
- 4:30:27that
- 4:30:30HTTP then your IP and your Port that is
- 4:30:368081 and my IP is this
- 4:30:38one fine this is where my schema
- 4:30:41registry is stop and
- 4:30:43run we are good now let's define the
- 4:30:46consumer
- 4:30:49configuration Now consumer also you need
- 4:30:51to define the bootstrap server this are
- 4:30:54the bootstrap server let me copy
- 4:30:58this next here also if You observe in
- 4:31:01the consumer part we need to Define d
- 4:31:04serializer properties right so I can
- 4:31:06Define here key der
- 4:31:09serializer is the same that is my string
- 4:31:12so I using the default Capa string der
- 4:31:15serializer but value der serializer
- 4:31:18should be from IO confl capka abro
- 4:31:22realizer just Define that class we are
- 4:31:25good here now next you need to Define
- 4:31:29Auto offset reset you can Define it as
- 4:31:33earliest or latest it is up to
- 4:31:35you next in consumer also you need to
- 4:31:39tell to that consumer okay where is your
- 4:31:41schema registry up and running If You
- 4:31:43observe the flow both abro serializer
- 4:31:47and abro der serializer is connecting to
- 4:31:49the schema registry right so you need to
- 4:31:52give information about schema registry
- 4:31:54to both your consumer and producer so we
- 4:31:57have given here in the producer now same
- 4:32:01let's give in the consumer so it need to
- 4:32:05be inside the code so that's fine I'll
- 4:32:09just define properties
- 4:32:15next just Define one additional field
- 4:32:18that is to tell to the consumer to read
- 4:32:21the abro messages okay so there is
- 4:32:24something
- 4:32:26called
- 4:32:28specific okay and
- 4:32:32abro
- 4:32:35reader equal to
- 4:32:37true that's it okay so the main thing we
- 4:32:41need to Def Define serializer and
- 4:32:44deserializer that is what we want to
- 4:32:46Define and also you just need to Define
- 4:32:49about your schema registry so these are
- 4:32:51the three component we are missing
- 4:32:53serializer der serializer and schema
- 4:32:55registry and these three we just
- 4:32:57configured in this yl file okay but one
- 4:33:01thing I observed here producer also need
- 4:33:04this property bootstrap hyphen server
- 4:33:06and consumer also need this property so
- 4:33:08better let's not Define in both the
- 4:33:10place I can Define it in
- 4:33:14a global label in capka itself bootstrap
- 4:33:19server this looks good so kapka having
- 4:33:23the bootstrap server for both producer
- 4:33:25and consumer and in producer we have
- 4:33:27Define the focus on this value
- 4:33:29serializer okay we have defined capka
- 4:33:32abro serializer and we are telling to
- 4:33:34the serializer about where is my schema
- 4:33:37registry up and running so that he can
- 4:33:39add the schema to this registry and and
- 4:33:41we have the Der serializer and we are
- 4:33:44telling the again about the schema
- 4:33:46registry and we are telling that here
- 4:33:48okay read the abro messages fine so we
- 4:33:52are good here now let me close
- 4:33:55everything just go to the
- 4:33:59producer I just want to add a key okay
- 4:34:02while sending the event I also want to
- 4:34:05add a key up type let me Zoom this for
- 4:34:08you up Type U ID I want to send some
- 4:34:11random string okay random uid do2
- 4:34:17string so while publishing the messages
- 4:34:20I'm sending the key along with that I'm
- 4:34:22sending the abro messages and if I'll
- 4:34:25open the
- 4:34:26consumer in consumer we are consuming
- 4:34:29both key and value key is nothing the
- 4:34:33string value is nothing my ab
- 4:34:35messages and we have defined all the
- 4:34:38component we understand in this
- 4:34:40architecture fine so I believe we are
- 4:34:42good let's start the application and
- 4:34:44we'll see how it works so let's go to
- 4:34:48the main
- 4:34:49class just start the
- 4:34:55application so we are getting some error
- 4:34:58here what is the error okay no group ID
- 4:35:01found in consumer config so that is what
- 4:35:03the mistake we have done we have not
- 4:35:05specify the group ID so I'll just add
- 4:35:09group ID as a some some random Bel okay
- 4:35:13Java
- 4:35:15Tey
- 4:35:17new all good let's restart it
- 4:35:22now so we are good now if You observe in
- 4:35:24the console the Java new is nothing my
- 4:35:27consumer group assigned to this
- 4:35:30particular topic javat abro with the
- 4:35:33three partition 0 1 and two can you see
- 4:35:36here let me clear this now I want to
- 4:35:39verify the same thing in the control
- 4:35:42center okay just go to the browser this
- 4:35:45is where the 9021 is my control center
- 4:35:48right now this is where the cluster if
- 4:35:50I'll open this it should have our topic
- 4:35:53info can you see here the topic name is
- 4:35:56there java abro with the three partition
- 4:35:59if I'll open this and if I'll check the
- 4:36:02messages you won't find any messages
- 4:36:04because you have not published
- 4:36:06anything so there is nothing it's still
- 4:36:09loading no messages right so so let me
- 4:36:12keep a duplicate page of this particular
- 4:36:15thing because until you are open to this
- 4:36:18page you can see only those messages so
- 4:36:20this is just a tool to monitor the event
- 4:36:23whether it is going through your
- 4:36:24producer and consumed or not for the
- 4:36:27simple purpose now if I'll check the
- 4:36:30schema
- 4:36:31here I don't have any schema maap to
- 4:36:33this particular topic so here when
- 4:36:36schema will be added when it will be
- 4:36:39serialized then only capka AB serialize
- 4:36:42will push that particular schema to
- 4:36:44schema registry that is what we
- 4:36:46understand in the theory right let's
- 4:36:48verify it so what I'll do I'll go to the
- 4:36:51postman and what is the end
- 4:36:54point let's verify in the controller
- 4:36:56class go to the
- 4:36:58controller events and what is the port
- 4:37:01number 8181 I believe 8181 yeah go to
- 4:37:05the
- 4:37:06postman
- 4:37:088181 events and what is the request body
- 4:37:12we have lot of field let me add those
- 4:37:24field so these are the field we have
- 4:37:27defined in our schema okay so we can
- 4:37:29only pass this field and then we'll play
- 4:37:32with the schema Evolution once we'll go
- 4:37:34through one happy
- 4:37:35scenario just format it before I send
- 4:37:38the request let's check the status of of
- 4:37:41schema registry whether we'll verify
- 4:37:44whether schema is added to the schema
- 4:37:46registry or not okay how can I validate
- 4:37:48that so the schema registry is running
- 4:37:51on Port 8081 right first let me okay let
- 4:37:55me
- 4:37:568081 I want to get all the topic
- 4:38:01name so topic is not created here not
- 4:38:04sure why it's not showing in the schema
- 4:38:05registry I believe after sending one
- 4:38:08object it will show but okay let me
- 4:38:10verify the latis
- 4:38:13schema okay we cannot do that since the
- 4:38:16topic is not showing here let me refresh
- 4:38:17it nothing but topic is there right
- 4:38:20that's fine let me trigger the
- 4:38:26request can you see here we got the
- 4:38:28success result message published now if
- 4:38:31you check in your console send messages
- 4:38:33this is from producer end these are the
- 4:38:35information we are sending and it went
- 4:38:37to the upset zero and this is from
- 4:38:41consumer end can you see
- 4:38:44here this is what theed message and this
- 4:38:47is what the messages we send from the
- 4:38:49producer now to validate that just go
- 4:38:52here can you see here in the messages we
- 4:38:55can see the
- 4:38:56message this is what the messages we
- 4:38:59have sent just now right and if you now
- 4:39:01validate the schema just check the
- 4:39:04schema you can see the schema here these
- 4:39:07are the field you have defined right ID
- 4:39:10first name last name Emil with the
- 4:39:12default field date of birth age and the
- 4:39:16abro class name package everything all
- 4:39:18the schema we have Define is imported
- 4:39:20here once the abro serializer is
- 4:39:22serialized and send the event to the
- 4:39:24topic now now let's verify the topic
- 4:39:27name yeah can you see here the topic
- 4:39:30name is this now to validate the schema
- 4:39:33in the schema registry what you can do
- 4:39:37Local Host 8081 subjects then topic name
- 4:39:40what is the topic name copy
- 4:39:43again just replace the topic
- 4:39:48name can you see here let me Zoom this
- 4:39:50for you subject is nothing the topic
- 4:39:53name I mean the value as a suffix will
- 4:39:56be appened by the kfka conflent and
- 4:39:58version is one ID is 61 this is what the
- 4:40:01schema we can see the schema in the
- 4:40:03schema registry so we concluded that
- 4:40:06schema registry has the schema
- 4:40:08information while abro serializer
- 4:40:10serialize it it added to the schema
- 4:40:12registry and since we are able to
- 4:40:14receive the messages from the consumer
- 4:40:16and ab dis realizer dis realize it okay
- 4:40:20and which version of schema we have the
- 4:40:23current version is one because there is
- 4:40:24no change this is the fres schema added
- 4:40:27to the schema registry so this is the
- 4:40:30happy scenario we validate right we can
- 4:40:32able to produce the message and we can
- 4:40:34able to consume the message by taking
- 4:40:36the help of capka bro serializer and der
- 4:40:38serializer and also we validate the
- 4:40:40schema
- 4:40:41in schema registry and we can able to
- 4:40:44visualize the messages and everything
- 4:40:45from the control center that's fine but
- 4:40:49Our intention here to check the schema
- 4:40:52evaluation related things so if I'm
- 4:40:54doing any change on my producer data
- 4:40:56whether it will impact to my consumer or
- 4:40:59not that is what we want to verify and
- 4:41:01that is why we are using this schema
- 4:41:03registry so let's jump into the
- 4:41:05presentation and then we'll perform
- 4:41:08these schema changes when I say schema
- 4:41:10changes we just want to do some
- 4:41:13modification on our object which we want
- 4:41:15to produce and consume so what we'll do
- 4:41:19we'll remove the date of birth field and
- 4:41:21age field and we'll rename this
- 4:41:24particular field email to email ID in
- 4:41:26our existing schema so just go to the
- 4:41:31schema so the schema is employee.
- 4:41:34absc here we want to remove the Dov and
- 4:41:38age just remove
- 4:41:40it and then I want to rename this email
- 4:41:44to email
- 4:41:45ID so we have done two changes here we
- 4:41:48have removed two field do and age and we
- 4:41:52rename one of the existing field now
- 4:41:55let's see if I run the application
- 4:41:57whether it will break my consumer or my
- 4:42:00consumer will work as it is without any
- 4:42:02impact okay that is what we want to
- 4:42:04verify but since we have done the
- 4:42:06changes on the schema we need to
- 4:42:08regenerate the object abro object let me
- 4:42:11stop the
- 4:42:14application then just execute the mavin
- 4:42:17goal
- 4:42:19install so the build is succeed now to
- 4:42:22validate whether those field are
- 4:42:24excluded from our abro object generated
- 4:42:27by the abro plugin we can go to the dto
- 4:42:31and we can validate can you see here
- 4:42:34those fields are not part of the Abra
- 4:42:37object
- 4:42:38now now let's run the application and
- 4:42:41we'll trigger the end point to publish
- 4:42:43the event then we'll validate whether
- 4:42:46the consumer is working or it is giving
- 4:42:48any error or exception okay so that is
- 4:42:51the reason we are using schema registry
- 4:42:53to manage the schema if I'm doing any
- 4:42:56change on my produce object or abro
- 4:42:59object okay so let's wait it to
- 4:43:02complete so it started now I'll go to
- 4:43:05the postman and I'll change the field
- 4:43:09actually we have removed these two field
- 4:43:10right d and and AG just remove it age
- 4:43:13and do is removed and then this will be
- 4:43:16email ID that is what the changes we
- 4:43:18have done cool now let's trigger the
- 4:43:23API we got the response message
- 4:43:25published now if you'll go to the
- 4:43:28console you can see here right ID first
- 4:43:32name last name and email ID this is what
- 4:43:35the object we have produced from the
- 4:43:37producer and this is what the serialized
- 4:43:40as part of the
- 4:43:42capka abro serializer and then if You
- 4:43:44observe here this are the message we
- 4:43:46have received in the consumer end and in
- 4:43:48both the place producer and consumer we
- 4:43:50don't have those field now if you'll go
- 4:43:53and check in the control
- 4:43:56center if I'll go on schema if I'll come
- 4:44:00back to the messages so message might
- 4:44:03not be present here so I can retrigger
- 4:44:05it to validate
- 4:44:06that so just send the request
- 4:44:10again
- 4:44:13you can see the messages here
- 4:44:15right those field is not part of the
- 4:44:18messages
- 4:44:19itself nothing is there right now if
- 4:44:22I'll go and check the schema Let me
- 4:44:25refresh this I'll come back to schema
- 4:44:28again can you see here schema is also
- 4:44:32updated ID first name last name email ID
- 4:44:36now if you'll check the version history
- 4:44:38there are two version the initial
- 4:44:40version and second version in initial
- 4:44:43version we have the fields right d age
- 4:44:47and email but in the current version
- 4:44:51those fields are not there so if you
- 4:44:53want to compare you can click on this
- 4:44:55turn on version difference and you can
- 4:44:57see here these are the field has been
- 4:45:01removed so you can see the red mark here
- 4:45:04and the value is renamed you can see
- 4:45:06this okay this is are the difference now
- 4:45:09if you want to validate the the same
- 4:45:11thing in the schema registry instead of
- 4:45:13this control center UI just trigger this
- 4:45:16request again can you see here the
- 4:45:19version is two and the schema is this
- 4:45:23this is are the latest schema okay I
- 4:45:26have done the changes on the schema but
- 4:45:29still my consumer is not breaking I'm
- 4:45:31able to consume the messages even the
- 4:45:33field has been removed from the payload
- 4:45:35or abro messages itself so that is the
- 4:45:38main purpose of using this schema
- 4:45:40registry If You observe here it
- 4:45:42maintains the schema history okay so it
- 4:45:45maintains the different version of your
- 4:45:47schema since it has the backward
- 4:45:50compatibility it cross validate your
- 4:45:52abro messages with all the schema
- 4:45:55present in the schema registry now let's
- 4:45:57do another changes and we'll validate
- 4:46:00whether the schema evolution is working
- 4:46:02with this particular implementation or
- 4:46:04not so what changes when I say changes
- 4:46:07let's add one more field okay so go go
- 4:46:11to the project go to the schema I want
- 4:46:14to add some field so as for the
- 4:46:17discussion we need to add some different
- 4:46:20field like middle name or something like
- 4:46:21that okay so let me add that field first
- 4:46:25name last name let me copy the
- 4:46:27same and I want to
- 4:46:30add middle
- 4:46:33name now since I have done the changes
- 4:46:35on the schema I need to rebuild the
- 4:46:37project let me stop the capka I just
- 4:46:41need to run the M install
- 4:46:45again so you can see here build is
- 4:46:48succeed now if I'll go to my dto class
- 4:46:51and if I'll open it let me see whether
- 4:46:54the field is added or
- 4:46:56not yeah can you see here the middle
- 4:46:59name is added here now let's do one
- 4:47:01thing let's start the
- 4:47:04server and we'll add this middle name as
- 4:47:07part of our
- 4:47:09payload
- 4:47:12just
- 4:47:13check okay it's still running let's writ
- 4:47:17it to
- 4:47:20complete yeah so it is started now let
- 4:47:22me clear this and now what I'll do I
- 4:47:26will trigger the request okay again we
- 4:47:28have done the schema changes we have
- 4:47:30added the new field which is not present
- 4:47:33in my both the schema I mean schema
- 4:47:35version one and two does not have the
- 4:47:37information about this new field now
- 4:47:39let's see whether it is working or we
- 4:47:40are getting any error then we'll try to
- 4:47:43find out the solution of it send the
- 4:47:49request we are getting 5 internal server
- 4:47:53error now if you go and check in your
- 4:47:55console the error itself
- 4:47:58self-explanatory invalid configuration
- 4:48:00exception let me Zoom this for you
- 4:48:03schema being registered is incompatible
- 4:48:06with an earlier schema error code
- 4:48:0949 it means means whatever the schema I
- 4:48:12have with me and the value we you are
- 4:48:15giving now is not compatible then how
- 4:48:19can I make my schema as a comforable to
- 4:48:22avoid this 49 error so what you can do
- 4:48:26if you are adding any new field then
- 4:48:29initially mark it as a
- 4:48:33default okay then what you can do since
- 4:48:36you have done the changes just build it
- 4:48:38again then we'll start the server and
- 4:48:40and
- 4:48:42validate so build is succeed now let me
- 4:48:45restart the
- 4:48:46server so it started now go to the
- 4:48:49postman and send the request
- 4:48:52again the message has been published now
- 4:48:55if you go and check in the messages can
- 4:48:57you see the message here with the new
- 4:49:01field middle name can you see here we
- 4:49:04got the response now if Will validate
- 4:49:06the schema
- 4:49:08again version history there are total
- 4:49:11three version in first version it was
- 4:49:13the happy scenario we have all the field
- 4:49:16in second version we have removed the do
- 4:49:18and age and we rename email ID then
- 4:49:21let's see what is there in the third
- 4:49:23version third version we have just added
- 4:49:26this field middle name now if you'll
- 4:49:28compare it you can see this field is
- 4:49:31newly added okay and we are not getting
- 4:49:35any error so this is how schema registry
- 4:49:38maintains the schema and valid against
- 4:49:41the abro messages what you are producing
- 4:49:44and what you are dis realizing or
- 4:49:46consuming and this way you can avoid the
- 4:49:49failure in the schema Evolution which
- 4:49:51will not break your entire system okay
- 4:49:54so this is what the straightforward
- 4:49:55approach we have demonstrate about the
- 4:49:57schema evalution using schema registry
- 4:50:00and Abra producer and consumer
- 4:50:03example that's all about this particular
- 4:50:06video guys thanks for watching this
- 4:50:09video meet you soon with A New
- 4:50:14Concept
About this transcript
This page contains the full transcript of π Apache Kafka Crash Course With Spring Boot 3.0.x | @Javatechie by Java Techie, generated from the public captions YouTube serves with the video. The transcript has 41,182 words across 6,162 segments, with the original timestamps preserved so you can click any line to jump to that moment in the embedded player.
What you can do with it
Use the transcript to take notes, quote the speaker, build a study guide, generate a summary with ChatGPT or Claude via the YouTube Summary tool, or export it as a timed subtitle file with YouTube to SRT. You can also re-open it in the transcriber to translate the transcript into 100+ languages.
Free YouTube transcript tool
YouTube2Text is a free YouTube transcript generator β no signup, no daily limit. Paste any YouTube link and get the full transcript instantly, with timestamps, click-to-jump, translation to 100+ languages, AI prompts for ChatGPT, Claude, and Gemini, and exports to TXT, SRT, VTT, or Markdown.