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,105 @@
# Run Analysis

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

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

Probe: [probe/main.go](probe/main.go), [probe/run.sh](probe/run.sh)

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

## Purpose

Measure kq's own capacity ceiling on the north-star execution-time
distribution with a minimal load generator that removes `kq-bench`'s per-event
JSON and Python accounting overhead, and measure the Kafka cost of concurrent
versus serial `Enqueue`.

## Observations

Capacity sweep. Every point used 512 concurrent `Enqueue` goroutines per
producer client:

| Requested RPS | Partitions | Members x concurrency | Median completions/s | Tail after production | Queue p50 / p95 / p99 / max | Kafka / producer / worker cores |
| ---: | ---: | ---: | ---: | ---: | ---: | ---: |
| 5,000 | 16 | 8 x 1000 | 4,985 | 1 s | 472 / 980 / 1,185 / 1,676 ms | 1.75 / 0.51 / 0.11 |
| 20,000 | 32 | 32 x 1000 | 20,094 | 2 s | 465 / 965 / 1,187 / 1,974 ms | 4.81 / 1.66 / 0.67 |
| 30,000 | 48 | 48 x 1000 | 30,093 | 2 s | 465 / 974 / 1,203 / 2,377 ms | 4.70 / 1.83 / 0.98 |
| 40,000 | 64 | 64 x 1000 | 39,838 | 2 s | 466 / 976 / 1,209 / 2,376 ms | 5.30 / 2.12 / 1.21 |
| 60,000 | 96 | 96 x 1000 | 60,117 | 2 s | 466 / 980 / 1,221 / 2,380 ms | 5.47 / 2.44 / 1.80 |
| 80,000 | 128 | 128 x 1000 | 80,073 | 2 s | 467 / 978 / 1,222 / 2,375 ms | 5.06 / 2.40 / 2.30 |
| 100,000 | 128 | 96 x 2000 | 99,816 | 3 s | 532 / 1,187 / 1,772 / 3,691 ms | 5.08 / 2.25 / 2.39 |
| 120,000 | 160 | 120 x 2000 | 120,278 | 2 s | 527 / 1,151 / 1,619 / 4,194 ms | 4.77 / 2.52 / 3.22 |
| 150,000 | 160 | 150 x 2000 | 149,580 | 6 s | 560 / 2,597 / 5,385 / 7,699 ms | 4.52 / 1.84 / 3.47 |

Concurrency at 20,000 RPS and 32 partitions:

| Members x concurrency | Slots | Queue p50 / p95 / p99 / max |
| ---: | ---: | ---: |
| 180 x 100 | 18,000 | 115 / 410 / 1,291 / 2,592 ms |
| 80 x 250 | 20,000 | 260 / 923 / 2,019 / 3,411 ms |
| 48 x 500 | 24,000 | 319 / 847 / 1,315 / 2,350 ms |
| 32 x 1000 | 32,000 | 465 / 965 / 1,187 / 1,974 ms |

Enqueue only, with no workers running:

| Clients x goroutines | Achieved enqueues/s | Kafka cores | Producer cores |
| ---: | ---: | ---: | ---: |
| 15 x 1 | 12,000 | 3.70 | 1.20 |
| 1 x 64 | 11,998 | 3.63 | 1.08 |
| 1 x 512 | 39,984 | 3.67 | 1.47 |
| 4 x 512 | 79,934 | 3.35 | 2.21 |

- All 39,750,000 tasks across the capacity and concurrency points completed
with zero duplicates and zero missing IDs. No enqueue errors occurred.
- Queue p50 stayed between 465 and 472 ms from 5,000 to 80,000 RPS at
concurrency 1000.
- At 150,000 RPS the median completion rate matched arrival, but queue p95 and
p99 rose to 2.6 s and 5.4 s and completions continued for 6 s after
production ended.
- At 100,000 RPS and above, completed tasks per worker process ranged from
about 630,000 to 1,010,000 within a single point.
- Kafka's idle baseline was about 0.5 cores.

