zgba Network

Kafka load testing: end-to-end latency, consumer lag and the clock trap

In a system built on Kafka, the bottleneck is often not the HTTP layer but the message path: producers slow down waiting for the broker’s acknowledgement, consumers fall behind, and consumer lag grows. A Kafka load test pushes that path at a realistic rate and answers three questions: At what latency do the brokers accept the throughput you need? Can the consumers read at the same rate? How does the number of messages left behind change as the load grows? Disclosure: I build Spitfire, a self-hosted load testing tool. The method below applies to any tool; the last section shows how Spitfire does it. What to measure Producer latency: the time from sending a message to the broker’s acknowledgement. It depends on acks: acks=all waits for every in-sync replica and is the slowest, acks=1 waits for the leader only, acks=0 does not wait and leaves nothing to measure. Use your production setting. End-to-end latency: the time from producing a record to consuming it. This is the delay users actually feel. While the producer’s ack time stays short, this one grows as consumers fall behind. Throughput: messages per second and, with the message size, bytes per second. 10,000 messages a second at 1 KB and the same count at 100 KB are very different loads. Consumer lag: the gap between the last offset a consumer group has committed in a partition and the end of that partition, in other words the messages not processed yet. A small, steady lag is normal; a lag that keeps growing under load means the consumers cannot keep up. Errors: rejected messages, timeouts, authorization failures and interruptions during a rebalance. The clock trap The obvious way to measure end-to-end latency is to put a timestamp in the message and subtract it from the time the consumer receives it. If the producer and the consumer run on different machines, that number is wrong: two machines’ clocks can differ by more than the latency you are trying to measure. Even with NTP, a few milliseconds of drift is normal, and your “p95 end-to-end latency” ends up showing the clock difference, not your system. Only compare times read from the same clock: produce and consume the same record on the same machine (and the same process, using a monotonic clock), and accept that you measure a sample of the records rather than all of them. Shaping the load On the producer side the target is usually a rate (“5,000 order events a second”), not a number of users. So instead of a fixed number of virtual users, use a constant arrival rate: the message rate does not drop when the broker slows down, and the latency shows as it is. In a closed model (fixed VUs), producers slow down with the broker and the problem hides as a quiet drop in throughput. On the consumer side the count matters: a consumer group runs at most as many active consumers as there are partitions. Put 20 consumers on a topic with 12 partitions and 8 of them sit idle. Match the consumer count to your production instances. Common mistakes Unrealistic messages. Sending the same small message every time flatters compression and caches. Produce messages whose size and content resemble production and change every time. Key distribution. Messages go to a partition by key. Few keys pile the load onto a few partitions (a hot partition); an empty key lets the partitioner spread messages. A shared topic. Do not write test messages to a topic that production consumers read. Use a separate topic and consumer group, and plan retention and cleanup up front. Average latency. Broker latency has a long tail: log segment rolls, replica sync and GC pauses cause rare but large spikes. Look at p95 and p99, not the average. A short test. To see whether lag grows, the load has to stay steady. Leave a ramp-up of a few minutes and a steady part of at least 10–15 minutes. Watching your real service’s lag The most useful lag in a load test is often not the test’s own consumer group but your real service’s group: does the service that processes these events fall behind under the load you generate, and by how much? You can read that without touching the service: Kafka’s admin API lists a partition’s end offset and a group’s committed offset, and lag is the difference. The load test never joins the group, commits or resets offsets, so watching it does not change how the service works. The principal needs Describe on the topic and on the group. With Spitfire In Spitfire a Kafka test is built in a web editor or written as JSON. The produce step waits for the broker’s acknowledgement (no linger), and the VUs of a runner share one producer client, like a real service. The consume step gives every VU its own consumer client; with a groupId they share the partitions like real consumer instances. For the two measurements above: Every produced record gets a spitfire-ts header. When a consume step reads a record stamped in the same run on the same runner, it records kafka_e2e_latency. Spitfire never compares two machines’ clocks; across N runners about 1/N of the records are measured, and the run page says so. A consume step with a fixed group watches that group’s lag, and any Kafka step can watch other groups with lagGroups, read-only, every 2 seconds, as kafka_consumer_lag. This test produces 500 order events a second for 15 minutes while 10 consumers read the topic, and watches the real order-service group: { “name”: “Kafka: produce and consume order events”, “scenarios”: [ { “name”: “producer”, “executor”: { “type”: “constant-arrival-rate”, “rate”: 500, “timeUnit”: “1s”, “duration”: “15m”, “preAllocatedVUs”: 20, “maxVUs”: 100 }, “steps”: [ { “id”: “produce”, “name”: “Produce order event”, “protocol”: “kafka”, “connection”: “kafka”, “kafka”: { “action”: “produce”, “topic”: “orders”, “key”: “order-{{$uuid}}”, “value”: ”{“orderId”:”{{$uuid}}”,“amount”:{{$randInt 10 5000}}}”, “lagGroups”: [ “order-service” ] } } ] }, { “name”: “consumer”, “executor”: { “type”: “constant-vus”, “vus”: 10, “duration”: “15m” }, “steps”: [ { “id”: “consume”, “name”: “Consume order event”, “protocol”: “kafka”, “connection”: “kafka”, “kafka”: { “action”: “consume”, “topic”: “orders”, “groupId”: “spitfire-loadtest”, “wait”: “5s” }, “checks”: [ { “type”: “jsonPath”, “path”: “$.orderId”, “op”: “exists” } ] } ] } ], “thresholds”: [ { “metric”: “req_duration”, “filter”: { “step”: “produce” }, “expr”: “p(95)<50” }, { “metric”: “req_failed”, “expr”: “rate<0.001” }, { “metric”: “kafka_e2e_latency”, “expr”: “p(95)<500” }, { “metric”: “kafka_consumer_lag”, “filter”: { “check”: “order-service” }, “expr”: “value<100” } ] } The run fails if the produce acknowledgement’s p95 goes over 50 ms, more than 0.1% of requests fail, end-to-end p95 goes over 500 ms, or the service’s group is left more than 100 records behind when the run ends. Spitfire installs on Docker or Kubernetes with one command and has a free edition. The full guide, with permissions, per-partition lag and how this lines up with your broker metrics, is here: Kafka load testing — Spitfire.

View original article