← All posts
Post 07August 22, 2026

Using Kafka topic lag as backpressure in an ingestion pipeline

When producers outrun consumers, durable topic lag can absorb bursts without blocking clients, as long as the broker has headroom. Measurements from log0 before and after fixing the ClickHouse writer, and under broker memory pressure.

Ashmit JaiSarita Gupta
Ashmit JaiSarita Gupta

Full-stack Software Engineer - (Builder of log0)

Using Kafka topic lag as backpressure in an ingestion pipeline

Post 5 showed a consumer that drained at 40 rows per second while the pipeline produced thousands. The obvious question is: where did all those extra rows go? They were not dropped, and the gateway never stopped returning 202. They went into the queue. A durable Kafka topic sitting between a fast producer and a slow consumer absorbed a 67,071-message backlog from a single burst, held every one of them on disk, and let the writer drain them at its own pace. Nothing was lost; detection fell behind and caught up. This post is about why that decoupling is the most useful property in the whole pipeline, and where it stops working.

This is post 7 in a series on building log0. Posts 4 through 6 were a connected arc: accept fast at the gateway (post 4), hit a wall at the ClickHouse writer (post 5), fix the writer with batching (post 6). This post steps back and looks at the thing that sat between the fast front and the slow back the entire time, and quietly made the wall survivable instead of fatal.


The mismatch has to go somewhere

Go back to the numbers from post 5. The gateway accepted thousands of requests per second. The naive ClickHouse writer drained about 40 rows per second. That is a producer running roughly 70 times faster than its consumer, sustained, under load.

In any system where the producer outpaces the consumer, the difference between the two rates does not disappear. It has to be accounted for somewhere. There are only three places it can go, and which one a system picks is one of the most consequential architecture decisions in its design:

  1. Drop it. Refuse the excess. The producer is told no, or its data is silently discarded.
  2. Block it. Make the producer wait for the consumer. The fast side is forced down to the speed of the slow side.
  3. Store it. Put the difference in a buffer and let the consumer work through it at its own pace.

log0 picks the third, and the buffer is a Kafka topic. That single choice is what turned Post 5 from an outage into a slow drain you could watch and then forget about.


Couple, or decouple

The first two options are the same mistake wearing two hats: they couple the producer to the consumer. If the gateway had to block until ClickHouse confirmed the write, then the gateway's throughput would be ClickHouse's throughput, 40 rows per second, and every log POST would crawl. If it dropped logs the writer could not keep up with, the system would lose data at exactly the moment something is going wrong and it is most needed. Either way, the slowest component in the chain sets the speed and the failure mode for everything in front of it.

A queue breaks that coupling.

The queue as a shock absorber. A fast ingestion-gateway returns 202 and writes into a durable raw-logs topic; a slow normalization-to-ClickHouse consumer drains it. Queue depth is the difference between the ingest rate in and the drain rate out, so the queue stores the mismatch instead of dropping it, turning a throughput problem into a lag problem. Three regimes: when out is greater than or equal to in the queue stays near empty and a 400-VU burst peaks at about 283 messages and never sustains a backlog; when out is less than in the durable log absorbs the burst and before batching a 300-VU burst peaked at 67,071, held on disk and drained at about 68 per second while the caller still got 202; and the absorber is finite, a sustained 3x800-VU load OOM-kills the broker at which point 202 finally stopsThe queue as a shock absorber. A fast ingestion-gateway returns 202 and writes into a durable raw-logs topic; a slow normalization-to-ClickHouse consumer drains it. Queue depth is the difference between the ingest rate in and the drain rate out, so the queue stores the mismatch instead of dropping it, turning a throughput problem into a lag problem. Three regimes: when out is greater than or equal to in the queue stays near empty and a 400-VU burst peaks at about 283 messages and never sustains a backlog; when out is less than in the durable log absorbs the burst and before batching a 300-VU burst peaked at 67,071, held on disk and drained at about 68 per second while the caller still got 202; and the absorber is finite, a sustained 3x800-VU load OOM-kills the broker at which point 202 finally stops

The gateway writes to the raw-logs topic and returns 202 the instant the producer accepts the record, as covered in post 4. It does not know or care whether normalization is keeping up. The consumer reads from raw-logs whenever it is ready. The two are connected only by a durable log, and the depth of that log is exactly the running difference between how fast events arrive and how fast they are drained. When the consumer is faster, the log stays near empty. When the consumer is slower, the log grows. Either way, the producer is never blocked and nothing is dropped.

That is the whole trick, and it is worth saying plainly: a queue converts a throughput problem into a lag problem. A throughput problem fails requests. A lag problem means the data is a little stale for a while. The second is enormously more recoverable than the first.


What the absorber looks like under a burst

This is measurable, and it is the same drain chart from post 5, now read for what it says about backpressure rather than about ClickHouse.

raw-logs consumer lag over time after a burst. Before batching, a 300-VU burst drives the lag almost vertically to a peak of 67,071 messages, which then drains slowly at about 68 messages per second. After batching, a 400-VU burst barely leaves the floor, peaking near 283 and never sustaining a backlograw-logs consumer lag over time after a burst. Before batching, a 300-VU burst drives the lag almost vertically to a peak of 67,071 messages, which then drains slowly at about 68 messages per second. After batching, a 400-VU burst barely leaves the floor, peaking near 283 and never sustaining a backlog