## Interpretation

- kq sustains the 20,000 tasks/s north-star arrival rate on this host with
more than 4 cores idle. The knee near 150,000 RPS coincides with host CPU
saturation, so it is a host limit rather than a demonstrated kq limit.
- Worker-side kq CPU is small: 0.67 cores at 20,000 tasks/s.
- Kafka CPU was roughly flat from 12,000 to 80,000 enqueues/s. Hypothesis:
broker cost here is dominated by request handling and background threads
rather than per-record work, and franz-go batches concurrent `Enqueue` calls.
- Throughput is not limited by the fixed-generation worker, but required slots
(18,000 to 32,000) far exceed the roughly 5,000 tasks in flight implied by
Little's law. Latency, not throughput, is the cost of generations.
- `group.share.max.size` (default 200) caps members per share group. At
concurrency 100, 20,000 RPS already needed 180 members, so lower concurrency
cannot scale further without raising that limit.
- The uneven per-process completion counts at high rates suggest share-group
assignment or fetch skew. This needs per-partition data to confirm.

## Suggestions

1. Rework `kq-bench` so the harness itself can validate rates above 20,000 RPS:
concurrent producer enqueues, compact events, and streaming accounting.
2. Pursue a worker design that does not wait for a full generation before
fetching again. See the queue-latency diagnosis run for evidence.
3. Record per-partition acquisition and per-member throughput to explain
assignment skew.
4. Treat the member limit as a design constraint for any sharded-consumer
approach.

## Caveats

- This is not a `kq-bench` run. The probe measures only the ready success
path: no retries, movers, observer, dead-lettering, or warm-up accounting.
- Single Dockerized Kafka broker on localhost, RF=1, no network latency. `acks=all`
therefore involves one replica.
- Handlers only sleep; payloads are 200 bytes.
- Each point ran once with a 30 or 60 s arrival window.
- CPU figures are cumulative process CPU over a 10 s mid-run window.
- Queue time uses a producer timestamp taken just before `Enqueue`, so it
includes enqueue latency, unlike the harness's definition, which starts at
enqueue completion.
Original file line number Diff line number Diff line change
@@ -0,0 +1,242 @@
//go:build ignore

// Minimal kq load probe used by this run. It drives kq.Client and kq.Worker
// unmodified and records only task IDs and queue times.
//
// go build -o probe ./docs/performance/runs/2026-09-23-adhoc-lean-probe-capacity-33c95bf/probe/main.go
// probe produce -clients 2 -goroutines 512 -rps 20000 -dur 60s
// probe work -workers 8 -concurrency 1000 [-exec northstar|const:<ms>] -out ids-0.bin
// probe check ids-*.bin
package main

import (
"context"
"encoding/binary"
"flag"
"fmt"
"math"
"math/rand/v2"
"os"
"os/signal"
"slices"
"strconv"
"strings"
"sync"
"sync/atomic"
"syscall"
"time"

"github.com/raydatray/kq"
)

const queue = "probe"

func cfg() kq.Config {
policy, _ := kq.NewRetryPolicy(0, func(int32, string) time.Duration { return 0 })
return kq.Config{Brokers: []string{"localhost:19092"}, Queue: queue, RetryPolicy: policy}
}

func main() {
switch os.Args[1] {
case "produce":
produce(os.Args[2:])
case "work":
work(os.Args[2:])
case "check":
check(os.Args[2:])
}
}

