Learn Labs
4. Kafka Consumers: Reading Data from Kafka

4.6 Configuring consumers

Minimum data the consumer wants to receive from the broker per fetch.

6.1 Fetch sizing — the latency/efficiency dial

fetch.min.bytes (default 1 byte)

Minimum data the consumer wants to receive from the broker per fetch. "If a broker receives a request for records from a consumer but the new records amount to fewer bytes than fetch.min.bytes, the broker will wait until more messages are available before sending records back."

Why raise it:

  • "reduces the load on both the consumer and the broker, as they have to handle fewer back-and-forth messages in cases where the topics don't have much new activity (or for lower-activity hours of the day)"
  • "if the consumer is using too much CPU when there isn't much data available"
  • "to reduce load on the brokers when you have a large number of consumers"

Cost: "increasing this value can increase latency for low-throughput cases."

fetch.max.wait.ms (default 500 ms)

The companion cap — how long the broker will wait to satisfy fetch.min.bytes.

fetch.min.bytes = 1 MBthe broker has 1 MB of datafetch.max.wait.ms = 100100 ms have elapsedbroker respondsWHICHEVER HAPPENS FIRST
Figure 4.6.1fetch.max.wait.ms (default 500 ms)

"This results in up to 500 ms of extra latency in case there is not enough data flowing... If you want to limit the potential latency (usually due to SLAs controlling the maximum latency of the application), set fetch.max.wait.ms to a lower value."

fetch.max.bytes (default 50 MB) — prefer this one

Maximum bytes Kafka returns per poll of a broker. "Used to limit the size of memory that the consumer will use to store data returned from the server, irrespective of how many partitions or messages were returned."

The progress guarantee: "records are sent to the client in batches, and if the first record-batch that the broker has to send exceeds this size, the batch will be sent and the limit will be ignored. This guarantees that the consumer can continue making progress."

