- Backpressure is a capacity-management problem, not simply a Kafka tuning problem.
- Bound consumer concurrency according to the safest downstream capacity, not the broker's maximum delivery rate.
- Consumer lag is a signal. Sustained lag growth, retries, pool saturation and downstream latency reveal the real bottleneck.
- Retries need backoff, jitter, finite limits and a replayable dead-letter strategy.
- Idempotency and correct offset handling are required because business processing can be attempted more than once.
Kafka can process enormous volumes of events. That strength also creates one of the most common production mistakes I see in event-driven systems: teams assume that because Kafka can deliver messages quickly, every downstream component can process them equally quickly.
It cannot. A consumer may be able to read 20,000 messages per second while a database safely sustains only a fraction of that write rate. An external API may have a contractual rate limit. A CPU-heavy enrichment step may become the real bottleneck. The gap between arrival rate and sustainable end-to-end processing rate is where backpressure begins.
Backpressure is not primarily a Kafka problem. It is a system-capacity coordination problem.
What backpressure actually means
Think about the complete path rather than the Kafka consumer in isolation:
If Kafka delivers 10,000 events per second but the complete business pipeline can safely sustain 2,500, the remaining work must be buffered somewhere. Kafka is excellent at being that durable buffer. Your application should not replace it with an uncontrolled in-memory queue.
Rule #1: Never allow uncontrolled concurrency
The first design control I look for is bounded concurrency. This is deliberately simple: the consumer accepts only a finite number of in-flight business operations at a time.
const CONCURRENCY = 20;
for (let i = 0; i < messages.length; i += CONCURRENCY) {
const batch = messages.slice(i, i + CONCURRENCY);
await Promise.all(
batch.map(message => this.processMessage(message)),
);
}In larger systems I normally prefer a worker pool, semaphore, queue or framework-level concurrency limiter rather than hand-written batching. The principle is the same: N is explicit and controlled.
How I choose concurrency in a real system
I do not choose 20, 50 or 100 because the number looks reasonable. I start from downstream capacity. If a database pool has 50 connections, other request traffic needs 20 and I want a safety margin of 10, then a Kafka worker count that can create 200 simultaneous database operations is obviously wrong.
The same reasoning applies to external APIs, Redis, CPU-bound transforms and file systems. The slowest critical dependency usually determines sustainable throughput.
Consumer lag is a health signal, not automatically a problem
Lag increasing during a short traffic burst is not necessarily unhealthy. A well-designed system can allow lag to grow temporarily and then catch up when the arrival rate drops. I become concerned when lag grows continuously while processing throughput remains below arrival rate.
| What I observe | Likely pressure point | First action |
|---|---|---|
| Lag ↑, CPU 90%+ | Consumer compute bottleneck | Profile processing, scale only after confirming partition capacity |
| Lag ↑, CPU low, DB pool 100% | Database bottleneck | Reduce concurrency, optimize queries/indexes and pool usage |
| Lag ↑, downstream latency ↑ | External dependency | Apply timeout, circuit breaker and retry-topic strategy |
| Memory ↑ rapidly | Unbounded in-flight work | Bound queue size and worker concurrency |
| DLQ ↑ suddenly | Bad data or dependency failure | Group failures by reason before replay |
Pause and resume consumption intentionally
When an internal queue crosses a high-water mark, pausing consumption is often safer than continuing to accumulate work. Resume only after it falls below a lower threshold so the consumer does not oscillate constantly.
const HIGH_WATER_MARK = 1_000;
const LOW_WATER_MARK = 400;
if (this.workQueue.size >= HIGH_WATER_MARK) {
consumer.pause([{ topic }]);
}
if (this.workQueue.size <= LOW_WATER_MARK) {
consumer.resume([{ topic }]);
}Retries must be bounded
Unlimited or immediate retries are one of the fastest ways to turn a dependency outage into a platform-wide incident. If 5,000 failed events each retry five times immediately, the application can produce 25,000 additional requests against a service that is already unhealthy.
I prefer exponential backoff with jitter, a finite attempt count, and a retry topic or scheduler that prevents the main consumer from blocking indefinitely.
function retryDelay(attempt: number, baseMs = 500): number {
const exponential = baseMs * Math.pow(2, attempt);
const jitter = Math.floor(Math.random() * 250);
return Math.min(exponential + jitter, 30_000);
}Dead Letter Queues are operational tools
A DLQ should not become a permanent garbage bin. Each failed event needs enough context to diagnose and replay it safely.
{
"eventId": "veh_983192",
"topic": "vehicle.events",
"eventType": "VehicleDetected",
"originalPayload": { "vehicleId": "V-291" },
"failureReason": "ValidationFailed",
"retryCount": 3,
"failedAt": "2026-08-15T12:10:00Z"
}Operations should be able to answer why the event failed, whether similar events are failing, whether the defect is fixed, and whether replay can create duplicate business actions.
Consumers must be idempotent
In practical at-least-once processing, the same business event can be attempted more than once. Idempotency is therefore a business requirement, not merely a Kafka setting.
async process(event: VehicleEvent): Promise<void> {
const exists = await this.processedEventRepository.exists({
where: { eventId: event.id },
});
if (exists) return;
await this.dataSource.transaction(async manager => {
await this.vehicleService.applyEvent(event, manager);
await manager.getRepository(ProcessedEvent).insert({
eventId: event.id,
processedAt: new Date(),
});
});
}Offset commits should reflect durable business completion
A useful mental model is: receive → validate → process → persist → commit. But offset management alone does not make business processing exactly-once. If persistence succeeds and the process crashes before the offset is committed, Kafka can deliver the event again. That is precisely why idempotency or a transactional pattern remains necessary.
Circuit breakers protect downstream dependencies
When repeated failures occur, the consumer should stop hammering the same failing dependency. A circuit breaker transitions from closed to open, waits, then allows limited half-open probes before normal traffic resumes.
Depending on the business case, events can move to a retry topic, be delayed, or be sent to the DLQ while the circuit remains open.
What I look for first during a production incident
When Kafka lag starts increasing, my first reaction is not to increase the replica count. I normally inspect the system in this order:
- Did the incoming event rate genuinely increase?
- Did per-event processing latency increase?
- Is CPU saturated?
- Is the database connection pool saturated?
- Are downstream APIs responding slowly?
- Did retry volume increase?
- Did partition distribution or consumer rebalancing change?
- Is one event type disproportionately expensive?
This sequence matters because it separates a consumer-compute problem from a dependency-capacity problem. The fix is different.
Production deployment controls
Kubernetes can help scale consumers, but resource requests, limits and graceful termination are just as important as replica count.
apiVersion: apps/v1
kind: Deployment
metadata:
name: vehicle-event-consumer
spec:
replicas: 3
selector:
matchLabels:
app: vehicle-event-consumer
template:
metadata:
labels:
app: vehicle-event-consumer
spec:
terminationGracePeriodSeconds: 45
containers:
- name: consumer
image: registry.example.com/vehicle-consumer:1.4.0
resources:
requests:
cpu: "250m"
memory: "512Mi"
limits:
cpu: "1"
memory: "1Gi"
env:
- name: KAFKA_MAX_IN_FLIGHT
value: "20"
- name: KAFKA_RETRY_LIMIT
value: "3"Metrics I expect on the dashboard
- consumer lag and lag growth rate
- events received and processed per second
- processing latency percentiles
- internal worker-queue depth
- retry and DLQ rates
- database pool usage and slow queries
- CPU, memory and Node.js event-loop lag
- downstream dependency latency and error rate
Five mistakes I would avoid
- Increasing replicas before finding the bottleneck. More consumers can amplify pressure on a saturated DB or API.
- Retrying immediately inside the main consumer. This creates retry storms and blocks useful work.
- Using unlimited Promise.all. Broker throughput becomes uncontrolled application concurrency.
- Committing before the business operation is durable. You can acknowledge work that never completed.
- Treating the DLQ as permanent storage. A DLQ needs ownership, monitoring, diagnosis and a replay procedure.
Final architecture checklist
- Bound in-flight work.
- Measure the slowest dependency.
- Use pause/resume or another load-shedding control.
- Make processing idempotent.
- Separate retry and DLQ paths.
- Protect remote dependencies with timeouts and circuit breakers.
- Commit only after durable business completion.
- Observe the full processing chain, not Kafka alone.
A system that remains controlled under overload is more valuable than one that wins a throughput benchmark under perfect conditions.
Your feedback helps prioritize deeper technical content.






