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.
"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.msto 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.bytesinstead, 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
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).
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."
The default for max.poll.records is 500.
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."
| Value | Behavior | Risk |
|---|---|---|
latest (default) | Start from the newest records (written after the consumer started) | SILENT DATA LOSS — everything between the lost offset and now is skipped |
earliest | Read all data in the partition from the very beginning | MASSIVE REPROCESSING — duplicates, possibly hours/days of it |
none | Throw an exception when attempting to consume from an invalid offset | Fails 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."
⚠ “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.”
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."
| Pool | T1-p0 | T1-p1 | T1-p2 | T2-p0 | T2-p1 | T2-p2 |
|---|---|---|---|---|---|---|
| Deal | C1 | C2 | C1 | C2 | C1 | C2 |
- 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
- "an assignment that is as balanced as possible"
- "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:
| Assignor | Balance | Rebalance cost | Notes |
|---|---|---|---|
RangeAssignor | Poor when consumers don't divide partitions evenly | Eager ("stop the world") | The default |
RoundRobinAssignor | Good (±1) when all subscribe to same topics | Eager | |
StickyAssignor | Good; better than RoundRobin for heterogeneous subscriptions | Eager, but minimizes movement | |
CooperativeStickyAssignor | Same as Sticky | Incremental — no stop-the-world | Requires 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."
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
| Config | Notes |
|---|---|
client.id | Any string; used by brokers to identify requests (e.g. fetch requests). Logging, metrics, and quotas |
group.instance.id | Any unique string → static group membership (§3) |
receive.buffer.bytes / send.buffer.bytes | TCP 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."
“It will behave like a BRAND-NEW consumer group with NO MEMORY of anything it consumed in the past.”
"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.
4.5 The poll loop
The timeout parameter controls how long poll() will block if data is not available in the consumer buffer.
4.7 Commits and offsets — where consumers actually break
Get this wrong by one and you reprocess one record per commit forever (harmless-ish) or skip one record per commit (data loss).