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
Expand Up @@ -9,6 +9,20 @@ 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.

## Sequence

This is the second of four runs published on 2026-09-23:

1. [Harness capacity](../2026-09-23-north-star-steady-v1-harness-capacity-33c95bf/analysis.md):
`kq-bench` saturated the host near 13,000 tasks/s, mostly on its own
producer and Python accounting overhead.
2. **This run.** Remove the harness overhead and find kq's own ceiling.
3. [Queue-latency diagnosis](../2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/analysis.md):
explain the flat roughly 465 ms queue p50 observed here.
4. [Share group with 1000 members](../2026-09-23-adhoc-share-group-1000-members-33c95bf/analysis.md):
revisit this run's 20,000 RPS concurrency series with the member cap
raised from 200 to 1000.

## Purpose

Measure kq's own capacity ceiling on the north-star execution-time
Expand Down
Original file line number Diff line number Diff line change
@@ -1,21 +1,21 @@
#!/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>
# env: EXEC=northstar|const:<ms> OUT=<dir> COMPOSE=<compose file> SETTLE=<seconds>
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=(docker compose -f ${COMPOSE:-$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
sleep ${SETTLE:-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"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,20 @@ Results: [result.json](result.json)
This document is generated by AI from the run metadata, machine results, and
retained diagnostics.

## Sequence

This is the third of four runs published on 2026-09-23:

1. [Harness capacity](../2026-09-23-north-star-steady-v1-harness-capacity-33c95bf/analysis.md):
`kq-bench` saturated the host near 13,000 tasks/s.
2. [Lean-probe capacity](../2026-09-23-adhoc-lean-probe-capacity-33c95bf/analysis.md):
kq reached 150,000 tasks/s, but queue p50 stayed near 465 ms regardless of
load at concurrency 1000.
3. **This run.** Find what causes that queue time.
4. [Share group with 1000 members](../2026-09-23-adhoc-share-group-1000-members-33c95bf/analysis.md):
follows suggestion 2 below by raising the member cap to 1000 and running
1000 members x concurrency 10 at 20,000 RPS.

## Purpose

Explain the roughly 465 ms queue p50 observed at worker concurrency 1000 by
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,122 @@
# Run Analysis

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

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

Kafka: [compose.yaml](compose.yaml) (`bench/compose.yaml` plus
`group.share.max.size=1000` and `group.share.partition.max.record.locks=4000`)

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

## Sequence

This is the fourth and latest of four runs published on 2026-09-23:

1. [Harness capacity](../2026-09-23-north-star-steady-v1-harness-capacity-33c95bf/analysis.md):
`kq-bench` saturated the host near 13,000 tasks/s.
2. [Lean-probe capacity](../2026-09-23-adhoc-lean-probe-capacity-33c95bf/analysis.md):
kq reached 150,000 tasks/s, but queue p50 stayed near 465 ms.
3. [Queue-latency diagnosis](../2026-09-23-adhoc-queue-latency-diagnosis-33c95bf/analysis.md):
queue time comes from the worker generation wait; many small generations
are fast, but the default 200-member cap limits their throughput.
4. **This run.** Raise the member cap to 1000 and test whether 1000 small
generations sustain 20,000 RPS with low queue time.

## Purpose

Test whether raising the share-group member limit to its maximum (1000) lets
many small worker generations (1000 members x concurrency 10, emulating a
sharded consumer) sustain 20,000 RPS with low queue time, and measure
sensitivity to ready partition count.

## Observations

20,000 RPS for 60 s (1,200,000 tasks). 10 processes x 100 workers x
concurrency 10 = 1000 share-group members and 10,000 slots:

| Ready partitions | Outcome | Median / peak completions/s | Queue p50 / p95 / p99 / max | Kafka / worker cores |
| ---: | --- | ---: | ---: | ---: |
| 2 | collapsed; 60,959 completed | 256 / 7,274 | 3,627 / 56,544 / 68,915 / 72,032 ms | 7.21 / 0.25 |
| 4 | collapsed; 230,552 completed | 1,136 / 16,123 | 2,354 / 40,253 / 55,050 / 60,916 ms | 7.56 / 0.29 |
| 8 | sustained | 20,027 / 21,002 | 14 / 104 / 203 / 621 ms | 5.80 / 1.62 |
| 16 | sustained | 20,012 / 20,933 | 24 / 89 / 179 / 1,170 ms | 5.88 / 1.63 |
| 32 | sustained | 20,032 / 20,896 | 43 / 158 / 272 / 665 ms | 5.82 / 1.62 |

- `group.share.max.size=1000` is the broker maximum; 1001 is rejected at
startup. `group.share.partition.max.record.locks` was raised to 4000, the
default value of its upper bound
(`group.share.max.partition.max.record.locks`).
- The member cap applies to the whole share group regardless of partition
count.
- With 2 and 4 partitions, completions peaked in the first seconds of
production (7,274/s and 16,123/s), then decayed continuously to about
200/s and 520/s while producers kept enqueuing 20,000/s. Completions
stopped entirely after about 72 to 76 s. Kafka used 7.2 to 7.6 cores.
- With 8, 16, and 32 partitions, all 3,600,000 tasks completed with zero
duplicates and zero missing IDs, and completions matched arrival.
- For comparison at 20,000 RPS: the lean-probe run's best result under the
default 200-member cap (180 members x concurrency 100) was
115 / 410 / 1,291 ms queue p50 / p95 / p99.
- 1000 members cost about 1 more Kafka core and 1 more worker core than the
32-member x concurrency 1000 topology in the lean-probe run (4.81 / 0.67).

## Interpretation

- Many small generations work once the member cap is raised. At 8 or 16
partitions, 1000 members x concurrency 10 sustained the 20,000 RPS north-star
rate with queue p50 of 14 to 24 ms and p99 of 179 to 203 ms, roughly an
order of magnitude better than the 200-member result.
- Partition count bounds throughput independently of member count.
Hypothesis: a share partition's in-flight window runs from its oldest
unacknowledged record and holds at most `record.locks` records, so each
partition sustains roughly `4000 / time-from-acquire-to-ack of the oldest
record`. That time is about 1 s on this workload, which gives about 8,000/s
for 2 partitions and 16,000/s for 4. Those match the observed peaks of
7,274 and 16,123 but are not confirmed from broker metrics.
- Exceeding that bound caused collapse rather than a plateau. Hypothesis:
most member fetches return few or no records and are retried at once. That
raises broker load, which slows acknowledgement and window advance, which
lowers the bound further. Rising Kafka CPU during falling throughput is
consistent with this. Why completions stopped entirely is unexplained.
- 32 partitions had higher latency than 8 or 16. Hypothesis: each member's
10-record fetches are spread across more partitions and fetch sessions.
Untested.
- Each member still waits for its generation of 10, so the remaining p99
(about 180 to 200 ms) likely includes the generation wait. The
acquisition-while-busy hypothesis from run 3 is also unresolved.
- At 1000 members the cap now binds at about 20,000 to 25,000 tasks/s per
queue for concurrency 10 on this workload. Higher rates need larger
generations (worse tail latency), multiple queues, or a worker that refills
continuously.

## Suggestions

1. Design the sharded consumer around 1000 members, about 10 slots per shard,
and a partition count sized so that
`partitions x 4000 / expected acquire-to-ack time` comfortably exceeds peak
arrival.
2. Document partition sizing as a correctness-of-operation requirement; the
failure mode is collapse, not gradual degradation.
3. Confirm the in-flight window model before relying on it: run 2 partitions
at 5,000 RPS (below the estimated bound) and at 20,000 RPS with 200 members,
and capture broker share-partition metrics during a collapse.
4. Sweep shard concurrency (5, 10, 20) at 8 and 16 partitions to locate the
latency optimum.
5. Add these share-group settings and a partition-sizing check to the
eventual GCP benchmark, since multi-broker behavior may differ.

## Caveats

- Single Dockerized Kafka broker on localhost, RF=1; one run per point with a
60 s arrival window.
- Same probe as runs 2 and 3: ready success path only, sleep-only handlers,
200-byte payloads, queue time measured from just before `Enqueue`.
- For collapsed points, `correctness.missing` counts tasks not completed when
the probe exited, including tasks still in Kafka. Worker processes exit after
15 s without completions, so these points do not show whether the system
would have recovered.
- 1000 members ran as 1000 `kq.Worker` instances across 10 processes; a real
sharded consumer inside one worker may have different client and connection
overhead.
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
services:
kafka:
image: apache/kafka:4.3.1
ports:
- "19092:19092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:19092,CONTROLLER://:19093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:19092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:19093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_SHARE_COORDINATOR_STATE_TOPIC_REPLICATION_FACTOR: 1
KAFKA_SHARE_COORDINATOR_STATE_TOPIC_MIN_ISR: 1
KAFKA_GROUP_SHARE_MIN_RECORD_LOCK_DURATION_MS: 1000
KAFKA_GROUP_SHARE_MAX_SIZE: 1000
KAFKA_GROUP_SHARE_PARTITION_MAX_RECORD_LOCKS: 4000
healthcheck:
test: ["CMD-SHELL", "/opt/kafka/bin/kafka-broker-api-versions.sh --bootstrap-server localhost:19092 >/dev/null 2>&1"]
interval: 2s
timeout: 5s
retries: 30
Loading
Loading