{"id":1152,"date":"2026-08-19T02:38:14","date_gmt":"2026-08-19T02:38:14","guid":{"rendered":"https:\/\/www.devopsschool.com\/tutorials\/?p=1152"},"modified":"2026-08-19T02:38:15","modified_gmt":"2026-08-19T02:38:15","slug":"kafka-master-tutorials-series-6-kafka-consumer-deep-dive","status":"publish","type":"post","link":"https:\/\/www.devopsschool.com\/tutorials\/kafka-master-tutorials-series-6-kafka-consumer-deep-dive\/","title":{"rendered":"Kafka Master Tutorials Series: 6 &#8211; Kafka Consumer Deep Dive"},"content":{"rendered":"\n<h2 class=\"wp-block-heading\">Consumer Groups, Parallelism, Offsets, Rebalancing, Failover, Consumption Patterns and Production Tuning<\/h2>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\"><strong>Audience:<\/strong>&nbsp;Students and freshers with no previous Kafka experience<br><strong>Goal:<\/strong>&nbsp;Start with \u201cWhat is a consumer?\u201d and finish with production-grade consumer design, reliability, scaling, recovery, and performance tuning<br><strong>Training environment:<\/strong>&nbsp;Confluent Kafka Cluster<br><strong>Technical version note:<\/strong>&nbsp;This tutorial is aligned with modern Apache Kafka 4.x consumer concepts. Where the newer Consumer rebalance protocol differs from the older Classic protocol, both are explained explicitly.<\/p>\n<\/blockquote>\n\n\n\n<figure class=\"wp-block-image\"><img decoding=\"async\" src=\"https:\/\/file+.vscode-resource.vscode-cdn.net\/Users\/rajeshkumar\/Downloads\/Kafka_Consumer_Deep_Dive.jpg\" alt=\"Kafka Consumers \u2014 Complete Deep Dive\"\/><\/figure>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\"><strong>Image note:<\/strong>&nbsp;The infographic is a learning map. The written tutorial below is the source of truth for protocol details, especially the differences between the newer Consumer group protocol and the older Classic protocol.<\/p>\n<\/blockquote>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">1. Learning Objectives<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">By the end of this tutorial, you should be able to explain:<\/p>\n\n\n\n<ul class=\"wp-block-list\">\n<li>What a Kafka Consumer is<\/li>\n\n\n\n<li>How a consumer reads records from Kafka<\/li>\n\n\n\n<li>Why consumers\u00a0<strong>pull<\/strong>\u00a0data rather than Kafka pushing records to them<\/li>\n\n\n\n<li>How topic partitions map to consumers<\/li>\n\n\n\n<li>What a Consumer Group is<\/li>\n\n\n\n<li>How Consumer Groups provide parallel processing<\/li>\n\n\n\n<li>How Kafka balances workload across consumers<\/li>\n\n\n\n<li>Why the number of partitions limits parallelism inside one consumer group<\/li>\n\n\n\n<li>What the Group Coordinator does<\/li>\n\n\n\n<li>What rebalancing means<\/li>\n\n\n\n<li>What causes a rebalance<\/li>\n\n\n\n<li>What happens when a consumer joins<\/li>\n\n\n\n<li>What happens when a consumer leaves or crashes<\/li>\n\n\n\n<li>What happens when topic partitions change<\/li>\n\n\n\n<li>How the modern Consumer rebalance protocol differs from the Classic protocol<\/li>\n\n\n\n<li>What Kafka offsets are<\/li>\n\n\n\n<li>Current position vs committed offset<\/li>\n\n\n\n<li>Automatic vs manual offset commits<\/li>\n\n\n\n<li><code>commitSync()<\/code>\u00a0vs\u00a0<code>commitAsync()<\/code><\/li>\n\n\n\n<li>How to choose an offset strategy<\/li>\n\n\n\n<li>How consumer failover and recovery work<\/li>\n\n\n\n<li>Fan-out consumption<\/li>\n\n\n\n<li>Load-balanced consumption<\/li>\n\n\n\n<li>At-most-once vs at-least-once vs exactly-once approaches<\/li>\n\n\n\n<li>Consumer lag<\/li>\n\n\n\n<li>Consumer performance tuning<\/li>\n\n\n\n<li>Consumer configuration settings that matter in production<\/li>\n\n\n\n<li>Common consumer mistakes<\/li>\n\n\n\n<li>A production-ready consumer checklist<\/li>\n<\/ul>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">2. What Is a Kafka Consumer?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A&nbsp;<strong>Kafka Consumer<\/strong>&nbsp;is a client application that reads records from Kafka topics.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The simplest mental model is:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Producer\n   |\n   v\nKafka Topic\n   |\n   v\nConsumer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Vehicle\n   |\n   v\nProducer\n   |\n   v\nTopic: vehicle-telemetry\n   |\n   v\nAnalytics Consumer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer receives records such as:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>{\n  \"vehicle_id\": \"CAR-101\",\n  \"speed\": 88,\n  \"battery\": 72\n}\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">and performs business work:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Read event\n   |\n   v\nDeserialize\n   |\n   v\nValidate\n   |\n   v\nProcess\n   |\n   +--&gt; Update database\n   +--&gt; Generate alert\n   +--&gt; Calculate analytics\n   +--&gt; Call another service\n   |\n   v\nRecord progress\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">3. Kafka Consumers Pull Data<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka consumers normally&nbsp;<strong>pull<\/strong>&nbsp;records from brokers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka does not continuously push records into your application without the consumer asking.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer\n   |\n   | \"Give me records from Partition 2\n   |  beginning at my current position.\"\n   v\nBroker\n   |\n   v\nRecord batch\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer repeatedly calls:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>consumer.poll(...)\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Think of&nbsp;<code>poll()<\/code>&nbsp;as:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">\u201cGive me the records that I am currently allowed to read from my assigned partitions.\u201d<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">This pull design gives the consumer control over:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>How quickly it reads\nHow many records it handles\nWhen it pauses\nWhen it resumes\nHow it tracks progress\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">4. Consumer End-to-End Flow<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">At a high level:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Application starts\n      |\n      v\nCreate KafkaConsumer\n      |\n      v\nConnect to Kafka\n      |\n      v\nJoin Consumer Group\n      |\n      v\nReceive Partition Assignment\n      |\n      v\nFetch Records\n      |\n      v\nDeserialize\n      |\n      v\nProcess Business Logic\n      |\n      v\nCommit Offset\n      |\n      v\npoll() again\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This loop continues for the life of the consumer.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">5. What Does a Consumer Read?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A Kafka consumer receives&nbsp;<code>ConsumerRecord<\/code>&nbsp;objects.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A record includes important information such as:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic\nPartition\nOffset\nTimestamp\nKey\nValue\nHeaders\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic     = vehicle-telemetry\nPartition = 2\nOffset    = 1050\nKey       = CAR-101\nValue     = {\"speed\":88}\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This information matters because a consumer must know:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Where did the record come from?\nWhat order was it stored in?\nHow far have I progressed?\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">6. Deserialization<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">The producer serialized application data into bytes.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer must reverse that process.<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Kafka bytes\n    |\n    v\nDeserializer\n    |\n    v\nApplication Object\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">For example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>key.deserializer=org.apache.kafka.common.serialization.StringDeserializer\nvalue.deserializer=org.apache.kafka.common.serialization.StringDeserializer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">With Confluent Schema Registry, a consumer may instead use schema-aware deserializers for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Avro\nProtobuf\nJSON Schema\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A producer and consumer must agree on how data is encoded.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">7. What Is a Consumer Group?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A&nbsp;<strong>Consumer Group<\/strong>&nbsp;is a group of consumers cooperating to read one or more topics.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Consumers belong to the same logical group by using the same:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.id\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.id = order-processing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Consumers:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A\nConsumer B\nConsumer C\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Together:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer Group: order-processing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka treats them as workers sharing one consumption workload.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">8. Why Consumer Groups Exist<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose the topic receives:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>100,000 events \/ second\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">One consumer may not be able to process all of them.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">We want:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer 1\nConsumer 2\nConsumer 3\nConsumer 4\n...\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">to divide the work.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka accomplishes this by distributing&nbsp;<strong>partitions<\/strong>&nbsp;among consumers in the same group.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">9. Consumer Group + Partitions = Parallel Processing<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic: orders\n\nPartitions:\nP0\nP1\nP2\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">One consumer:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer Group\n\nConsumer A\n   |\n   +--&gt; P0\n   +--&gt; P1\n   +--&gt; P2\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">There is only one consumer process doing all the work.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Add three consumers:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 --&gt; Consumer A\nP1 --&gt; Consumer B\nP2 --&gt; Consumer C\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Now the group can process three partitions in parallel.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is one of Kafka&#8217;s most important scaling mechanisms.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">10. Golden Rule of Consumer Group Parallelism<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Inside a normal Kafka consumer group:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">One partition is assigned to at most one consumer in that group at a time.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Therefore:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Maximum useful partition-level consumer parallelism\napproximately equals\nnumber of partitions available to the group.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>6 partitions\n1 consumer\n\nConsumer 1 -&gt; P0 P1 P2 P3 P4 P5\n<\/code><\/pre>\n\n\n\n<pre class=\"wp-block-code\"><code>6 partitions\n3 consumers\n\nConsumer 1 -&gt; P0 P3\nConsumer 2 -&gt; P1 P4\nConsumer 3 -&gt; P2 P5\n<\/code><\/pre>\n\n\n\n<pre class=\"wp-block-code\"><code>6 partitions\n6 consumers\n\n1 partition per consumer\n<\/code><\/pre>\n\n\n\n<pre class=\"wp-block-code\"><code>6 partitions\n10 consumers\n\n6 consumers can receive partitions\n4 consumers have no partition work\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Adding consumers beyond the available partitions does not create additional partition-level parallelism.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">11. One Consumer Can Read Multiple Partitions<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Do not misunderstand the previous rule.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka does&nbsp;<strong>not<\/strong>&nbsp;require:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>1 consumer = 1 partition\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A consumer can own many partitions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>12 partitions\n3 consumers\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A possible assignment:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A -&gt; P0 P3 P6 P9\nConsumer B -&gt; P1 P4 P7 P10\nConsumer C -&gt; P2 P5 P8 P11\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This is normal.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">12. Why Partition Count Is So Important<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Topic partition count affects:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Producer parallelism\nBroker distribution\nConsumer parallelism\nScaling ceiling\nOrdering boundaries\nOperational overhead\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">For consumers:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Too few partitions\n        |\n        v\nCannot add useful consumers beyond partition count\n        |\n        v\nLimited parallel processing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">But:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Huge number of partitions\n        |\n        v\nMore metadata\nMore files\nMore leader\/replica management\nMore operational overhead\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Do not create partitions without capacity planning.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">13. How Kafka Balances Work Across Consumers<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka does not balance individual records directly.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">It balances&nbsp;<strong>partition ownership<\/strong>.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 P1 P2 P3 P4 P5\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">and:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A\nConsumer B\nConsumer C\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka determines an assignment such as:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A -&gt; P0 P3\nConsumer B -&gt; P1 P4\nConsumer C -&gt; P2 P5\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Each consumer fetches records from the leaders of its assigned partitions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">So load balancing is really:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partition Assignment\n        |\n        v\nConsumer Work Distribution\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">14. Important Limitation: Partitions May Not Have Equal Work<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka can distribute partitions reasonably, but it cannot magically guarantee equal CPU work.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Imagine:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 = 1,000 events\/sec\nP1 = 1,100 events\/sec\nP2 = 50,000 events\/sec\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A -&gt; P0\nConsumer B -&gt; P1\nConsumer C -&gt; P2\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Consumer C still has much more work.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Why?<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Because Kafka assigns&nbsp;<strong>partitions<\/strong>, not perfectly equal units of business computation.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is why producer key design and partition distribution affect consumer performance.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">15. Consumer Group ID<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\"><code>group.id<\/code>&nbsp;identifies the consumer group.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Different&nbsp;<code>group.id<\/code>&nbsp;values create independent consumption.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic: orders\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Group 1:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.id = analytics\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Group 2:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.id = fraud-detection\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Group 3:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.id = notifications\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">All three groups can independently read the same&nbsp;<code>orders<\/code>&nbsp;topic.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This becomes our&nbsp;<strong>fan-out pattern<\/strong>&nbsp;later.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">16. The Group Coordinator<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A Kafka broker acts as the&nbsp;<strong>Group Coordinator<\/strong>&nbsp;for a consumer group.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The coordinator manages important group responsibilities such as:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Group membership\nPartition assignment coordination\nConsumer liveness\nRebalance coordination\nOffset commits\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer Group\n      |\n      v\nGroup Coordinator\n      |\n      v\nKafka\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">17. How Kafka Finds the Group Coordinator<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka has an internal topic:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>__consumer_offsets\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer group&#8217;s&nbsp;<code>group.id<\/code>&nbsp;maps to a partition of this internal topic.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The broker that leads that offsets partition becomes the coordinator for the group.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.id\n   |\n   v\nhash \/ mapping\n   |\n   v\n__consumer_offsets partition\n   |\n   v\nLeader Broker\n   |\n   v\nGroup Coordinator\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This distributes consumer-group coordination across brokers.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">18. What Is Rebalancing?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A&nbsp;<strong>rebalance<\/strong>&nbsp;is the process of changing partition assignments among consumers in a group.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example before:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A -&gt; P0 P1\nConsumer B -&gt; P2 P3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Consumer C joins.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">After rebalance:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A -&gt; P0\nConsumer B -&gt; P1 P2\nConsumer C -&gt; P3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The exact assignment depends on the selected protocol and assignment strategy.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The key idea is:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">Group membership or subscribed-partition changes may require Kafka to redistribute partition ownership.<\/p>\n<\/blockquote>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">19. Why Rebalancing Exists<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Without rebalancing:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer B crashes\n      |\n      v\nP2 and P3 have no active consumer\n      |\n      v\nProcessing stops forever\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">With rebalancing:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer B crashes\n      |\n      v\nKafka detects membership change\n      |\n      v\nPartitions reassigned\n      |\n      v\nConsumer A \/ C take over\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Rebalancing is one of the mechanisms that gives consumer groups fault tolerance and elasticity.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">20. Major Rebalance Triggers<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A rebalance or reassignment can happen when:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>A consumer joins\nA consumer leaves gracefully\nA consumer crashes or becomes unreachable\nConsumer membership expires\nA subscribed topic gains partitions\nA subscribed topic changes\nA regex subscription matches a new topic\nSubscription metadata changes\nAdministrative\/manual group changes occur\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The exact protocol messages differ between the modern&nbsp;<strong>Consumer<\/strong>&nbsp;protocol and the older&nbsp;<strong>Classic<\/strong>&nbsp;protocol.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">21. Consumer Joins a Group<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Imagine:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Group currently:\n\nConsumer A\nConsumer B\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A new instance starts:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer C\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka must decide:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Which partitions should Consumer C receive?\n\nWhich existing consumers should give up partitions?\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">That causes a new assignment.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">22. Consumer Leaves Gracefully<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A consumer can shut down cleanly.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer B\n   |\n   v\nclose()\n   |\n   v\nLeaves group\n   |\n   v\nIts partitions become available\n   |\n   v\nAssignment updated\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Graceful shutdown is preferable to simply killing processes because the system can respond more cleanly.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">23. Consumer Crashes<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A crash is different.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer cannot politely say:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"I am leaving.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka must detect that the member is no longer healthy.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">After failure detection:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer B declared unavailable\n      |\n      v\nIts partitions reassigned\n      |\n      v\nAnother consumer resumes processing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This is consumer failover.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">24. Topic Partition Count Changes<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic: orders\nPartitions: 3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">becomes:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partitions: 6\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The group now has new work.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka needs to assign the new partitions to consumers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Therefore adding partitions can cause new partition assignments.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is why changing topic partitions is an operational event, not just a metadata edit.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">25. Modern Consumer Protocol vs Classic Protocol<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">This is extremely important for modern Kafka students.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka now has two consumer-group protocol families you may encounter:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>1. Consumer protocol\n2. Classic protocol\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The&nbsp;<strong>Consumer protocol<\/strong>&nbsp;is the newer generation of Kafka&#8217;s group rebalance protocol.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The&nbsp;<strong>Classic protocol<\/strong>&nbsp;is the older model that many existing tutorials and applications still describe.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">You should understand both because real organizations may run both during migration.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">26. Classic Consumer Group Protocol \u2014 Conceptual Model<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">In the Classic protocol, the group has a client-side leader involved in partition assignment.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A simplified flow is:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumers\n   |\n   v\nJoin Group\n   |\n   v\nCoordinator\n   |\n   v\nGroup leader participates in assignment\n   |\n   v\nSync Group\n   |\n   v\nConsumers receive assignments\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Classic assignment strategies include client-side assignors such as:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Range\nRound Robin\nSticky\nCooperative Sticky\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Some Classic strategies use eager movement of partitions; cooperative strategies can reduce how much work must stop and move.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">27. Modern Consumer Rebalance Protocol<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Modern Kafka&#8217;s Consumer protocol moves more group-assignment responsibility to the broker-side coordinator.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Key ideas:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Broker-side assignment\nIncremental assignment changes\nNo client group leader needed for assignment\nReduced global synchronization\nFaster \/ less disruptive rebalances\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The assignment strategy is managed differently from the older Classic client-side assignor model.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">For students:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">Do not memorize only the old \u201cone consumer becomes leader and assigns everything\u201d explanation as if it describes every modern Kafka consumer group.<\/p>\n<\/blockquote>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">28. Which Protocol Should Students Learn?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Learn the architecture in this order:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>First:\nConsumer Group concepts\nPartitions\nCoordinator\nOffsets\nRebalancing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Then:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Modern Consumer protocol\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Then:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Classic protocol\nfor compatibility and existing systems\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The business concepts are the same:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumers share partitions\nMembership changes\nAssignments change\nOffsets allow recovery\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The protocol mechanics differ.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">29. Rebalancing and Application Processing<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Rebalances matter because consumers may temporarily change ownership of partitions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Imagine:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A owns P0\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Then:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 is moved to Consumer B\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The application must correctly handle:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Finish or stop work on P0\nCommit safe progress\nRelease partition-specific resources\nBegin processing new assignment\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Incorrect rebalance handling can cause:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Duplicate work\nIncorrect offset commits\nLost application progress\nProcessing delays\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">30. ConsumerRebalanceListener<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">The Java consumer provides rebalance callbacks.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Important concepts:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>onPartitionsRevoked(...)\nonPartitionsAssigned(...)\nonPartitionsLost(...)\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">These callbacks let applications react when partition ownership changes.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Typical uses include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Commit safe offsets\nFlush local state\nClose partition-specific resources\nInitialize state for new partitions\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\"><code>onPartitionsLost<\/code>&nbsp;is especially important because in some failure cases the consumer may discover it has already lost ownership and should not assume it can safely perform the same actions as a normal revoke.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">31. What Is a Kafka Offset?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A Kafka&nbsp;<strong>offset<\/strong>&nbsp;is a numerical position inside a partition log.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partition 0\n\nOffset 0 -&gt; Record A\nOffset 1 -&gt; Record B\nOffset 2 -&gt; Record C\nOffset 3 -&gt; Record D\nOffset 4 -&gt; Record E\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Offsets belong to partitions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Therefore:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 offset 100\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">and:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P1 offset 100\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">are different positions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Offsets are not global message IDs across the entire topic.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">32. Consumer Position \u2014 Current Position<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka consumers have a&nbsp;<strong>current position<\/strong>.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A simple definition:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">The current position is the offset of the next record the consumer will fetch\/return from that partition.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partition:\n\n0 1 2 3 4 5 6 7\n            ^\n            |\nCurrent position = 6\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This means the consumer has already advanced past earlier records and expects to continue from around offset 6.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The current position normally lives in the running consumer&#8217;s state and advances as&nbsp;<code>poll()<\/code>&nbsp;returns records.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">33. Committed Offset<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">The&nbsp;<strong>committed offset<\/strong>&nbsp;is the saved recovery checkpoint for a consumer group.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partition:\n\n0 1 2 3 4 5 6 7\n        ^   ^\n        |   |\nCommitted Current\n   = 4     = 6\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer may currently have fetched through offset 5 but only safely committed progress up to the point represented by committed offset 4.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If the consumer crashes:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Current in-memory position is lost\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka uses:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Committed offset\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">to decide where the group should resume.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">34. Critical Offset Rule: Commit the NEXT Record to Read<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">This is one of the most important details in Kafka.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If you successfully process records:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Offset 10\nOffset 11\nOffset 12\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">your commit should normally indicate:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>13\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">because committed offset means:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">\u201cThe next record I intend to read\/process is offset 13.\u201d<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Do not think of the committed value merely as:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"The last record I processed.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Think:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"The next position I should resume from.\"\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">35. Where Are Committed Offsets Stored?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka stores consumer-group commits in the internal compacted topic:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>__consumer_offsets\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer\n   |\n   v\nOffset Commit\n   |\n   v\nGroup Coordinator\n   |\n   v\n__consumer_offsets\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This allows offsets to survive a consumer process restart.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">36. Why&nbsp;<code>__consumer_offsets<\/code>&nbsp;Is Compacted<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka usually only needs the latest committed position for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group\n+\ntopic\n+\npartition\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>order-processing \/ orders \/ P0 -&gt; 1050\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Later:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>order-processing \/ orders \/ P0 -&gt; 1100\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The newest commit becomes the important recovery state.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A compacted internal topic is a good fit for this kind of state.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">37. Four Positions Worth Understanding<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Students should eventually understand these concepts:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>1. Current consumer position\n2. Committed group offset\n3. High watermark \/ readable end\n4. Log end \/ latest written position\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Simplified:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Committed       Current          Readable End\n    |              |                 |\n    v              v                 v\n\n0 1 2 3 4 5 6 7 8 9 10 11 12 ...\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">In transactional consumption there is another important concept:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Last Stable Offset (LSO)\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">which matters for&nbsp;<code>read_committed<\/code>.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">We will cover transactions separately.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">38. What Is Consumer Lag?<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Consumer lag answers:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">\u201cHow far behind is the consumer group?\u201d<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Simplified:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>End of available data\n-\ngroup's consumed\/committed progress\n=\nlag\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Latest\/end position = 10,000\nCommitted position  = 9,500\n\nLag \u2248 500 records\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">High lag means work is accumulating faster than the consumer group is completing it.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">39. Lag Is a Symptom, Not the Root Cause<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">If lag grows, possible causes include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer processing too slowly\nToo few consumers\nToo few partitions\nHot partition\nSlow database\nSlow external API\nLarge records\nDeserialization cost\nLong GC pauses\nNetwork latency\nBroker throttling\nBad fetch configuration\nFrequent rebalances\nConsumer errors\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Do not respond to lag by blindly adding consumers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic has 4 partitions\nGroup already has 4 active consumers\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">adding ten more consumers will not create more partition-level parallelism.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">40. Offset Commit Strategies<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">There are two major approaches:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Automatic Commit\nManual Commit\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Neither is \u201calways correct.\u201d<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The right strategy depends on:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Business correctness\nFailure semantics\nProcessing time\nIdempotency\nThroughput\nOperational complexity\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">41. Automatic Offset Commit<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">With:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>enable.auto.commit=true\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">the consumer periodically commits offsets in the background.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The interval is controlled by:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>auto.commit.interval.ms\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This is easy to use.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">But it is important to understand what you are promising.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The application must ensure that records returned from polling are actually processed safely before later polling\/commit behavior can move the recovery point past work that is not finished.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Automatic commit is convenient, but it gives the application less explicit control over the exact business-processing boundary.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">42. Automatic Commit Does NOT Mean \u201cCommit Every Record\u201d<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A common misunderstanding:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>enable.auto.commit=true\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">does not mean:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Kafka commits every record immediately.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The commits happen periodically.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Therefore the gap between:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>current processing position\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">and:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>committed recovery position\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">can exist.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If a failure happens, some records may be read again.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is one reason duplicate-tolerant\/idempotent processing is valuable.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">43. Manual Offset Commit<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">With:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>enable.auto.commit=false\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">the application controls commits.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Common APIs include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>commitSync()\ncommitAsync()\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The core production rule is:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">Commit progress only when the corresponding business work is safely complete according to your required delivery semantics.<\/p>\n<\/blockquote>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">44.&nbsp;<code>commitSync()<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\"><code>commitSync()<\/code>&nbsp;waits for the commit result.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Process records\n      |\n      v\ncommitSync()\n      |\n      v\nWait for broker response\n      |\n      v\nContinue\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Advantages:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Simple failure handling\nKnown commit completion point\nGood at controlled boundaries\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Tradeoff:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Blocking can reduce throughput if used too frequently\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Do not call synchronous commit after every individual record unless the workload genuinely requires that behavior and the performance cost is acceptable.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">45.&nbsp;<code>commitAsync()<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\"><code>commitAsync()<\/code>&nbsp;sends the commit without blocking the main processing flow.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Process\n   |\n   v\ncommitAsync()\n   |\n   +----&gt; Commit request continues\n   |\n   v\nKeep processing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Advantages:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Lower blocking overhead\nBetter throughput\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">But you must handle:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Commit callback errors\nOrdering of commit attempts\nShutdown\/rebalance boundaries\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A common production pattern is:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Periodic async commits during normal processing\n\nPLUS\n\na final controlled sync commit at important shutdown\/rebalance boundaries\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">when that model matches the application&#8217;s semantics.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">46. Why Commit Timing Matters<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Consider two sequences.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Strategy A<\/h2>\n\n\n\n<pre class=\"wp-block-code\"><code>Commit offset\n     |\n     v\nProcess database write\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Failure scenario:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Commit succeeds\nApplication crashes\nDatabase write never happens\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka may resume&nbsp;<strong>after<\/strong>&nbsp;the record.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">From the business point of view:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>record may be lost\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This resembles&nbsp;<strong>at-most-once<\/strong>&nbsp;processing.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">47. Process Then Commit<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Now reverse it:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Process database write\n     |\n     v\nCommit offset\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Failure scenario:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Database write succeeds\nApplication crashes before commit\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">After restart:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Kafka resumes from old committed offset\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The record is processed again.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This creates possible duplicates.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This resembles&nbsp;<strong>at-least-once<\/strong>&nbsp;processing.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">48. At-Most-Once<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Commit\n   |\n   v\nProcess\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Potential outcome:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>No duplicate processing\nbut possible loss if failure happens after commit and before processing.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Use cases may include data where losing an occasional record is preferable to duplicate work.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">For important business events, this is often not desirable.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">49. At-Least-Once<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Process\n   |\n   v\nCommit\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Potential outcome:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>No intentional loss of successfully fetched work\nbut duplicates are possible after failures.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Therefore the downstream processing should preferably be:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Idempotent\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Meaning:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Processing the same event twice\ndoes not create an incorrect business result.\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">50. Exactly-Once \u2014 Be Precise<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Do not teach students:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">\u201cManual commits give exactly once.\u201d<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">They do not.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Exactly-once requires a coordinated processing design.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka supports transactional approaches for Kafka-to-Kafka workflows, together with transactional producers and consumers configured to read committed transactional data.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If your consumer writes to:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>External SQL database\nExternal REST API\nEmail provider\nPayment provider\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka cannot magically create one atomic transaction across all of those systems.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">You may need:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Idempotency keys\nDatabase transactions\nOutbox \/ inbox patterns\nDeduplication\nKafka transactions where applicable\nApplication-specific coordination\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Exactly-once is an architecture topic, not one checkbox.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">51.&nbsp;<code>auto.offset.reset<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">What happens when there is&nbsp;<strong>no usable committed offset<\/strong>?<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka uses:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>auto.offset.reset\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Common strategies include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>earliest\nlatest\nnone\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Modern Kafka also supports a duration-based reset option.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Important:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\"><code>auto.offset.reset<\/code>&nbsp;is not the normal \u201cwhere should my healthy consumer read next?\u201d setting.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">It is primarily used when Kafka cannot find a valid committed position for the group\/partition.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">52.&nbsp;<code>earliest<\/code><\/h1>\n\n\n\n<pre class=\"wp-block-code\"><code>auto.offset.reset=earliest\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">means:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">Start from the earliest retained offset available.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Useful for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>New analytics consumers\nBackfill processing\nReplay use cases\nTraining\nRebuilding state\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Remember:<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka may already have deleted older records according to retention.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">\u201cEarliest\u201d means earliest&nbsp;<strong>still retained<\/strong>, not necessarily the first record ever produced.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">53.&nbsp;<code>latest<\/code><\/h1>\n\n\n\n<pre class=\"wp-block-code\"><code>auto.offset.reset=latest\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">means:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">Start from the latest position and consume new records from that point forward.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Useful for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Real-time consumers\nwhere old history is not needed\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Danger:<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A new group configured with&nbsp;<code>latest<\/code>&nbsp;will normally not process the existing retained history before its starting point.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Know whether that is acceptable.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">54.&nbsp;<code>none<\/code><\/h1>\n\n\n\n<pre class=\"wp-block-code\"><code>auto.offset.reset=none\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">means:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">If Kafka has no valid committed offset, fail rather than automatically choosing earliest\/latest.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">This can be useful when silently choosing a new position would be dangerous.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">For critical applications, failing visibly can be safer than unexpectedly skipping or replaying data.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">55. Duration-Based Reset<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Modern Kafka can also reset based on a duration from the current time.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Start from approximately N hours\/days ago\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This is useful for some recovery or time-window processing cases.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Treat it as an advanced option after students master:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>earliest\nlatest\nnone\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">56. Right Offset Strategy \u2014 Decision Guide<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Ask:<\/p>\n\n\n\n<h3 class=\"wp-block-heading\">Question 1<\/h3>\n\n\n\n<p class=\"wp-block-paragraph\">Can duplicate processing happen safely?<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If yes:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>At-least-once\n+\ncommit after processing\n+\nidempotent downstream\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">is often a strong practical design.<\/p>\n\n\n\n<h3 class=\"wp-block-heading\">Question 2<\/h3>\n\n\n\n<p class=\"wp-block-paragraph\">Can any record be lost?<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If no:<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Avoid committing before successful processing.<\/p>\n\n\n\n<h3 class=\"wp-block-heading\">Question 3<\/h3>\n\n\n\n<p class=\"wp-block-paragraph\">What should a new group do?<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Choose intentionally:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>earliest\nlatest\nnone\nduration-based\n<\/code><\/pre>\n\n\n\n<h3 class=\"wp-block-heading\">Question 4<\/h3>\n\n\n\n<p class=\"wp-block-paragraph\">Can the business operation and offset be committed atomically?<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If Kafka-to-Kafka:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Kafka transactions may help.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If external side effects:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Use application\/database patterns appropriate to that system.\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">57. Consumer Failover and Recovery<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 -&gt; Consumer A\nP1 -&gt; Consumer B\nP2 -&gt; Consumer C\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Consumer B crashes.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Before:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>A -&gt; P0\nB -&gt; P1\nC -&gt; P2\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">After detection and reassignment:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>A -&gt; P0 P1\nC -&gt; P2\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">or another valid assignment.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The new owner starts from the group&#8217;s committed position for P1.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">That is why committed offsets are fundamental to recovery.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">58. Failure Window and Duplicate Processing<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Committed offset = 100\nCurrent position = 110\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer has processed records near 100-109 but has not committed the newer progress.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Then it crashes.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Replacement consumer resumes from:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>100\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Records between the committed point and the crash-time processing position may be processed again.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is expected for at-least-once designs.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Do not call this a Kafka bug.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">It is the consequence of the chosen commit boundary.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">59. Consumer Liveness<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka needs to distinguish:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Healthy consumer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">from:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Dead or stuck consumer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">There are two related concerns:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>1. Membership\/liveness heartbeat\n2. Application processing progress \/ poll interval\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">These concepts are important but the exact configuration differs between the modern Consumer protocol and the older Classic protocol.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">60. Heartbeats \u2014 Concept<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Consumers maintain group membership by demonstrating that they are alive.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer\n   |\n   v\nCoordinator\n\n\"I'm alive.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If Kafka determines the consumer is no longer alive:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Remove member\n      |\n      v\nReassign partitions\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">61. Classic vs Consumer Protocol Timeout Configuration<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Do not blindly copy old timeout advice.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Classic protocol<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Client-side settings historically include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>heartbeat.interval.ms\nsession.timeout.ms\n<\/code><\/pre>\n\n\n\n<h2 class=\"wp-block-heading\">Modern Consumer protocol<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Heartbeat\/session behavior is controlled more by broker-side group configuration.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is one of the important operational differences between the protocols.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Always check which group protocol your client and cluster are using before tuning these settings.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">62.&nbsp;<code>max.poll.interval.ms<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Another failure mode is:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer is alive\nbut application takes too long between poll() calls.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka uses:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>max.poll.interval.ms\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">to place an upper bound on how long the consumer can go without making expected poll progress under group management.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>poll()\n   |\n   v\nReceive 500 records\n   |\n   v\nProcessing takes 10 minutes\n   |\n   v\nmax.poll.interval.ms = 5 minutes\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The group may treat the consumer as unable to keep up and reassign its partitions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This can create a rebalance loop.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">63. Long Processing Is a Major Consumer Design Problem<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">If each record requires:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>10 seconds\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">and:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>max.poll.records=500\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">one poll could theoretically represent a huge amount of processing time.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Possible solutions:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Reduce max.poll.records\nOptimize business processing\nBatch downstream operations\nIncrease safe poll interval where justified\nPause\/resume partitions\nUse carefully designed worker-thread architecture\nIncrease partitions and consumers if appropriate\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Do not simply increase every timeout to enormous values.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">That can make real failures take too long to recover.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">64. Consumer Is Not Automatically Your Worker Thread Pool<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A very important Java design point:<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A&nbsp;<code>KafkaConsumer<\/code>&nbsp;instance is not intended to be freely used from many application threads concurrently.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A common safe model is:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>One consumer thread\n      |\n      v\npoll()\n      |\n      v\nprocess assigned records\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If you hand records to worker threads:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer Thread\n      |\n      v\nQueue\n   \/  |  \\\n  v   v   v\nW1  W2  W3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">you must solve:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Ordering\nBackpressure\nOffset coordination\nFailure handling\nRebalances\nShutdown\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This is an advanced architecture, not \u201cfree extra parallelism.\u201d<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">65. Fan-Out Consumption Pattern<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Fan-out means:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">Multiple independent applications\/groups consume the same topic.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>                   -&gt; Group: analytics\n                  \/\nTopic: orders -----&gt; Group: fraud\n                  \\\n                   -&gt; Group: notifications\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Each group gets its own independent view of the topic.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Inside each group, partitions are shared among that group&#8217;s consumers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This lets the same event power many independent systems.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">66. Fan-Out Example<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Topic:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>vehicle-telemetry\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Groups:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.id = realtime-dashboard\ngroup.id = anomaly-detection\ngroup.id = data-warehouse\ngroup.id = billing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">All four groups independently consume vehicle telemetry.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Producer does not need to send the event four times.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka retains the stream, and each group tracks its own offsets.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">67. Load-Balanced Consumption Pattern<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Load balancing means:<\/p>\n\n\n\n<blockquote class=\"wp-block-quote is-layout-flow wp-block-quote-is-layout-flow\">\n<p class=\"wp-block-paragraph\">Multiple consumers use the same&nbsp;<code>group.id<\/code>&nbsp;and divide partitions.<\/p>\n<\/blockquote>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic: orders\n\nConsumer Group: order-workers\n\nConsumer 1\nConsumer 2\nConsumer 3\nConsumer 4\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Partitions are distributed across those consumers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Use when:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>One logical application\nneeds more processing capacity.\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">68. Fan-Out vs Load Balance<\/h1>\n\n\n\n<h2 class=\"wp-block-heading\">Fan-Out<\/h2>\n\n\n\n<pre class=\"wp-block-code\"><code>Same topic\nDifferent group IDs\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Result:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Each group independently gets the data.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Use for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Analytics\nAlerts\nBilling\nFraud\nData warehouse\n<\/code><\/pre>\n\n\n\n<h2 class=\"wp-block-heading\">Load Balance<\/h2>\n\n\n\n<pre class=\"wp-block-code\"><code>Same topic\nSame group ID\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Result:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumers share the partitions\/work.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Use for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Scaling one application horizontally.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This distinction is foundational.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">69. Combining Fan-Out and Load Balancing<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Real production systems use both.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic: orders\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Analytics group:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>analytics\n  Consumer A\n  Consumer B\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Fraud group:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>fraud\n  Consumer A\n  Consumer B\n  Consumer C\n  Consumer D\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Notification group:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>notifications\n  Consumer A\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka allows each application to scale independently.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">70. Manual&nbsp;<code>assign()<\/code>&nbsp;vs Group&nbsp;<code>subscribe()<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Most consumer applications use:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>subscribe(...)\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">which participates in consumer-group partition management.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka can also let an application explicitly assign partitions:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>assign(...)\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer manually owns:\n\nP0\nP4\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">With manual assignment:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Kafka group rebalancing does not manage those partition assignments for you.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Use manual assignment only when you intentionally want application-controlled partition ownership.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Do not confuse it with normal Consumer Group load balancing.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">71. Consumer Fetching<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Consumers fetch records in batches.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Important settings include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>fetch.min.bytes\nfetch.max.wait.ms\nfetch.max.bytes\nmax.partition.fetch.bytes\nmax.poll.records\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">These influence:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Network efficiency\nThroughput\nLatency\nMemory\nAmount of work returned per poll\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">72.&nbsp;<code>fetch.min.bytes<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">This asks the broker to try to return at least a certain amount of data before responding, subject to wait limits.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Higher values may:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Increase batch efficiency\nReduce request overhead\nImprove throughput\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">but may also:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Increase latency when traffic is low\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Another throughput\/latency tradeoff.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">73.&nbsp;<code>fetch.max.wait.ms<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">This limits how long the broker may wait while trying to satisfy the fetch-size conditions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Think:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>fetch.min.bytes\n=\n\"Try to give me this much data.\"\n\nfetch.max.wait.ms\n=\n\"But don't wait longer than this.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">These settings work together.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">74.&nbsp;<code>fetch.max.bytes<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Controls approximately how much data a fetch response may contain across partitions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This matters for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Network usage\nConsumer memory\nLarge workloads\nLarge records\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Do not set fetch sizes blindly without understanding record sizes and available memory.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">75.&nbsp;<code>max.partition.fetch.bytes<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">This controls the amount of data fetched per partition in a request.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If messages\/batches are large, this setting becomes important.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Always coordinate large-message configuration end-to-end:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Producer\nTopic\/Broker\nConsumer\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">76.&nbsp;<code>max.poll.records<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">This controls how many records&nbsp;<code>poll()<\/code>&nbsp;returns to the application at one time.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Important:<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">It does not necessarily change how much data Kafka fetches internally from brokers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">It limits how many fetched records are returned to the application per poll.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is extremely useful for controlling:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Processing batch size\nPer-poll work\nmax.poll.interval risk\nApplication memory\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">77. Example:&nbsp;<code>max.poll.records<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Each record takes 100 ms to process\nmax.poll.records = 500\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Worst-case serial work:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>500 \u00d7 100 ms\n=\n50 seconds\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If each record instead takes:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>2 seconds\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">then:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>500 \u00d7 2 sec\n=\n1000 sec\n=\n16+ minutes\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Now poll timing becomes dangerous.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Tuning must reflect&nbsp;<strong>business processing time<\/strong>, not only Kafka throughput.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">78.&nbsp;<code>pause()<\/code>&nbsp;and&nbsp;<code>resume()<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Consumers can temporarily pause assigned partitions without giving up ownership.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 P1 P2 assigned\n\nP1 downstream overloaded\n        |\n        v\npause(P1)\n        |\n        v\nContinue P0\/P2\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Later:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>resume(P1)\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This can help implement controlled backpressure.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">But offsets, processing state and rebalance handling must still be correct.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">79. Static Membership \u2014 Concept<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">In some deployments, consumer instances are stable and restart temporarily.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka supports static membership concepts through:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.instance.id\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This can reduce unnecessary membership churn in certain restart scenarios.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">But it is not a replacement for:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Good failure detection\nCorrect rebalancing\nOffset correctness\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Use it when your deployment model benefits from stable member identities.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">80. Rebalance Cost<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A rebalance can cost:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Temporary reduction in consumption\nPartition movement\nCache warm-up\nState reconstruction\nDatabase\/client initialization\nDuplicate processing near commit boundaries\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Therefore:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Constant rebalancing\n=\nBad operational health\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Monitor rebalance frequency.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">81. Rebalance Storm Example<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Imagine:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer starts\n    |\n    v\nProcessing too slow\n    |\n    v\nmax.poll.interval exceeded\n    |\n    v\nConsumer removed\n    |\n    v\nRebalance\n    |\n    v\nConsumer rejoins\n    |\n    v\nRebalance again\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This can repeat.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Symptoms:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>High lag\nLow throughput\nFrequent assignment changes\nLots of logs\nConsumers appear healthy but never catch up\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Fix the root cause.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">82. Offset Commit During Rebalance<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">When a consumer is about to give up a partition, the application may need to commit the safest completed progress for that partition.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 currently owned by Consumer A\n      |\n      v\nP0 will move\n      |\n      v\nCommit safe completed position\n      |\n      v\nRelease resources\n      |\n      v\nConsumer B receives P0\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Never commit work that has not actually completed.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">83.&nbsp;<code>onPartitionsLost<\/code>&nbsp;vs Normal Revoke<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Normal revoke means the consumer is being told:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"You are about to give up these partitions.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A lost-partition callback means:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"You may already have lost ownership.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The second situation requires more caution.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Do not assume it is safe to commit arbitrary state after ownership has already been lost.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">84. Offset Commit Best Practices<\/h1>\n\n\n\n<h2 class=\"wp-block-heading\">Best Practice 1<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Commit only&nbsp;<strong>safe completed progress<\/strong>.<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Process\n   |\n   v\nBusiness operation succeeds\n   |\n   v\nCommit\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Best Practice 2<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Commit the&nbsp;<strong>next offset to read<\/strong>.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If records through offset 120 are complete:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>commit 121\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Best Practice 3<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Do not commit every record unless necessary.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Prefer reasonable batches:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Process N records\n      |\n      v\nCommit safe position\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Balance:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Duplicate window\nvs\ncommit overhead\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Best Practice 4<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Make downstream processing idempotent where possible.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Event ID = ORDER-123:PAYMENT-CAPTURED\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Database stores that event ID.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Second processing attempt:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Already processed\n-&gt; do not create duplicate payment\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This turns expected at-least-once replay into safe recovery.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Best Practice 5<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Handle commit errors.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Do not write:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>consumer.commitAsync();\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">and assume:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"Everything is definitely saved.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe failures where required.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Best Practice 6<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Handle rebalance boundaries.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Before giving up a partition:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Finish safe work\nCommit safe progress\nRelease resources\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">when your application design and protocol callback allow it.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Best Practice 7<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Monitor lag AND commit health.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A consumer can be processing quickly but failing commits.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Then:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Current position moves\nCommitted position stays behind\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A crash may cause large replay.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Monitor both.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">85. Consumer Failover Example<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Topic:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>payments\nP0 P1 P2 P3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Group:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>payment-risk\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Before:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A -&gt; P0 P1\nConsumer B -&gt; P2 P3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Consumer B crashes.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">After reassignment:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer A -&gt; P0 P1 P2 P3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P2 committed = 900\nP2 current before crash = 920\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">new processing resumes from committed progress around:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>900\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Records near 900-919 may be processed again.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">That is why payment processing must be idempotent.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">86. Fault Tolerance Depends on More Than Kafka<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka can reassign partitions.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">But your application may depend on:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Database\nHTTP service\nCache\nFile system\nMachine learning service\nPayment gateway\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If those are unhealthy:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer is technically alive\nbut business processing cannot complete.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Production health must include downstream dependencies.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">87. Consumer Performance Tuning \u2014 Start With Measurement<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Never start by changing ten configs.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Use:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>1. Define SLO\n2. Measure baseline\n3. Measure lag\n4. Measure processing time\n5. Identify bottleneck\n6. Change one controlled variable\n7. Load test\n8. Failure test\n9. Measure again\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Questions:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Is Kafka fetch slow?\nIs processing slow?\nIs database slow?\nAre partitions skewed?\nAre consumers rebalancing?\nAre commits slow?\nIs one partition hot?\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">88. Goal 1 \u2014 Increase Throughput<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Potential levers:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Increase useful consumer count\nIncrease topic partitions when architecture requires it\nIncrease fetch efficiency\nTune fetch.min.bytes\nTune fetch.max.wait.ms\nTune fetch sizes\nTune max.poll.records\nBatch database writes\nParallelize downstream work carefully\nReduce serialization\/deserialization overhead\nAvoid frequent commits\nAvoid frequent rebalances\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">But always protect:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Ordering\nCorrect commits\nMemory\nFailure recovery\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">89. Goal 2 \u2014 Lower Latency<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Potential levers:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Lower fetch waiting\nSmaller processing batches\nFast downstream services\nAppropriate consumer locality\nAvoid huge max.poll.records\nReduce rebalance frequency\nFast deserialization\nMonitor broker throttling\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Tradeoff:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Smaller batches\n=\nmore requests \/ overhead\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Again:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Latency\nvs\nThroughput\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">90. Goal 3 \u2014 Improve Reliability<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Focus on:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Correct offset commit timing\nIdempotent business processing\nSafe rebalance handling\nFailure testing\nIntentional auto.offset.reset\nCommit monitoring\nConsumer lag monitoring\nRetries in downstream operations\nDead-letter\/error strategy where appropriate\nTransactions where appropriate\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The consumer is reliable only if the&nbsp;<strong>business result<\/strong>&nbsp;is reliable.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">91. Goal 4 \u2014 Improve Availability<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Focus on:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Multiple consumer instances\nEnough partitions\nHealthy group coordination\nFast but sensible failure detection\nStable deployments\nControlled rolling restarts\nAvoid rebalance storms\nRedundant downstream services\nMonitoring and alerting\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Availability is not just:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"Consumer process is running.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">It means:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>The group continues making useful progress.\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">92. Consumer Configuration Study Map<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Important settings include:<\/p>\n\n\n\n<figure class=\"wp-block-table\"><table class=\"has-fixed-layout\"><thead><tr><th class=\"has-text-align-left\" data-align=\"left\">Setting<\/th><th class=\"has-text-align-left\" data-align=\"left\">What it controls<\/th><th class=\"has-text-align-left\" data-align=\"left\">Why it matters<\/th><\/tr><\/thead><tbody><tr><td><code>bootstrap.servers<\/code><\/td><td>Initial Kafka endpoints<\/td><td>Cluster discovery<\/td><\/tr><tr><td><code>group.id<\/code><\/td><td>Consumer group identity<\/td><td>Work sharing + offsets<\/td><\/tr><tr><td><code>key.deserializer<\/code><\/td><td>Key decoding<\/td><td>Correct application types<\/td><\/tr><tr><td><code>value.deserializer<\/code><\/td><td>Value decoding<\/td><td>Correct application data<\/td><\/tr><tr><td><code>enable.auto.commit<\/code><\/td><td>Automatic commits<\/td><td>Delivery semantics<\/td><\/tr><tr><td><code>auto.commit.interval.ms<\/code><\/td><td>Auto commit interval<\/td><td>Replay\/commit window<\/td><\/tr><tr><td><code>auto.offset.reset<\/code><\/td><td>Start point when no valid commit exists<\/td><td>Recovery\/new-group behavior<\/td><\/tr><tr><td><code>max.poll.records<\/code><\/td><td>Records returned per poll<\/td><td>Processing batch<\/td><\/tr><tr><td><code>max.poll.interval.ms<\/code><\/td><td>Max time between expected poll progress<\/td><td>Rebalance\/failure detection<\/td><\/tr><tr><td><code>fetch.min.bytes<\/code><\/td><td>Minimum fetch target<\/td><td>Throughput vs latency<\/td><\/tr><tr><td><code>fetch.max.wait.ms<\/code><\/td><td>Max broker fetch wait<\/td><td>Throughput vs latency<\/td><\/tr><tr><td><code>fetch.max.bytes<\/code><\/td><td>Fetch response size<\/td><td>Network\/memory<\/td><\/tr><tr><td><code>max.partition.fetch.bytes<\/code><\/td><td>Per-partition fetch amount<\/td><td>Large records\/memory<\/td><\/tr><tr><td><code>isolation.level<\/code><\/td><td>Transaction visibility<\/td><td>Exactly-once transactional reads<\/td><\/tr><tr><td><code>group.protocol<\/code><\/td><td>Group protocol selection<\/td><td>Modern vs Classic behavior<\/td><\/tr><\/tbody><\/table><\/figure>\n\n\n\n<p class=\"wp-block-paragraph\">Protocol-specific timeout\/assignment settings must be interpreted according to whether the group is using the modern Consumer protocol or Classic protocol.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">93. Modern Group Protocol Configuration Awareness<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">For modern Kafka, students should recognize:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>group.protocol=consumer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">as the setting associated with using the newer Consumer group protocol in clients that support it.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">With this protocol:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Assignment becomes broker-driven\nGroup behavior is more incremental\nSome Classic client-side assignment\/heartbeat settings are no longer used the same way\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Do not mix configuration advice from different group protocols.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">94. Classic Group Assignment Strategies<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">If working with Classic groups, common assignment strategies include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>RangeAssignor\nRoundRobinAssignor\nStickyAssignor\nCooperativeStickyAssignor\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Each has different distribution and rebalance behavior.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Do not say:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>\"Kafka always uses Round Robin.\"\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Assignment depends on protocol and configuration.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">95. Modern Consumer Protocol Assignment<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">With the newer Consumer protocol:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Assignment logic is broker-side\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The server provides available assignment strategies, and the consumer can participate according to the modern group protocol configuration.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The important teaching point is:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Classic:\nclient-led assignment mechanics\n\nModern Consumer protocol:\nbroker-side coordinator-driven assignment mechanics\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">96. Consumer Lag Monitoring<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Monitor lag per:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer Group\nTopic\nPartition\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Why partition-level?<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Because:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Overall group lag = 10,000\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">may hide:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>P0 lag = 0\nP1 lag = 0\nP2 lag = 10,000\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">That often indicates:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Hot partition\nSlow record type\nBad key distribution\nStuck processing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Partition-level visibility is essential.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">97. Metrics Worth Monitoring<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Examples of useful consumer-health signals:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Records consumed rate\nBytes consumed rate\nFetch latency\nFetch rate\nRecords per request\nAssigned partitions\nConsumer lag\nCommit latency\nCommit rate\nCommit failures\nRebalance rate\nTime since last successful poll\nProcessing latency\nApplication error rate\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka metrics alone are not enough.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Also measure:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Database latency\nAPI latency\nQueue depth\nCPU\nMemory\nGC\nThread pool saturation\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">98. Common Consumer Mistakes<\/h1>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 1 \u2014 More Consumers Than Partitions<\/h2>\n\n\n\n<pre class=\"wp-block-code\"><code>4 partitions\n20 consumers\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">You do not get twenty-way partition processing.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Many consumers remain idle.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 2 \u2014 Commit Before Processing<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Can create business data loss after a crash.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 3 \u2014 Never Commit<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Consumer restarts may replay huge amounts of work.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 4 \u2014 Auto Commit Without Understanding It<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Easy configuration can hide incorrect business semantics.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 5 \u2014 Slow Processing Between Polls<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Can cause consumer-group instability and repeated rebalances.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 6 \u2014 Long Blocking API Calls in Poll Thread<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>poll()\n  |\n  v\nCall external API\nwait 8 minutes\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Dangerous if poll progress requirements are exceeded.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 7 \u2014 Ignoring Rebalance Callbacks<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Can lose safe local state or commit the wrong progress.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 8 \u2014 Assuming Offset = Message ID<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Offset is a partition position.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 9 \u2014 Assuming Consumer Lag Means \u201cAdd Consumers\u201d<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">If the topic has too few partitions, extra consumers do nothing.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">If one partition is hot, extra consumers do not split that partition.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 10 \u2014 Confusing Different Consumer Groups<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Different groups do&nbsp;<strong>not<\/strong>&nbsp;share work with each other.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">They independently consume the topic.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 11 \u2014 Using One Group ID for Unrelated Applications<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">If analytics and notifications accidentally use the same&nbsp;<code>group.id<\/code>:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>They will split partitions\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">instead of both receiving all events.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This is a serious architecture mistake.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Mistake 12 \u2014 Ignoring Protocol Version<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Applying Classic timeout\/assignor advice to a modern Consumer-protocol group can be wrong.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Know which protocol you are operating.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">99. Java Consumer Example<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A simplified teaching example:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Properties props = new Properties();\n\nprops.put(\"bootstrap.servers\", \"&lt;BOOTSTRAP_SERVER&gt;\");\nprops.put(\"group.id\", \"vehicle-analytics\");\n\nprops.put(\n    \"key.deserializer\",\n    \"org.apache.kafka.common.serialization.StringDeserializer\"\n);\n\nprops.put(\n    \"value.deserializer\",\n    \"org.apache.kafka.common.serialization.StringDeserializer\"\n);\n\nprops.put(\"enable.auto.commit\", \"false\");\nprops.put(\"auto.offset.reset\", \"earliest\");\n\n<em>\/\/ For modern Kafka clients\/clusters where you intentionally use<\/em>\n<em>\/\/ the newer Consumer group protocol:<\/em>\n<em>\/\/ props.put(\"group.protocol\", \"consumer\");<\/em>\n\nKafkaConsumer&lt;String, String&gt; consumer =\n    new KafkaConsumer&lt;&gt;(props);\n\nconsumer.subscribe(List.of(\"vehicle-telemetry\"));\n\ntry {\n    while (true) {\n\n        ConsumerRecords&lt;String, String&gt; records =\n            consumer.poll(Duration.ofMillis(1000));\n\n        for (ConsumerRecord&lt;String, String&gt; record : records) {\n\n            System.out.println(\n                \"topic=\" + record.topic()\n                + \" partition=\" + record.partition()\n                + \" offset=\" + record.offset()\n                + \" key=\" + record.key()\n                + \" value=\" + record.value()\n            );\n\n            <em>\/\/ Business processing here<\/em>\n        }\n\n        <em>\/\/ Teaching example:<\/em>\n        <em>\/\/ commit after the batch is successfully processed.<\/em>\n        consumer.commitSync();\n    }\n}\nfinally {\n    consumer.close();\n}\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This example favors clarity.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A production application should additionally implement:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Error handling\nControlled shutdown\nRebalance handling\nMetrics\nRetry strategy\nIdempotent processing\nSecurity configuration\nAppropriate commit strategy\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">100. Confluent Kafka Cluster Connection<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">For a Confluent-managed cluster, the consumer also needs authentication\/security configuration.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Conceptually:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>bootstrap.servers=&lt;CONFLUENT_BOOTSTRAP&gt;\n\nsecurity.protocol=SASL_SSL\nsasl.mechanism=PLAIN\n\nsasl.jaas.config=&lt;CREDENTIAL_CONFIGURATION&gt;\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Then:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer\n   |\n   v\nAuthenticate\n   |\n   v\nDiscover brokers\n   |\n   v\nFind group coordinator\n   |\n   v\nJoin group\n   |\n   v\nReceive partition assignment\n   |\n   v\nFetch from partition leaders\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Keep credentials in secure configuration\/secrets, not hardcoded source code.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">101. Complete Consumer Journey<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Now combine everything:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>APPLICATION START\n      |\n      v\nCreate KafkaConsumer\n      |\n      v\nbootstrap.servers\n      |\n      v\nAuthentication\n      |\n      v\nCluster Metadata\n      |\n      v\ngroup.id\n      |\n      v\nFind Group Coordinator\n      |\n      v\nJoin Consumer Group\n      |\n      v\nGroup Protocol\n      |\n      v\nPartition Assignment\n      |\n      v\nFind Partition Leaders\n      |\n      v\nFetch Requests\n      |\n      v\nRecord Batches\n      |\n      v\nDeserialize Key \/ Value\n      |\n      v\npoll() returns ConsumerRecords\n      |\n      v\nBusiness Processing\n      |\n      v\nCurrent Position Advances\n      |\n      v\nCommit Safe Offset\n      |\n      v\n__consumer_offsets\n      |\n      v\npoll() again\n\nIf membership changes:\n      |\n      v\nRebalance \/ Reassignment\n      |\n      v\nPartitions Move\n      |\n      v\nResume From Committed Progress\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Students should be able to explain every box.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">102. Lab 1 \u2014 One Consumer, Multiple Partitions<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Create:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Topic: orders\nPartitions: 6\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Start:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>1 consumer\ngroup.id = order-workers\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>One consumer owns all six partitions.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Questions:<\/p>\n\n\n\n<ol class=\"wp-block-list\">\n<li>Can this application process partitions in parallel internally?<\/li>\n\n\n\n<li>What limits its throughput?<\/li>\n\n\n\n<li>What happens if the process crashes?<\/li>\n<\/ol>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">103. Lab 2 \u2014 Scale to Three Consumers<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Start:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>3 consumer instances\nsame group.id\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe the assignment.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Expected concept:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>6 partitions\n3 consumers\n\u2248 2 partitions per consumer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Questions:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Did a rebalance occur?\nWhich partitions moved?\nDid lag fall?\nDid throughput improve?\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">104. Lab 3 \u2014 More Consumers Than Partitions<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Keep:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>6 partitions\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Start:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>10 consumers\nsame group.id\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Some consumers receive no partitions.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This demonstrates:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partitions determine partition-level parallelism.\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">105. Lab 4 \u2014 Consumer Failure<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Start three consumers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Kill one abruptly.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Failure detection\nReassignment\nRemaining consumers take partitions\nConsumption resumes\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Record:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Rebalance duration\nLag increase\nDuplicate processing\nNew assignments\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">106. Lab 5 \u2014 Graceful Shutdown vs Crash<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Experiment A:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>consumer.close()\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Experiment B:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>kill -9 \/ abrupt container stop\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Compare:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Time to group recovery\nLogs\nRebalance behavior\nDuplicate window\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This teaches why controlled shutdown matters.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">107. Lab 6 \u2014 Automatic Commit<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Configure:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>enable.auto.commit=true\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Process records slowly.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Stop the consumer at different times.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Restart it.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Which records repeat?\nWhere does consumption restart?\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">The objective is not to prove \u201cauto commit is bad.\u201d<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The objective is to understand its failure window.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">108. Lab 7 \u2014 Manual Commit<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Configure:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>enable.auto.commit=false\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Flow:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>poll\nprocess all records successfully\ncommitSync\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Crash:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>after processing\nbefore commit\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe duplicates after restart.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Then explain:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Why duplicates are expected\nWhy idempotency matters\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">109. Lab 8 \u2014 Commit Before Processing<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">In a disposable training environment only:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>poll\ncommit\nthen process\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Force a crash after commit but before processing finishes.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer restarts after committed position\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Explain the potential lost-processing window.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">This demonstrates at-most-once semantics.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">110. Lab 9 \u2014 Fan-Out<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Topic:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>orders\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Create:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Group A = analytics\nGroup B = notifications\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Both groups consume the same events independently.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Observe separate committed offsets.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">111. Lab 10 \u2014 Load Balance<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Create:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Group = order-workers\nConsumers = 3\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partitions divided among consumers.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Compare with Lab 9.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Students should be able to explain:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Different groups = fan-out\nSame group = workload sharing\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">112. Lab 11 \u2014 Consumer Lag<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Make the consumer intentionally slow.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Observe lag growing.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Then test:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Add a consumer\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If enough partitions exist, lag recovery may improve.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Then create a scenario where one partition is hot.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Observe that adding consumers may not solve a single hot-partition bottleneck.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">113. Lab 12 \u2014&nbsp;<code>max.poll.records<\/code><\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Test:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>max.poll.records = 500\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Measure processing time.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Then:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>max.poll.records = 50\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Poll frequency\nProcessing batch size\nMemory\nLatency\nThroughput\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Explain why the right value depends on processing cost.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">114. Lab 13 \u2014 Rebalance During Processing<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Run several consumers.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">While processing:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Start another consumer.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Observe:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Partition assignment changes.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Add rebalance callbacks and print:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>revoked partitions\nassigned partitions\nlost partitions\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">This turns rebalancing from an abstract concept into something visible.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">115. Production Offset Strategy Example<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Suppose we process payment events.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Requirement:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Never charge twice.\nDo not silently lose a payment event.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A reasonable architecture might include:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>At-least-once Kafka consumption\n      |\n      v\nStable event\/payment ID\n      |\n      v\nIdempotent database\/payment operation\n      |\n      v\nCommit Kafka offset after safe business completion\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka offsets alone are not enough.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">The business operation itself must be designed for retries.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">116. Production Analytics Example<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Analytics event:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>page-view\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Duplicate processing may be less harmful.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">A different commit\/batching strategy may be acceptable.<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">Therefore:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Payment consumer configuration\n!=\nAnalytics consumer configuration\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">There is no universal \u201cbest Kafka consumer config.\u201d<\/p>\n\n\n\n<p class=\"wp-block-paragraph\">There is only:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Best config for this workload and correctness requirement.\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">117. Production Consumer Checklist<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Before calling a consumer production-ready:<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Group Design<\/h2>\n\n\n\n<ul class=\"wp-block-list\">\n<li>[ ] Intentional\u00a0<code>group.id<\/code><\/li>\n\n\n\n<li>[ ] Fan-out vs load-balancing pattern understood<\/li>\n\n\n\n<li>[ ] Consumer count matched to expected partition-level parallelism<\/li>\n\n\n\n<li>[ ] Partition count capacity planned<\/li>\n\n\n\n<li>[ ] Hot partitions considered<\/li>\n\n\n\n<li>[ ] Group protocol known: Consumer or Classic<\/li>\n<\/ul>\n\n\n\n<h2 class=\"wp-block-heading\">Offset Strategy<\/h2>\n\n\n\n<ul class=\"wp-block-list\">\n<li>[ ]\u00a0<code>enable.auto.commit<\/code>\u00a0chosen intentionally<\/li>\n\n\n\n<li>[ ] Commit boundary matches business correctness<\/li>\n\n\n\n<li>[ ] Committed value represents next offset to process<\/li>\n\n\n\n<li>[ ]\u00a0<code>auto.offset.reset<\/code>\u00a0chosen intentionally<\/li>\n\n\n\n<li>[ ] Duplicate processing behavior understood<\/li>\n\n\n\n<li>[ ] Data-loss window understood<\/li>\n\n\n\n<li>[ ] Commit failures monitored<\/li>\n<\/ul>\n\n\n\n<h2 class=\"wp-block-heading\">Processing<\/h2>\n\n\n\n<ul class=\"wp-block-list\">\n<li>[ ] Processing is idempotent where practical<\/li>\n\n\n\n<li>[ ] Slow downstream dependencies handled<\/li>\n\n\n\n<li>[ ] Retry strategy defined<\/li>\n\n\n\n<li>[ ] Dead-letter\/error handling defined if needed<\/li>\n\n\n\n<li>[ ] Long-processing design tested<\/li>\n\n\n\n<li>[ ] Ordering requirements documented<\/li>\n<\/ul>\n\n\n\n<h2 class=\"wp-block-heading\">Rebalancing<\/h2>\n\n\n\n<ul class=\"wp-block-list\">\n<li>[ ] Rebalance callbacks handled where needed<\/li>\n\n\n\n<li>[ ] Safe progress committed before normal revoke<\/li>\n\n\n\n<li>[ ] Lost-partition behavior understood<\/li>\n\n\n\n<li>[ ] Rolling deployment behavior tested<\/li>\n\n\n\n<li>[ ] Consumer crash tested<\/li>\n\n\n\n<li>[ ] Frequent rebalance alerting configured<\/li>\n<\/ul>\n\n\n\n<h2 class=\"wp-block-heading\">Performance<\/h2>\n\n\n\n<ul class=\"wp-block-list\">\n<li>[ ] Lag monitored<\/li>\n\n\n\n<li>[ ] Lag monitored per partition<\/li>\n\n\n\n<li>[ ]\u00a0<code>max.poll.records<\/code>\u00a0load tested<\/li>\n\n\n\n<li>[ ] Fetch settings load tested<\/li>\n\n\n\n<li>[ ] Memory measured<\/li>\n\n\n\n<li>[ ] CPU measured<\/li>\n\n\n\n<li>[ ] Deserialization cost measured<\/li>\n\n\n\n<li>[ ] Database\/API latency measured<\/li>\n\n\n\n<li>[ ] Consumer throughput measured<\/li>\n<\/ul>\n\n\n\n<h2 class=\"wp-block-heading\">Availability<\/h2>\n\n\n\n<ul class=\"wp-block-list\">\n<li>[ ] Multiple consumer instances used where required<\/li>\n\n\n\n<li>[ ] Failure detection behavior tested<\/li>\n\n\n\n<li>[ ] Consumers distributed appropriately<\/li>\n\n\n\n<li>[ ] Downstream dependencies have availability strategy<\/li>\n\n\n\n<li>[ ] Shutdown is graceful<\/li>\n\n\n\n<li>[ ] Restart\/recovery tested<\/li>\n<\/ul>\n\n\n\n<h2 class=\"wp-block-heading\">Security<\/h2>\n\n\n\n<ul class=\"wp-block-list\">\n<li>[ ] Credentials stored securely<\/li>\n\n\n\n<li>[ ] Topic READ permissions minimal<\/li>\n\n\n\n<li>[ ] Group permissions correctly scoped<\/li>\n\n\n\n<li>[ ] TLS\/SASL configuration validated<\/li>\n<\/ul>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">118. Troubleshooting Checklist<\/h1>\n\n\n\n<h2 class=\"wp-block-heading\">Symptom: Consumer Lag Growing<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Check:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer count\nPartition count\nHot partitions\nProcessing latency\nDatabase latency\nAPI latency\nFetch latency\nCPU\nMemory\nGC\nRebalances\nCommit latency\nBroker throttling\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Symptom: Constant Rebalances<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Check:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer crashes\nDeployment restarts\nmax.poll.interval\nLong processing\nMembership timeouts\nNetwork stability\nProtocol configuration\nRebalance callbacks\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Symptom: Duplicate Processing<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Check:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Crash after processing but before commit\nCommit frequency\nRetry behavior\nRebalance boundaries\nIdempotency\nOffset reset\nManual seeks\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">Duplicates are often expected in at-least-once systems.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Symptom: Missing Business Work<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Check immediately:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Were offsets committed before processing?\nWas auto.offset.reset=latest used for a new group?\nWas seek() used incorrectly?\nDid application drop errors?\nWas downstream work asynchronous and not awaited?\n<\/code><\/pre>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Symptom: Consumers Idle<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Check:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Number of partitions\nNumber of consumers in same group\nTopic subscription\nAuthorization\nAssignment\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">If:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>consumers &gt; partitions\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">idle members can be completely normal.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">119. Interview Questions<\/h1>\n\n\n\n<h2 class=\"wp-block-heading\">What is a Kafka consumer?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">A client application that fetches and processes records from Kafka topic partitions.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What is a Consumer Group?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">A set of consumers sharing the same group identity and cooperating to divide partition consumption.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Can two consumers in the same group read the same partition simultaneously?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Under normal group assignment, one partition is assigned to at most one consumer in that group at a time.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What if consumers are in different groups?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Each group consumes independently, enabling fan-out.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What limits consumer parallelism?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Primarily the number of partitions available to the group at the Kafka partition-assignment level.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What is a consumer offset?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">A numerical position in a partition used to track reading progress.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Current vs committed offset?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Current position is where the running consumer expects to read next. Committed offset is the saved recovery checkpoint for the group.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Where are group offsets stored?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka stores committed group offsets in the internal&nbsp;<code>__consumer_offsets<\/code>&nbsp;topic.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What is rebalancing?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">The process of changing partition ownership among consumers when group membership or subscribed partition metadata changes.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What happens when a consumer crashes?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Kafka detects the membership failure, reassigns its partitions, and replacement consumers resume using committed group progress.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What is fan-out?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Different consumer groups independently reading the same topic.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What is load balancing?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Consumers in the same group sharing partition ownership.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Why is commit timing important?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Because it determines whether failures may lead to duplicates or lost processing.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">At-least-once?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Process first, then commit; failures may cause reprocessing.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">At-most-once?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">Commit before processing; failures can cause processing to be skipped.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">Does manual commit guarantee exactly once?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">No. Exactly-once requires coordinated transaction\/idempotency design.<\/p>\n\n\n\n<h2 class=\"wp-block-heading\">What is consumer lag?<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">How far a consumer\/group is behind the available end of the partition stream.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">120. Master Mental Model<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">A beginner says:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer reads messages.\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A Kafka engineer says:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Consumer starts\n    |\n    v\nConnects to Kafka\n    |\n    v\nFinds Group Coordinator\n    |\n    v\nJoins Consumer Group\n    |\n    v\nParticipates in Group Protocol\n    |\n    v\nReceives Partition Assignment\n    |\n    v\nFinds Partition Leaders\n    |\n    v\nFetches Record Batches\n    |\n    v\nTracks Current Position\n    |\n    v\nDeserializes\n    |\n    v\nProcesses Business Logic\n    |\n    v\nCommits Safe Next Offset\n    |\n    v\n__consumer_offsets\n    |\n    v\nHandles Rebalances\n    |\n    v\nRecovers From Failures\n    |\n    v\nContinues Processing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">A production Kafka engineer asks:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Is partition distribution balanced?\n\nCan processing keep up?\n\nWhat is the lag?\n\nWhat happens when this consumer crashes?\n\nWhat happens if commit succeeds but processing fails?\n\nWhat happens if processing succeeds but commit fails?\n\nCan duplicate processing hurt us?\n\nCan any record be skipped?\n\nHow long does a rebalance take?\n\nAre rebalances happening too often?\n\nDo we have enough partitions?\n\nAre we using the correct group protocol?\n\nAre downstream systems idempotent?\n\nAre offsets and business state coordinated correctly?\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">That is Kafka consumer mastery.<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">121. Final Takeaways<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">Remember these twelve ideas:<\/p>\n\n\n\n<ol class=\"wp-block-list\">\n<li>Kafka consumers\u00a0<strong>pull<\/strong>\u00a0records from partition leaders.<\/li>\n\n\n\n<li>Consumers in the same group\u00a0<strong>share partitions<\/strong>.<\/li>\n\n\n\n<li>Different consumer groups independently consume the same data.<\/li>\n\n\n\n<li>Partition count defines the main ceiling for partition-level consumer parallelism.<\/li>\n\n\n\n<li>Kafka balances workload by assigning\u00a0<strong>partitions<\/strong>, not individual messages.<\/li>\n\n\n\n<li>Rebalancing changes partition ownership when membership or metadata changes.<\/li>\n\n\n\n<li>Modern Kafka has a newer Consumer group protocol as well as the older Classic protocol.<\/li>\n\n\n\n<li>Current position and committed offset are different.<\/li>\n\n\n\n<li>Committed offset represents the\u00a0<strong>next position to resume\/read<\/strong>.<\/li>\n\n\n\n<li>Commit timing determines duplicate-vs-loss behavior during failures.<\/li>\n\n\n\n<li>At-least-once designs should expect reprocessing and use idempotent business logic.<\/li>\n\n\n\n<li>Gold-standard consumer operations require monitoring\u00a0<strong>lag, commits, rebalances, processing time, and downstream health<\/strong>.<\/li>\n<\/ol>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h1 class=\"wp-block-heading\">122. Suggested Next Consumer Tutorials<\/h1>\n\n\n\n<p class=\"wp-block-paragraph\">After this foundation, continue with:<\/p>\n\n\n\n<pre class=\"wp-block-code\"><code>Kafka Consumer Internals\n        |\n        v\nConsumer Group Protocol Deep Dive\n        |\n        v\nClassic vs Modern Consumer Rebalancing\n        |\n        v\nPartition Assignment Strategies\n        |\n        v\nOffsets and Delivery Semantics\n        |\n        v\nAt-Least-Once \/ At-Most-Once \/ Exactly-Once\n        |\n        v\nConsumer Performance Tuning Lab\n        |\n        v\nConsumer Lag Monitoring\n        |\n        v\nRebalance Troubleshooting\n        |\n        v\nProduction Consumer Failure Testing\n<\/code><\/pre>\n\n\n\n<p class=\"wp-block-paragraph\">That path takes a student from \u201cI can read a Kafka message\u201d to \u201cI can design, operate and troubleshoot a production Kafka consumer system.\u201d<\/p>\n\n\n\n<hr class=\"wp-block-separator has-alpha-channel-opacity\"\/>\n\n\n\n<h2 class=\"wp-block-heading\">Technical Source Basis<\/h2>\n\n\n\n<p class=\"wp-block-paragraph\">This tutorial was checked against current Apache Kafka 4.3 consumer documentation\/Javadocs and current Confluent consumer documentation for consumer groups, offset tracking, group coordination, rebalancing, fetch behavior, commit APIs, and consumer lag. URLs are intentionally omitted from this learning document.<\/p>\n","protected":false},"excerpt":{"rendered":"<p>Consumer Groups, Parallelism, Offsets, Rebalancing, Failover, Consumption Patterns and Production Tuning Audience:&nbsp;Students and freshers with no previous Kafka experienceGoal:&nbsp;Start with \u201cWhat is a consumer?\u201d and finish with&#8230; <\/p>\n","protected":false},"author":1,"featured_media":0,"comment_status":"closed","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[1],"tags":[],"class_list":["post-1152","post","type-post","status-publish","format-standard","hentry","category-uncategorized"],"_links":{"self":[{"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/posts\/1152","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/comments?post=1152"}],"version-history":[{"count":1,"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/posts\/1152\/revisions"}],"predecessor-version":[{"id":1153,"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/posts\/1152\/revisions\/1153"}],"wp:attachment":[{"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/media?parent=1152"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/categories?post=1152"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/www.devopsschool.com\/tutorials\/wp-json\/wp\/v2\/tags?post=1152"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}