func produce(args []string) {
fs := flag.NewFlagSet("produce", flag.ExitOnError)
clients := fs.Int("clients", 1, "kq.Client instances")
perClient := fs.Int("goroutines", 512, "concurrent Enqueue callers per client")
rps := fs.Int("rps", 12000, "target total enqueues/s")
dur := fs.Duration("dur", 30*time.Second, "duration")
fs.Parse(args)

total := uint64(*rps) * uint64(dur.Seconds())
callers := *clients * *perClient
var next, done, errs atomic.Uint64
var wg sync.WaitGroup
start := time.Now().Add(200 * time.Millisecond)
for c := 0; c < *clients; c++ {
client, err := kq.NewClient(cfg())
if err != nil {
panic(err)
}
defer client.Close()
for g := 0; g < *perClient; g++ {
wg.Add(1)
go func() {
defer wg.Done()
payload := make([]byte, 200)
for {
id := next.Add(1) - 1
if id >= total {
return
}
due := start.Add(time.Duration(float64(id) / float64(*rps) * float64(time.Second)))
if d := time.Until(due); d > 0 {
time.Sleep(d)
}
binary.LittleEndian.PutUint64(payload[0:], id)
binary.LittleEndian.PutUint64(payload[8:], uint64(time.Now().UnixNano()))
if _, err := client.Enqueue(context.Background(), kq.Task{Type: "probe", Payload: payload}); err != nil {
errs.Add(1)
continue
}
done.Add(1)
}
}()
}
}
_ = callers
wg.Wait()
el := time.Since(start).Seconds()
fmt.Printf("produce: requested=%d/s achieved=%.0f/s total=%d errs=%d\n", *rps, float64(done.Load())/el, done.Load(), errs.Load())
}

// north-star: lognormal, mean 250 ms, p95 ~475 ms, p99 ~646 ms.
var (
sigma = math.Log(646.0/475.0) / (2.326 - 1.645)
median = 250.0 / math.Exp(sigma*sigma/2)
)

func execTime(id uint64) time.Duration {
r := rand.New(rand.NewPCG(id, 0x9e3779b97f4a7c15))
ms := median * math.Exp(sigma*r.NormFloat64())
return time.Duration(ms * float64(time.Millisecond))
}

func pick(mode string, id uint64) time.Duration {
if strings.HasPrefix(mode, "const:") {
ms, _ := strconv.Atoi(strings.TrimPrefix(mode, "const:"))
return time.Duration(ms) * time.Millisecond
}
return execTime(id)
}

func work(args []string) {
fs := flag.NewFlagSet("work", flag.ExitOnError)
workers := fs.Int("workers", 4, "kq.Worker instances in this process")
concurrency := fs.Int("concurrency", 1000, "Worker concurrency")
out := fs.String("out", "ids.bin", "completed ids output")
idle := fs.Duration("idle", 15*time.Second, "exit after this long with no completions (after first)")
execMode := fs.String("exec", "northstar", "northstar | const:<ms>")
fs.Parse(args)

var mu sync.Mutex
ids := make([]uint64, 0, 1<<20)
queueMS := make([]float32, 0, 1<<20)
var completed atomic.Int64
var lastDone atomic.Int64

handler := func(ctx context.Context, task kq.Task) error {
id := binary.LittleEndian.Uint64(task.Payload[0:])
enq := int64(binary.LittleEndian.Uint64(task.Payload[8:]))
q := float32(time.Now().UnixNano()-enq) / 1e6
select {
case <-time.After(pick(*execMode, id)):
case <-ctx.Done():
return ctx.Err()
}
mu.Lock()
ids = append(ids, id)
queueMS = append(queueMS, q)
mu.Unlock()
completed.Add(1)
lastDone.Store(time.Now().UnixNano())
return nil
}

ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer cancel()
var wg sync.WaitGroup
for range *workers {
w, err := kq.NewWorker(kq.WorkerConfig{Config: cfg(), Concurrency: *concurrency}, handler)
if err != nil {
panic(err)
}
wg.Add(1)
go func() {
defer wg.Done()
if err := w.Run(ctx); err != nil {
fmt.Fprintln(os.Stderr, "worker error:", err)
}
}()
}
fmt.Fprintln(os.Stderr, "ready")

var prev int64
var peak float64
ticker := time.NewTicker(time.Second)
for range ticker.C {
cur := completed.Load()
rate := float64(cur - prev)
prev = cur
if rate > peak {
peak = rate
}
fmt.Fprintf(os.Stderr, "%s completed=%d rate=%.0f/s\n", time.Now().Format("15:04:05"), cur, rate)
if l := lastDone.Load(); l > 0 && time.Since(time.Unix(0, l)) > *idle {
break
}
if ctx.Err() != nil {
break
}
}
cancel()
wg.Wait()

mu.Lock()
defer mu.Unlock()
f, _ := os.Create(*out)
buf := make([]byte, 8)
for _, id := range ids {
binary.LittleEndian.PutUint64(buf, id)
f.Write(buf)
}
f.Close()
qf, _ := os.Create(*out + ".queue")
for _, q := range queueMS {
binary.LittleEndian.PutUint32(buf, math.Float32bits(q))
qf.Write(buf[:4])
}
qf.Close()
fmt.Printf("work: completed=%d peak_1s_rate=%.0f/s\n", len(ids), peak)
}

