Kafka load testing: producers, consumers, throughput, end-to-end latency and consumer lag
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 looks for answers to three questions: at what latency do the brokers accept the throughput you need, can the consumers read at the same rate, and how does the number of messages left behind change as the load grows.
What are we measuring?
- Producer latency: the time from sending a message to the broker's acknowledgement (ack). It depends on
acks:acks=allwaits for every in-sync replica and is the slowest,acks=1waits for the leader only,acks=0does not wait and leaves nothing to measure. Use the production setting in the test. - 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. Mind the clocks: two machines' clocks can differ by more than the latency you measure, so a measurement that reads the send time on one machine and the receive time on another shows the clock difference, not the milliseconds.
- Throughput: messages produced and consumed 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 read 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: messages the broker rejects, timeouts, authorization failures, and interruptions during a rebalance.
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 fixing the number of virtual users (VUs), 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 (constant VUs) producers slow down with the broker and the problem hides as a quiet drop in throughput (load test types).
On the consumer side the count matters: a consumer group runs at most as many consumers at once 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; if you want to test scaling, take the partition count into account too.
Common mistakes
- Unrealistic messages. Sending the same small message every time flatters compression and caches. Produce messages whose size and content resemble production and that change every time.
- Key distribution. Messages go to a partition by key. Using few keys piles the load onto a few partitions (a hot partition); an empty key lets the partitioner spread messages. Stay close to the key variety of production.
- A shared topic. Do not write test messages to a topic that production consumers read. Use a separate topic and consumer group for the test, 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 (the p95 and p99 guide).
- A short test. To see whether lag grows, the load has to stay steady for a few minutes. Leave a ramp-up of a few minutes and a steady part of at least 10–15 minutes.
With Spitfire
First add a Kafka connection under Connections: the broker list, a client id, SASL if needed (PLAIN, SCRAM-SHA-256, SCRAM-SHA-512) and TLS, and acks (all by default; leader or none). Passwords are stored encrypted and the test names only the connection; the Test button checks the brokers are reachable. In the test the Kafka step has two actions:
- produce: sends a message with a topic, key, value and headers, and waits for the broker's acknowledgement. The step's duration is the time to that acknowledgement; Spitfire does not wait to fill a batch (no linger), though messages of concurrent VUs are still grouped into one request. The VUs of a runner share one producer client per connection, like a real service. Values can use
{{$uuid}},{{$randInt 10 5000}},{{$timestamp}}or variables from a CSV data file; an empty key spreads messages across partitions. - consume: every VU keeps its own consumer client and reads one message per step. With a
groupIdthe VUs share the partitions like real consumer instances; a new group starts from the beginning (it reads the backlog too) and an existing group resumes from its committed offsets. Without a group, only messages that arrive during the test are read. The step's duration is the wait for the next message; when the wait (10 s by default) runs out, the step fails with a timeout. The message value goes to checks and extraction (with JSONPath when it is JSON), and topic, partition, offset, timestamp and key can be read as headers.
A test that produces 500 order events a second while 10 consumers read the same topic; the produce acknowledgement's p95 must stay under 50 ms and the error rate under 0.1%. The full file is on the examples page.
{
"name": "Kafka: sipariş olayı üret ve tüket",
"scenarios": [
{ "name": "uretici",
"executor": { "type": "constant-arrival-rate", "rate": 500, "timeUnit": "1s",
"duration": "5m", "preAllocatedVUs": 20, "maxVUs": 100 },
"steps": [ { "id": "produce", "protocol": "kafka", "connection": "kafka",
"kafka": { "action": "produce", "topic": "orders", "key": "order-{{$uuid}}",
"value": "{\"orderId\":\"{{$uuid}}\",\"amount\":{{$randInt 10 5000}}}" } } ] },
{ "name": "tuketici",
"executor": { "type": "constant-vus", "vus": 10, "duration": "5m" },
"steps": [ { "id": "consume", "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" }
]
}During the run you watch per-step requests per second (the produce and consume rates), p95, p99, total data sent and received, and error kinds (authorization, refused connection, broker rejection, timeout) live. The consume step's duration is not end-to-end latency: when the consumers are not behind, the record is already waiting and the duration is short; when they are behind, you still only see the time to take the next record. For end-to-end latency and consumer lag Spitfire takes two more measurements 0.17.0+.
End-to-end latency
Every record a produce step sends gets a spitfire-ts header (the send time and its source). When a consume step reads a record stamped in the same run on the same runner, it records the time from send to consume as kafka_e2e_latency (ms; p50, p95, p99). The produce step's Stamp records option turns the stamp off ("noStamp": true); a header of the same name you set yourself is left alone.
Why only the same runner? Because of the clock difference above, Spitfire never compares two machines' clocks: a record is measured when the runner that produced it also consumes it (both ends read the same monotonic clock). As a result, in a run spread over N runners about 1/N of the records are measured, and the run page says so. Records produced by other tools, other runs or other runners are consumed normally but not measured. Keep the produce and consume steps in the same test; your real service's consumers are outside this measurement, so use the lag below for them.
Consumer lag: the test's group and your real service's group
A step that consumes with a fixed consumer group watches that group's lag on its topic (Watch this consumer group's lag; "noLag": true turns it off). Any Kafka step can also watch other groups with lagGroups; the most useful is your real service's group: you see whether, and by how much, your service falls behind under the load the test produces. Lag is the partition's end offset minus the group's committed offset (in records); a partition without a commit counts every record it holds. It is reported per group as kafka_consumer_lag and per partition, labelled group/partition, as kafka_consumer_lag_partition.
- Read-only. Lag is read with Kafka's admin API: it lists offsets and committed offsets; Spitfire never joins the group, commits or resets offsets. Watching your real service's group does not change how it works.
- Every 2 seconds, from one runner. One runner per run (the one holding the first share of the load; the CLI itself for local runs) polls every 2 seconds, so each value is counted once. Swings shorter than that may not show.
- Permissions. The connection's principal needs Describe on the topic (metadata, ListOffsets) and Describe on every watched consumer group (OffsetFetch). Without them the run goes on, logs one warning, and lag is not reported.
- Fixed names. Topic and group must be fixed names, not templates (lag is not watched with
{{topic}}). A group watched by two steps is reported once; per-partition series stop at 64 partitions per group, while the total still covers them all.
The produce step of the test above also watches the real service's order-service group; the thresholds cap the p95 of end-to-end latency, the highest lag of any watched group, and the lag left in the service's group when the run ends. The test's own spitfire-loadtest group is already watched from the consume step. The test as changed was checked with spitfire validate.
{ "id": "produce", "name": "Sipariş olayı üret", "protocol": "kafka", "connection": "kafka",
"kafka": { "action": "produce", "topic": "orders", "key": "order-{{$uuid}}",
"value": "{\"orderId\":\"{{$uuid}}\",\"amount\":{{$randInt 10 5000}}}",
"lagGroups": [ "order-service" ] } }"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", "expr": "max<10000" },
{ "metric": "kafka_consumer_lag", "filter": { "check": "order-service" }, "expr": "value<100" }
]In thresholds, max checks that no watched group ever lagged more, value the lag left when the run ended; "filter": { "check": "order-service" } narrows the threshold to one group. The run page's Kafka card shows the end-to-end percentiles second by second and each group's lag live through the run; once the run is stored, the final and peak lag are added with the per-partition breakdown. Both are in reports, the CLI summary and comparisons (end-to-end p95, peak lag).
Alongside broker and exporter metrics
Spitfire's lag measurement covers the run and the groups you watch. If you already write lag to Prometheus (for example with a Kafka exporter), add that query to a Prometheus connection under Observability: the run page's Backend tab shows it and your broker metrics aligned with the load stages, and can draw it over the run's own charts. When the run breaks, series with queue, lag or backlog in their name that at least doubled to 10 or more rank high among the resources (OpenTelemetry root cause). To find the rate at which consumers fall behind, use a breakpoint test on the producer scenario with the threshold on your service's group lag.
Spitfire installs on Docker or Kubernetes with one command; every testing feature and protocol is open in the free edition.