Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
# Run Analysis

Run: [run.json](run.json)

Results: [result.json](result.json)

This document is generated by AI from the run metadata, machine results, and
retained diagnostics.

## Purpose

Explain the roughly 465 ms queue p50 observed at worker concurrency 1000 by
varying execution time and worker concurrency at low load with ample idle
slots.

## Observations

| Members x concurrency | RPS | Execution | Queue p50 / p95 / p99 / max |
| ---: | ---: | --- | ---: |
| 8 x 1000 | 2,000 | 0 ms | 0 / 1 / 6 / 310 ms |
| 8 x 1000 | 2,000 | constant 250 ms | 131 / 245 / 254 / 367 ms |
| 8 x 1000 | 2,000 | north-star | 402 / 873 / 1,105 / 1,673 ms |
| 180 x 10 | 2,000 | north-star | 20 / 81 / 119 / 300 ms |
| 180 x 1 | 500 | north-star | 1 / 70 / 191 / 344 ms |
| 180 x 1 | 500 | 0 ms | 1 / 1 / 9 / 292 ms |

- All 270,000 tasks completed with zero duplicates and zero missing IDs.
- At 2,000 RPS the 8 x 1000 topology had 8,000 slots for roughly 500 tasks in
flight, yet north-star queue p50 was 402 ms.
- With constant 250 ms execution, queue times were spread from 0 to about
250 ms, with a maximum of 367 ms.

## Interpretation

- Kafka delivery and kq's worker loop add little latency on their own: with
zero execution time, queue p99 was 6 to 9 ms.
- Queue time is dominated by waiting for a worker's in-progress generation to
finish. `Worker.Run` polls up to `Concurrency` records, waits for every one,
flushes acknowledgements, and only then polls again. Constant 250 ms
execution produced waits bounded by about one generation.
- A generation lasts as long as its slowest record. Larger generations sample
further into the lognormal tail, which is why north-star waits grow with
concurrency. The 20,000 RPS concurrency series in the lean-probe run shows
the same trend.
- Idle slots in other members did not absorb new work. Hypothesis: the
franz-go share consumer acquires records into a member whose generation is
still running, so those records wait behind it. The 180 x 1 north-star
point (p95 70 ms despite 180 members for about 125 tasks in flight) is
consistent with this, but franz-go fetch behavior was not inspected.
- franz-go finalizes unresolved share records when the next poll begins
(`docs/plans/0003-worker-concurrency.md`), so a single consumer cannot
refill freed slots while earlier records are unresolved.

## Suggestions

1. Design the next worker change around removing the generation wait, for
example a sharded consumer (several share consumers per worker, each
running small independent generations) or explicit acknowledgement across
polls if franz-go gains support for it.
2. Account for the share-group member limit in a sharded design. Each shard is
a share-group member, and `group.share.max.size` defaults to 200. The
180 x 10 point is effectively 180 shards of 10 and reached 20 / 81 / 119 ms,
but 200 members of 10 cannot reach 20,000 tasks/s on this workload. The
20,000 RPS point with 180 members of 100 had a 1,291 ms p99. Measure shard
size against both latency and the member limit, or raise the broker limit.
3. Confirm the acquisition-while-busy hypothesis with franz-go fetch and
acquisition tracing before choosing between designs.
4. Add a queue-latency comparison scenario to `kq-bench` covering zero,
constant, and north-star execution time so the effect is regression-tested.

## Caveats