Look at the purple line again, but this time count what did not happen. A 300-VU burst drove the backlog to 67,071 messages. During all of that, the gateway kept returning 202, and not a single request was refused. The 67,071 messages were not a failure; they were the shock absorber doing its job, holding the burst on disk while the slow writer worked it off at about 68 messages a second. Detection lagged, the incidents took longer to appear, but every log that arrived was safe and eventually processed.

The green line is the post-batching world, and it is a different, slightly heavier burst: 400 VUs instead of 300. With the consumer now faster than the producer, even the heavier burst barely registers: the lag peaks at roughly 283 messages and never sustains a backlog, drifting back to zero within about fifteen seconds. The absorber is still there, still doing the same job; there is almost nothing to absorb, because out is now greater than in.

That contrast is the point, and it holds even though the two bursts are not identical, the after run is the larger one. The queue did not behave differently before and after batching. It absorbed the mismatch in both cases. What changed was the size of the mismatch. Backpressure is what let the system survive the bad version long enough for me to find and fix the writer, instead of falling over the moment the consumer fell behind.


Why nothing dropped: the producer never waits

The reason the gateway could keep accepting through a 67,000-message backlog is the non-blocking producer from post 4. Worth showing the one line that matters again:

java
kafkaTemplate.send(KafkaTopics.RAW_LOGS, event.getTenantId(), event)
    .whenComplete((result, ex) -> { /* on failure, route to raw-logs-dlq */ });

The producer hands the record to the broker and returns. It never blocks on the consumer, because it has no idea the consumer exists. The consumer's lag is invisible from the front door. A growing backlog on raw-logs does not slow down a single POST, because nothing on the request path reads that backlog. The broker accepts the write, the gateway returns 202, and the depth of the queue is somebody else's problem, specifically, the consumer's, who will get to it.

This is what makes "backpressure is a feature" literally true here. The pressure of a slow consumer is absorbed by the queue and never transmitted back to the client. The client experiences a system that always accepts, even while, three hops downstream, a writer is badly behind.


Where the absorber stops absorbing

A shock absorber has a travel limit, and pretending otherwise would be the dishonest version of this post. The queue is durable on disk, but the broker still holds the working set in memory, and memory is finite, especially on the single laptop these numbers come from.

I found the limit by looking for it. Under a deliberately abusive load, a sustained sweep of three rounds at 800 virtual users plus repeated bursts, the broker ran out of memory and the kernel killed it. Redpanda exited with code 137, OOMKilled=true. And here is the genuinely nasty part: the moment the broker is down, the gateway's producer can no longer resolve or reach it, so POST /api/v1/logs starts to hang, the producer blocks waiting for a broker that is gone, while /actuator/health still cheerfully returns 200 because it never checked broker connectivity in the first place. The accept-fast property and the zero-drops property both depend on the broker being up. When the absorber itself dies, every guarantee built on top of it dies with it, and the health check is the last to find out.

That failure is its own post (post 12, on running Redpanda on one laptop), but it belongs here too, because it is the honest boundary of this one. Backpressure-by-buffering is a feature right up until the buffer's host runs out of memory, and then it is a cliff. The lesson is not "queues make a system safe." It is "queues move the failure from the front door to the broker, so that is now the thing to watch."


What is not done

  • There is no backpressure signal to the client. This is the same gap named in post 4, and it is the direct cost of the absorb-everything design. log0 never returns 429; it absorbs until it cannot. A production system wants the gateway to push back when the broker is under memory pressure, shedding load deliberately rather than discovering the limit by being OOM-killed. Today the queue absorbs silently, which is great until it is not.
  • Absorbing silently means lag can grow unwatched. Converting a throughput problem into a lag problem only helps if someone is watching the lag. A backlog that drains is fine; a backlog that grows without bound is an outage in slow motion. The thing to monitor in this architecture is consumer lag and broker memory, not request success rate, because request success rate stays green right up until the broker dies.
  • The queue is shared across all tenants. raw-logs is one topic for everyone. The absorber does not know whose burst it is holding, which means one noisy tenant's flood and another tenant's trickle sit in the same buffer and drain through the same consumer. That is fine for total throughput and a real fairness problem for detection latency, and it is exactly where the next post goes.
  • All numbers are single-node (Docker Desktop, 512 MB per service, single Redpanda node, k6). The 67,071 absorb, the ~68/s drain, and the OOM cliff all characterize this configuration. More broker memory moves the cliff; it does not remove it.

Next: post 8, multi-tenancy as a security boundary, not a column.. The queue absorbs bursts, but it is one shared queue, so whose burst is it absorbing? The answer starts with a single decision, keying every record by tenantId, and runs all the way down to what isolates one customer's logs from another's when they share a topic, a database, and a consumer.


Try log0

log0 is the platform this series is built on, an open, multi-tenant incident pipeline you can run yourself or use hosted.

Written by Ashmit JaiSarita Gupta. Find me on LinkedIn, GitHub, and X, and read the rest of the series on Hashnode.

Ashmit JaiSarita Gupta

Full-stack Software Engineer and the builder of log0. I write about backend systems, distributed systems, and the physics-flavored corners of engineering.

← Back to all posts

Turn log chaos into incident clarity

Get started
log0© 2026 log0, Inc.