04 / 05

How should a Kafka system behave when downstream processing is slower than incoming traffic?

Difficulty: 8/10
Event-driven platform, Event contracts, Idempotency

Behavior When Downstream Processing Is Slower Than Incoming Traffic

When downstream processing is slower than incoming traffic, the Kafka system should use lag as a buffer and protect the downstream from overload. Kafka's log is a buffer: producers write to the log at their rate, and consumers read at their rate. If the consumer is slower, lag grows, but the producers are not affected because Kafka decouples them. This is a key benefit of Kafka: it absorbs bursts and allows the consumer to catch up later. The first behavior is to let lag grow as a buffer, up to the retention limit. The second behavior is to protect the downstream: if the downstream is slow because it is overloaded, the consumer should not hammer it with more requests. Use bounded concurrency (a fixed number of in-flight requests), backpressure (pause polling when the downstream is slow), and circuit breakers (stop calling the downstream if it is failing). The third behavior is to monitor and alert: lag growing beyond a threshold should trigger an alert, and the team should decide whether to scale the consumer, optimize the downstream, or accept the lag. The trade-off is between throughput and protection. Pushing the downstream harder may increase throughput in the short term but can cause failures and cascading outages; protecting the downstream reduces throughput but keeps the system stable.

The mechanism for using lag as a buffer is that Kafka retains records for the retention period (e.g., 7 days). If the consumer is slower than the producer, lag grows, but as long as the lag is within the retention window, the consumer can catch up when the downstream recovers. If the lag exceeds the retention, records are deleted, and the consumer cannot catch up without data loss. This is why monitoring lag and retention is critical: if lag is approaching the retention limit, you must act. The mechanism for protecting the downstream is backpressure. In Kafka, the consumer controls the pace: it can pause partitions, reduce max.poll.records, or introduce a delay between polls. If the downstream is slow, the consumer should slow down rather than buffer unbounded work in memory. A common pattern is to use a bounded queue between the poll loop and the processing threads: when the queue is full, the consumer pauses polling. This prevents memory exhaustion and protects the downstream. The trade-off is between latency and stability. A consumer that pushes the downstream hard has lower latency when the downstream is fast but causes failures when it is slow. A consumer that applies backpressure has higher latency but is more stable. Version note: Kafka does not provide built-in backpressure; you implement it in the consumer. Some frameworks, such as Spring Kafka, provide pause/resume APIs. Kafka Streams has its own backpressure mechanism based on the max.poll.interval and the processing rate.

A common mistake is to let the consumer buffer unbounded work in memory, which leads to OutOfMemoryError and crashes. Another mistake is to ignore the retention limit: if lag exceeds retention, the consumer loses data. A third mistake is to retry failed downstream calls in a tight loop, which overloads the downstream and makes the problem worse. The trade-off is between throughput and protection. For a system where the downstream is critical and fragile, protection is more important; for a system where the downstream can handle bursts, pushing harder may be acceptable. The right behavior depends on the downstream's characteristics and the business requirements. A good practice is to define a maximum lag threshold and a maximum processing rate, and to alert when either is exceeded. Version note: if you use Kafka Streams, the processing is bounded by max.poll.records and the state store; if the downstream is slow, the stream task will lag. You can use the num.stream.threads and max.poll.records to control the processing rate. For a custom consumer, you control the rate directly. In all cases, the key is to monitor lag and downstream health and to have a plan for when the downstream is slow.

javascript
  1. 1

    Kafka's log is a buffer: lag grows but producers are not affected.

  2. 2

    Lag can grow up to the retention limit; beyond that, data is lost.

  3. 3

    Protect the downstream with bounded concurrency, backpressure, and circuit breakers.

  4. 4

    Do not buffer unbounded work in memory; it causes OutOfMemoryError.

  5. 5

    Pause partitions or reduce max.poll.records to slow down the consumer.

  6. 6

    Monitor lag and retention; alert when lag approaches retention.

  7. 7

    Retrying failed downstream calls in a tight loop makes the problem worse.

  8. 8

    Define a maximum lag threshold and a plan for when the downstream is slow.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.