(Nice design detail: the limit is advisory when honoring it would deadlock the consumer. Contrast with Ch. 2's stuck-consumer bug, where a mismatched fetch.message.max.bytes genuinely does stall.)

Matching broker config exists too: "requests for large amounts of data can result in large reads from disk and long sends over the network, which can cause contention and increase load on the broker."

max.partition.fetch.bytes (default 1 MB) — usually the wrong knob

Max bytes the server returns per partition. When poll() returns ConsumerRecords, the object uses at most this per partition assigned to the consumer.

"controlling memory usage using this configuration can be quite complex, as you have no control over how many partitions will be included in the broker response. Therefore, we highly recommend using fetch.max.bytes instead, unless you have special reasons to try and process similar amounts of data from each partition."

max.partition.fetch.bytes = 1 MB

  • assigned 10 partitions → up to 10 MB
  • rebalance → assigned 40 partitions → up to 40 MB

► Your memory footprint changes when the GROUP SIZE changes. That's why fetch.max.bytes (an absolute cap) is preferred.

max.poll.records

"the maximum number of records that a single call to poll() will return. Use this to control the amount of data (but not the size of data) your application will need to process in one iteration of the poll loop."

Its real job is bounding the time between poll() calls, so max.poll.interval.ms isn't tripped — see below.

6.2 Liveness and timeouts

session.timeout.ms (default 10 s) + heartbeat.interval.ms
session.timeout.ms = 3 st = 0heartbeatheartbeatheartbeatheartbeat.interval.ms = 1 s
  • heartbeat.interval.ms = how OFTEN the consumer sends a heartbeat. session.timeout.ms = how LONG it can go WITHOUT sending one.
  • RULE: heartbeat.interval.ms < session.timeout.ms.
  • RULE OF THUMB: heartbeat.interval.ms = session.timeout.ms / 3 (session 3 s → heartbeat 1 s).
Figure 4.6.3session.timeout.ms (default 10 s) + heartbeat.interval.ms

The trade-off:

Effect
Lower session.timeout.ms"allows consumer groups to detect and recover from failure sooner but may also cause unwanted rebalances"
Higher session.timeout.ms"reduces the chance of accidental rebalance but also means it will take longer to detect a real failure"
max.poll.interval.ms (default 5 minutes)

The backstop for a hung main thread with a healthy heartbeat thread.

"The easiest way to know whether the consumer is still processing records is to check whether it is asking for more records. However, the intervals between requests for more records are difficult to predict and depend on the amount of available data, the type of processing done by the consumer, and sometimes on the latency of additional services."

How to size it: "It has to be an interval large enough that it will very rarely be reached by a healthy consumer but low enough to avoid significant impact from a hanging consumer."

What happens on breach: "the background thread will send a 'leave group' request to let the broker know the consumer is dead and the group must rebalance, and then stop sending heartbeats."

The relationship with max.poll.records: "In applications that need to do time-consuming processing on each record returned, max.poll.records is used to limit the amount of data returned and therefore limit the duration before the application is available to poll() again. Even with max.poll.records defined, the interval between calls to poll() is difficult to predict, and max.poll.interval.ms is used as a fail-safe or backstop."

slow per-record processinge.g. a 200 ms DB writemax.poll.records = 500100 seconds between poll() callsmax.poll.records = 50001000 seconds between poll() callswell under the default5-min max.poll.interval.ms — fineEVICTED FROM GROUPexceeds max.poll.interval.msrebalancenew owner reprocessesREBALANCE STORMpossibly evicted too

The default for max.poll.records is 500.

Figure 4.6.4max.poll.interval.ms (default 5 minutes)
default.api.timeout.ms (default 1 minute)

Applies to (almost) all API calls when you don't specify an explicit timeout. "Since it is higher than the request timeout default, it will include a retry when needed. The notable exception is the poll() method, that always requires an explicit timeout."

request.timeout.ms (default 30 s)

Max time the consumer waits for a broker response. On breach: "the client will assume the broker will not respond at all, close the connection, and attempt to reconnect."

"It is recommended NOT to lower it. It is important to leave the broker with enough time to process the request before giving up — there is little to gain by resending requests to an already overloaded broker, and the act of disconnecting and reconnecting adds even more overhead."*

That's an anti-thundering-herd argument: shortening this timeout under load makes an overloaded broker worse.

6.3 auto.offset.reset — the data-loss/replay switch

When it applies: the consumer starts reading a partition for which it has no committed offset, or the committed offset is invalid — "usually because the consumer was down for so long that the record with that offset was already aged out of the broker."

ValueBehaviorRisk
latest (default)Start from the newest records (written after the consumer started)SILENT DATA LOSS — everything between the lost offset and now is skipped
earliestRead all data in the partition from the very beginningMASSIVE REPROCESSING — duplicates, possibly hours/days of it
noneThrow an exception when attempting to consume from an invalid offsetFails loudly — often the right choice for correctness-critical apps

This config, combined with Ch. 2's retention discussion, is exactly how "we lost a day of data and nobody noticed" happens: a consumer down longer than retention + auto.offset.reset=latest = silent skip.

6.4 enable.auto.commit (default true)

"Set it to false if you prefer to control when offsets are committed, which is necessary to minimize duplicates and avoid missing data." Interval controlled by auto.commit.interval.ms. Full discussion in §7.

6.5 partition.assignment.strategy — the four assignors

Setup for all examples: consumers C1, C2 subscribed to topics T1, T2, each topic with 3 partitions.

Range (default — RangeAssignor)

"Assigns to each consumer a consecutive subset of partitions from each topic it subscribes to."

T1-p0T1-p1T1-p2T2-p0T2-p1T2-p2C14 partitionsC22 partitionsTopic T1Topic T2

⚠ “Because each topic has an uneven number of partitions and the assignment is done FOR EACH TOPIC INDEPENDENTLY, the first consumer ends up with MORE partitions than the second. This happens whenever Range assignment is used and the number of consumers does not divide the number of partitions in each topic neatly.”

Figure 4.6.5Range (default — RangeAssignor)

Note this is the default and it is systematically unbalanced. The imbalance compounds per topic.

RoundRobin (RoundRobinAssignor)

"Takes all the partitions from all subscribed topics and assigns them to consumers sequentially, one by one."

PoolT1-p0T1-p1T1-p2T2-p0T2-p1T2-p2
DealC1C2C1C2C1C2
  • C1: T1-p0, T1-p2, T2-p1 ← 3
  • C2: T1-p1, T2-p0, T2-p2 ← 3

“In general, if all consumers are subscribed to the same topics (a very common scenario), RoundRobin will end up with all consumers having the same number of partitions (or at most one partition difference).”

Sticky (StickyAssignor) — two goals
  1. "an assignment that is as balanced as possible"
  2. "in case of a rebalance, leave as many assignments as possible in place, minimizing the overhead associated with moving partition assignments from one consumer to another"

"In the common case where all consumers are subscribed to the same topic, the initial assignment from Sticky will be as balanced as RoundRobin. Subsequent assignments will be just as balanced but will reduce the number of partition movements. In cases where consumers in the same group subscribe to different topics, the assignment achieved by Sticky is more balanced than RoundRobin."

Cooperative Sticky (CooperativeStickyAssignor)

"Identical to Sticky but supports cooperative rebalances in which consumers can continue consuming from the partitions that are not reassigned."

⚠️ "if you are upgrading from a version older than 2.3, you'll need to follow a specific upgrade path in order to enable the cooperative sticky assignment strategy, so pay extra attention to the upgrade guide."

Summary table:

AssignorBalanceRebalance costNotes
RangeAssignorPoor when consumers don't divide partitions evenlyEager ("stop the world")The default
RoundRobinAssignorGood (±1) when all subscribe to same topicsEager
StickyAssignorGood; better than RoundRobin for heterogeneous subscriptionsEager, but minimizes movement
CooperativeStickyAssignorSame as StickyIncremental — no stop-the-worldRequires a careful upgrade path from <2.3

You can also implement your own and point partition.assignment.strategy at your class.

6.6 client.rack — fetch from the closest replica

"By default, consumers will fetch messages from the leader replica of each partition. However, when the cluster spans multiple datacenters or multiple cloud availability zones, there are advantages both in performance and in cost to fetching from a replica located in the same zone as the consumer."

Default — leader-only fetch
cross-AZ traffic $$$P0 LEADERAZ-aP0 followerAZ-bconsumerin AZ-bthe follower in its own AZ goes unused
With client.rack + RackAwareReplicaSelector
P0 LEADERAZ-aconsumerin AZ-aP0 followerAZ-bconsumerin AZ-blocal! cheap, fast
Figure 4.6.76.6 client.rack — fetch from the closest replica

Two-part setup — you need BOTH:

Client side
client.rack = <the zone the client is in>
Broker side
replica.selector.class = org.apache.kafka.common.replica.RackAwareReplicaSelector

You can also "implement your own replica.selector.class with custom logic for choosing the best replica to consume from, based on client metadata and partition metadata."

(Recall Ch. 1: consumers MAY fetch from a follower; producers MUST use the leader. This config is what cashes in that asymmetry.)

6.7 The rest

ConfigNotes
client.idAny string; used by brokers to identify requests (e.g. fetch requests). Logging, metrics, and quotas
group.instance.idAny unique string → static group membership (§3)
receive.buffer.bytes / send.buffer.bytesTCP socket buffers; -1 = OS defaults. "It can be a good idea to increase these when producers or consumers communicate with brokers in a different datacenter" (higher latency, lower bandwidth)

6.8 ⚠️ offsets.retention.minutes — a BROKER config that silently resets your consumers

"This is a broker configuration, but it is important to be aware of it due to its impact on consumer behavior."

Group HAS ACTIVE MEMBERSactively heartbeatingcommitted offsets RETAINEDretrievable on reassignment or restartGroup BECOMES EMPTYkept only for offsets.retention.minutesoffsets are DELETEDDEFAULT: 7 DAYSGroup becomes active againa BRAND-NEW group with NO MEMORYauto.offset.reset decideslatestskip everything that arrivedearliestreprocess the entire retained log

“It will behave like a BRAND-NEW consumer group with NO MEMORY of anything it consumed in the past.”

Figure 4.6.86.8 ⚠️ offsets.retention.minutes — a BROKER config that silently resets your consumers

"Note that this behavior changed a few times, so if you use versions older than 2.1.0, check the documentation for your version for the expected behavior."

Real scenario: a batch job's consumer group runs monthly. Between runs the group is empty for ~30 days > 7 days → offsets vanish → next run either reprocesses everything or skips everything. This is a classic and very confusing production surprise.


On this page