- Same probe, host, and single-broker Kafka setup as the lean-probe capacity
run; one run per point with a 30 s arrival window.
- The 180 x 1 points ran at 500 RPS rather than 2,000 because 180 members at
concurrency 1 cannot sustain 2,000 tasks/s with north-star execution.
- Queue time includes enqueue latency because the timestamp is taken before
`Enqueue`.
Original file line number Diff line number Diff line change
@@ -0,0 +1,214 @@
{
"correctness": {
"produced": 270000,
"completed": 270000,
"duplicates": 0,
"missing": 0
},
"points": [
{
"label": "x-c1000-zero",
"requested_rps": 2000,
"duration_seconds": 30,
"ready_partitions": 16,
"worker_processes": 2,
"workers_per_process": 4,
"worker_concurrency": 1000,
"share_group_members": 8,
"total_slots": 8000,
"producer_clients": 1,
"enqueue_goroutines_per_client": 64,
"execution": "const:0",
"achieved_enqueue_rps": 2000,
"enqueue_errors": 0,
"completion_rps_median_1s": 2000,
"completion_rps_max_1s": 2007,
"seconds_with_completions_after_production": 0,
"queue_p50_ms": 0,
"queue_p95_ms": 1,
"queue_p99_ms": 6,
"queue_max_ms": 310,
"correctness": {
"produced": 60000,
"completed": 60000,
"duplicates": 0,
"missing": 0
},
"cpu_cores_10s_window": {
"kafka_docker_vm": 2.59,
"producer": 0.26,
"workers": 0.4
}
},
{
"label": "x-c1000-const250",
"requested_rps": 2000,
"duration_seconds": 30,
"ready_partitions": 16,
"worker_processes": 2,
"workers_per_process": 4,
"worker_concurrency": 1000,
"share_group_members": 8,
"total_slots": 8000,
"producer_clients": 1,
"enqueue_goroutines_per_client": 64,
"execution": "const:250",
"achieved_enqueue_rps": 2000,
"enqueue_errors": 0,
"completion_rps_median_1s": 2030,
"completion_rps_max_1s": 2747,
"seconds_with_completions_after_production": 0,
"queue_p50_ms": 131,
"queue_p95_ms": 245,
"queue_p99_ms": 254,
"queue_max_ms": 367,
"correctness": {
"produced": 60000,
"completed": 60000,
"duplicates": 0,
"missing": 0
},
"cpu_cores_10s_window": {
"kafka_docker_vm": 1.02,
"producer": 0.22,
"workers": 0.03
}
},
{
"label": "x-c1000-ns",
"requested_rps": 2000,
"duration_seconds": 30,
"ready_partitions": 16,
"worker_processes": 2,
"workers_per_process": 4,
"worker_concurrency": 1000,
"share_group_members": 8,
"total_slots": 8000,
"producer_clients": 1,
"enqueue_goroutines_per_client": 64,
"execution": "northstar",
"achieved_enqueue_rps": 2000,
"enqueue_errors": 0,
"completion_rps_median_1s": 1996,
"completion_rps_max_1s": 2241,
"seconds_with_completions_after_production": 1,
"queue_p50_ms": 402,
"queue_p95_ms": 873,
"queue_p99_ms": 1105,
"queue_max_ms": 1673,
"correctness": {
"produced": 60000,
"completed": 60000,
"duplicates": 0,
"missing": 0
},
"cpu_cores_10s_window": {
"kafka_docker_vm": 1.0,
"producer": 0.22,
"workers": 0.05
}
},
{
"label": "x-c10-ns",
"requested_rps": 2000,
"duration_seconds": 30,
"ready_partitions": 16,
"worker_processes": 6,
"workers_per_process": 30,
"worker_concurrency": 10,
"share_group_members": 180,
"total_slots": 1800,
"producer_clients": 1,
"enqueue_goroutines_per_client": 64,
"execution": "northstar",
"achieved_enqueue_rps": 2000,
"enqueue_errors": 0,
"completion_rps_median_1s": 2002,
"completion_rps_max_1s": 2099,
"seconds_with_completions_after_production": 1,
"queue_p50_ms": 20,
"queue_p95_ms": 81,
"queue_p99_ms": 119,
"queue_max_ms": 300,
"correctness": {
"produced": 60000,
"completed": 60000,
"duplicates": 0,
"missing": 0
},
"cpu_cores_10s_window": {
"kafka_docker_vm": 1.61,
"producer": 0.25,
"workers": 0.21
}
},
{
"label": "x-c1-ns",
"requested_rps": 500,
"duration_seconds": 30,
"ready_partitions": 16,
"worker_processes": 6,
"workers_per_process": 30,
"worker_concurrency": 1,
"share_group_members": 180,
"total_slots": 180,
"producer_clients": 1,
"enqueue_goroutines_per_client": 64,
"execution": "northstar",
"achieved_enqueue_rps": 500,
"enqueue_errors": 0,
"completion_rps_median_1s": 500,
"completion_rps_max_1s": 536,
"seconds_with_completions_after_production": 1,
"queue_p50_ms": 1,
"queue_p95_ms": 70,
"queue_p99_ms": 191,
"queue_max_ms": 344,
"correctness": {
"produced": 15000,
"completed": 15000,
"duplicates": 0,
"missing": 0
},
"cpu_cores_10s_window": {
"kafka_docker_vm": 1.43,
"producer": 0.16,
"workers": 0.26
}
},
{
"label": "x-c1-zero",
"requested_rps": 500,
"duration_seconds": 30,
"ready_partitions": 16,
"worker_processes": 6,
"workers_per_process": 30,
"worker_concurrency": 1,
"share_group_members": 180,
"total_slots": 180,
"producer_clients": 1,
"enqueue_goroutines_per_client": 64,
"execution": "const:0",
"achieved_enqueue_rps": 500,
"enqueue_errors": 0,
"completion_rps_median_1s": 500,
"completion_rps_max_1s": 501,
"seconds_with_completions_after_production": 0,
"queue_p50_ms": 1,
"queue_p95_ms": 1,
"queue_p99_ms": 9,
"queue_max_ms": 292,
"correctness": {
"produced": 15000,
"completed": 15000,
"duplicates": 0,
"missing": 0
},
"cpu_cores_10s_window": {
"kafka_docker_vm": 1.45,
"producer": 0.15,
"workers": 0.23
}
}
]
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
{
"run_id": "2026-09-23-adhoc-queue-latency-diagnosis-33c95bf",
"created_at": "2026-09-23T02:30:00Z",
"revision": "33c95bf74b5c",
"purpose": "Explain the ~465 ms queue p50 observed at worker concurrency 1000 by varying execution time and worker concurrency at low load with ample idle slots.",
"scenario": null,
"profile": null,
"runner": "../2026-09-23-adhoc-lean-probe-capacity-33c95bf/probe/main.go with -exec northstar | const:<ms>",
"command": "EXEC=<mode> probe/run.sh <label> <rps> 30 16 <worker processes> <workers per process> <concurrency> 1 64",
"environment": {
"architecture": "arm64",
"ci_provider": "",
"cpu_count": 12,
"go_version": "go version go1.27.0 darwin/arm64",
"kafka_version": "apache/kafka:4.3.1",
"memory_bytes": 22118400,
"operating_system": "Darwin",
"python_version": "3.13.12",
"runner": "rays-MacBook-Pro.local",
"docker_vm_cpus": 12,
"docker_vm_memory_bytes": 8216776704
},
"topology": {
"retry_partitions": 0,
"movers": 0,
"notes": "Per-point topology in result.json."
},
"workload": {
"arrival": {
"target_rps_sweep": [
500,
2000
],
"duration_seconds": 30
},
"execution_time_ms": {
"average": 250,
"p95": 475,
"p99": 646,
"distribution": "lognormal, sigma=ln(646/475)/(2.326-1.645), mean 250 ms, seeded by task ID"
},
"failures": {
"rate": 0,
"mode": "none",
"duration_ms": 0
},
"payload_bytes": 200
},
"notes": "Same probe, Kafka setup, and queue-time definition as the lean probe capacity run. One run per point."
}
Loading