func check(files []string) {
seen := map[uint64]int{}
var queueMS []float64
var maxID uint64
for _, name := range files {
data, _ := os.ReadFile(name)
for i := 0; i+8 <= len(data); i += 8 {
id := binary.LittleEndian.Uint64(data[i:])
seen[id]++
maxID = max(maxID, id)
}
qd, _ := os.ReadFile(name + ".queue")
for i := 0; i+4 <= len(qd); i += 4 {
queueMS = append(queueMS, float64(math.Float32frombits(binary.LittleEndian.Uint32(qd[i:]))))
}
}
dups := 0
for _, n := range seen {
if n > 1 {
dups += n - 1
}
}
missing := int(maxID+1) - len(seen)
slices.Sort(queueMS)
p := func(q float64) float64 {
if len(queueMS) == 0 {
return 0
}
return queueMS[int(q*float64(len(queueMS)-1))]
}
fmt.Printf("check: unique=%d dups=%d missing(<=max id)=%d queue p50/p95/p99/max=%.0f/%.0f/%.0f/%.0f ms\n",
len(seen), dups, missing, p(0.5), p(0.95), p(0.99), p(1))
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
#!/bin/zsh
# usage: run.sh LABEL RPS DUR PARTITIONS WORKER_PROCS WORKERS_PER_PROC CONCURRENCY [PRODUCER_CLIENTS] [GOROUTINES]
# env: EXEC=northstar|const:<ms> OUT=<dir>
set -u
HERE=${0:A:h}
REPO=${HERE:h:h:h:h:h}
LABEL=$1 RPS=$2 DUR=$3 PART=$4 WP=$5 WPP=$6 C=$7 PC=${8:-2} PG=${9:-512} EXEC=${EXEC:-northstar}
D=${OUT:-${TMPDIR:-/tmp}/kq-probe}/$LABEL; rm -rf $D; mkdir -p $D
BIN=$D/probe
(cd $REPO && go build -o $BIN $HERE/main.go)
KC=(docker compose -f $REPO/bench/compose.yaml -p kqprobe)
$KC up -d --wait >/dev/null 2>&1
$KC exec -T kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:19092 --create --topic probe-ready --partitions $PART --replication-factor 1 >/dev/null
$KC exec -T kafka /opt/kafka/bin/kafka-configs.sh --bootstrap-server localhost:19092 --alter --entity-type groups --entity-name kq.probe.workers --add-config share.auto.offset.reset=earliest >/dev/null
for i in $(seq 0 $((WP-1))); do
$BIN work -workers $WPP -concurrency $C -exec $EXEC -out $D/ids-$i.bin > $D/work-$i.out 2> $D/work-$i.err &
done
sleep 12 # let share group assignment settle
$BIN produce -clients $PC -goroutines $PG -rps $RPS -dur ${DUR}s > $D/produce.out 2>&1
wait
echo "== $LABEL rps=$RPS dur=$DUR part=$PART procs=$WP x workers=$WPP x C=$C exec=$EXEC"
cat $D/produce.out $D/work-*.out
$BIN check $D/ids-*.bin
$KC down -v >/dev/null 2>&1
Loading
Loading