A rolling deploy goes out at noon, the consumer group rebalances three times, and by one o'clock the billing service is forty minutes behind.

Most scaling problems in an event-driven Java system start with Kafka partitions: how many a topic has, which key decides where a record lands, and what happens to them when consumers come and go. Each can be measured, and each fix can be checked against the same measurement.

How many Kafka partitions a topic needs

A consumer group spreads a topic's partitions across its members, and each partition is read by one member at a time. Sagar Deepak Joshi's InfoQ account of a Java contact center platform (June 2026) puts the consequence plainly: "the maximum number of active consumers per topic equals the number of partitions." Twelve partitions means the thirteenth pod sits idle, whatever your autoscaler thinks.

The sizing rule comes from Jun Rao of Confluent and still holds: with a target throughput t, and per-partition throughput p on the producer and c on the consumer, a topic needs at least max(t/p, t/c) partitions. The consumer side is almost always the limit, because c is your code: a listener that calls a database or a REST API per record processes hundreds of records a second, while Confluent puts the producer side at tens of megabytes a second per partition.

InputIllustrative valueHow to get it
Peak throughput t6,000 records/sProduction metrics at peak, times the growth you plan for over the next two years
Per-partition producer rate p20,000 records/skafka-producer-perf-test.sh with your record size and acks
Per-partition consumer rate c400 records/sYour real listener against a backlog, one consumer, one partition
Partitionsmax(6,000/20,000, 6,000/400) = 15, so 16 or moreRound up and leave headroom for catching up after an outage
# What one partition takes on the producer side, with your record size and acks
bin/kafka-producer-perf-test.sh --topic sizing-test --num-records 1000000 --record-size 1024 \
    --throughput -1 --producer-props bootstrap.servers=broker:9092 acks=all

# What your consumer code processes per second: run your real listener against a backlog
# and read the records-consumed-rate metric, or time it from the group's committed offsets

Leave the headroom at creation time, because the count only moves one way. We asked a Kafka 4.3 broker to shrink a six-partition topic to three, and it refused: "The topic orders currently has 6 partition(s); 3 would not be an increase." Adding partitions is allowed, but it changes which partition each key hashes to, so per-key ordering breaks for records produced around the change. Confluent's advice is to size for the throughput you expect in one or two years.

Keys, ordering and hot partitions

Kafka keeps order only within a partition, and the record key decides the partition. Key by customer and every event for one customer is processed in order; key by country and one partition gets most of your traffic. That hot partition shows up as lag on one partition while the others sit near zero, and adding consumers does nothing for it, because one partition can only have one reader in a group.

The fix is a key with more distinct values than you have partitions and no single value carrying a large share of the traffic. When the business needs ordering at a coarse level (all events of one large merchant), keep that merchant's events on one partition and make the consumer fast for it, or split the work into an ordered step and an unordered one.

Kafka rebalancing: what a deploy costs, and what KIP-848 changes

A rebalance moves partitions between members when one joins, leaves or stops sending heartbeats. A rolling deploy of six pods triggers one for every pod that stops and one for every pod that starts. There are three ways a Java consumer can take part:

  • The classic protocol with an eager assignor such as RangeAssignor: every member gives up all its partitions, and the group is reassigned from scratch.
  • The classic protocol with CooperativeStickyAssignor: only the partitions that move are revoked, over two rounds.
  • The new consumer protocol, generally available since Kafka 4.0 (KIP-848): the broker computes assignments with a "fully incremental design, which no longer relies on a global synchronization barrier". Clients opt in with group.protocol=consumer; the default is still classic.

An eager rebalance stops every partition in the group until all members have rejoined, so it waits for the slowest member to finish the batch it is processing. The incremental protocols keep the unmoved partitions running; only the partitions that change owner pause, until the old owner has given them up and the new one has received them.

Under the consumer protocol a member learns its assignment from its heartbeats, and the broker setting group.consumer.heartbeat.interval.ms defaults to 5 seconds. It is a broker setting that needs a restart, and the client-side heartbeat.interval.ms is not supported under the new protocol, so tuning the handover is a conversation with whoever runs the cluster.

Which protocol pauses less depends on group size, batch duration and those settings. Measure your own group under a rolling deploy, per partition, before and after switching.

For rolling restarts on the classic protocol, static membership avoids most of the churn: set group.instance.id to a stable name per pod and raise the session timeout. The Kafka documentation describes it as a way "to avoid group rebalances caused by transient unavailability (e.g. process restarts)".

Kafka consumer lag: measuring it and stopping one slow consumer from stalling the rest

Consumer lag is the distance between the newest offset in a partition and the last offset the group committed. It is the number to alert on, per partition, because the total hides a hot partition. The commands below ran against our test broker:

# Lag per partition: LOG-END-OFFSET minus CURRENT-OFFSET
bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group billing

GROUP    TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG   CONSUMER-ID  HOST  CLIENT-ID
billing  orders  0          10061           15066           5005  -            -     -
billing  orders  1          9589            15031           5442  -            -     -

# Is the group stable, rebalancing, or empty?
bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 --describe --group billing --state

# Skip a partition past a poison message (after --dry-run, with the group stopped)
bin/kafka-consumer-groups.sh --bootstrap-server broker:9092 --group billing \
    --topic orders:2 --reset-offsets --to-offset 9895 --execute

Lag that grows on every partition means the group is too slow for the input, and the answer is more partitions and consumers, or a faster listener. Lag that grows on one partition means a hot key or a poison message.

The cascade in Joshi's account came from the listener itself. Bulk provisioning of up to ten thousand agents went through a topic with three partitions, each record made a synchronous REST call, and "consumer lag built to over thirty minutes". The fix was to take slow work off the consumer thread: the listener wrote each request to a Redis queue and returned, a separate worker pool made the calls, and lag dropped by about half.

The other common stall is a record that fails every time. In Spring Kafka, a DefaultErrorHandler with a DeadLetterPublishingRecoverer retries a few times and then moves the record to a dead-letter topic, so the partition keeps moving:

@Configuration
class KafkaErrorHandling {

    // Three attempts one second apart, then the record goes to orders-dlt and the partition moves on.
    @Bean
    DefaultErrorHandler errorHandler(KafkaTemplate<Object, Object> template) {
        var recoverer = new DeadLetterPublishingRecoverer(template);
        var handler = new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 2L));
        handler.addNotRetryableExceptions(IllegalArgumentException.class);   // bad data will not get better
        return handler;
    }
}
spring.kafka.consumer.properties.group.protocol=consumer
spring.kafka.listener.concurrency=3

We ran this against the test broker with three messages, one of them unparseable. The two good ones were processed, and the bad one arrived on orders-demo-dlt with the exception message in its headers, while the listener used the new consumer protocol.

Where consumer state belongs, and how many schemas a stream needs

A consumer that keeps state in memory, such as a cache built from the events it has read, ties that state to the partitions it owns. Every rebalance moves partitions without their state. In Joshi's platform each pod kept its own cache, pods disagreed, and a cold start replayed events for about five minutes per pod, which made autoscaling useless. Moving the state to Redis, rebuilt from Kafka by a background thread, cut startup delay by sixty percent.

Schemas multiply in a similar way. Spoorthi Basu's InfoQ article (May 2026) describes a ride-sharing pipeline with four event types and three ride types, which became twelve schemas, twelve Schema Registry entries and twelve tables. Her fix is one schema per domain with discriminator fields (eventType, rideType) and nullable blocks for the variant data, with the registry set to FULL or FULL_TRANSITIVE compatibility so a new enum value is caught before it reaches a consumer.

Kafka exactly once, and what it does not cover

The producer is idempotent by default "if no conflicting configurations are set", so broker retries no longer write duplicates. Transactions extend this to a read, process and write cycle inside Kafka: consume from one topic, produce to another and commit the offsets atomically.

Neither covers the world outside Kafka. A payment API called from a listener, or a row written to a database, is not rolled back when the transaction aborts, and it runs again when the record is redelivered. That side needs an idempotent consumer: an operation ID in every record, checked before the side effect. Joshi's team used exactly that for provisioning calls.

Kafka performance checklist

  1. Size partitions with max(t/p, t/c), using a measured consumer rate, for the throughput expected in one to two years.
  2. Check the key distribution: no key value should carry a large share of the traffic.
  3. Alert on lag per partition; the group total hides a hot partition.
  4. Keep slow I/O off the consumer thread, or bound it with timeouts and a worker pool.
  5. Send records that fail repeatedly to a dead-letter topic instead of retrying forever.
  6. Measure your rebalance pause under a rolling deploy, then try the cooperative assignor, the consumer protocol and a shorter heartbeat, and keep what the numbers support.
  7. Use static membership on the classic protocol for pods that restart in place.
  8. Keep consumer state outside the consumer, or make it cheap to rebuild.
  9. Make every side effect idempotent with an operation ID.

How we approach Kafka performance work

We ran Kafka at the center of Recostream, the recommendation engine we built and operated from 2019 until GetResponse acquired it in 2022. Recommendation requests were processed asynchronously through Kafka, and the real-time engine returned a new recommendation in 20 to 30 ms under very high event volumes.

A performance engagement on an event-driven system starts with measurement: a reproducible load, lag and pause numbers per partition before any change, and the same numbers after each one. If your consumers fall behind at peak or after every deploy, our performance engineering service starts with that measurement, scoped on a thirty-minute call.