From a3d7b7f2887075e5d0afd1f4eef4512fb9520b0d Mon Sep 17 00:00:00 2001 From: Ray Liu Date: Wed, 23 Sep 2026 00:41:20 -0400 Subject: [PATCH] [3/6][load] publish queue latency diagnosis run --- .../analysis.md | 78 +++++++ .../result.json | 214 ++++++++++++++++++ .../run.json | 50 ++++ 3 files changed, 342 insertions(+) create mode 100644 docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/analysis.md create mode 100644 docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/result.json create mode 100644 docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/run.json diff --git a/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/analysis.md b/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/analysis.md new file mode 100644 index 0000000..c016349 --- /dev/null +++ b/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/analysis.md @@ -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`. diff --git a/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/result.json b/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/result.json new file mode 100644 index 0000000..9cca224 --- /dev/null +++ b/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/result.json @@ -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 + } + } + ] +} diff --git a/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/run.json b/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/run.json new file mode 100644 index 0000000..fa2553b --- /dev/null +++ b/docs/performance/runs/2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/run.json @@ -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:", + "command": "EXEC= probe/run.sh