verae-nats-cluster/scripts/bench-report.py

192 lines
6.8 KiB
Python
Executable file
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""Turn nats bench text logs + latency JSON into markdown."""
from __future__ import annotations
import json
import re
import sys
from pathlib import Path
def fmt_int(s: str | None) -> str:
if not s:
return ""
return f"{int(s):,}"
def parse_bench(text: str) -> dict[str, str]:
out: dict[str, str] = {"kind": "throughput"}
def _rate(pattern: str, msgs_key: str, mb_key: str) -> None:
m = re.search(pattern, text)
if not m:
return
out[msgs_key] = m.group(1).replace(",", "")
n = float(m.group(2))
unit = m.group(3).upper()
if unit == "KB":
n = n / 1024.0
elif unit == "GB":
n = n * 1024.0
out[mb_key] = f"{n:.2f}"
_rate(
r"(?m)^\s*Pub stats:\s*([0-9,]+)\s*msgs/sec\s*~\s*([0-9.]+)\s*(KB|MB|GB)/sec",
"pub_msgs",
"pub_mb",
)
_rate(
r"(?m)^\s*Sub stats:\s*([0-9,]+)\s*msgs/sec\s*~\s*([0-9.]+)\s*(KB|MB|GB)/sec",
"sub_msgs",
"sub_mb",
)
_rate(
r"NATS Pub/Sub stats:\s*([0-9,]+)\s*msgs/sec\s*~\s*([0-9.]+)\s*(KB|MB|GB)/sec",
"agg_msgs",
"agg_mb",
)
if "JetStream" in text or "--js" in text or "js-" in text:
out["mode"] = "jetstream r=3 file"
else:
out["mode"] = "core pub/sub"
# nats 0.1.6 prints min/avg/max as msgs/sec across publishers, not µs delay
m = re.search(
r"min\s+([0-9,]+)\s*\|\s*avg\s+([0-9,]+)\s*\|\s*max\s+([0-9,]+)\s*\|\s*stddev\s+([0-9,]+)\s*msgs",
text,
)
if m:
out["pub_spread"] = f"{m.group(1)}{m.group(3)} (avg {m.group(2)})"
return out
def parse_lat(text: str) -> dict[str, str] | None:
for line in text.splitlines():
line = line.strip()
if line.startswith("{") and "p99_us" in line:
d = json.loads(line)
mode = d.get("mode") or ""
return {
"kind": "latency",
"count": str(d.get("count", "")),
"pubs": str(d.get("pubs", "")),
"size": str(d.get("size", "")),
"mode": str(mode),
"min": d.get("min", ""),
"avg": d.get("avg", ""),
"p50": d.get("p50", ""),
"p90": d.get("p90", ""),
"p99": d.get("p99", ""),
"max": d.get("max", ""),
}
return None
def lat_mode(run: str, recorded: str) -> str:
if recorded in ("ping", "flood"):
return recorded
if "ping" in run:
return "ping"
return "flood"
def thru_sort(p: dict[str, str]) -> tuple:
return (0 if p.get("mode", "").startswith("core") else 1, p.get("run", ""))
def lat_sort(p: dict[str, str]) -> tuple:
mode = lat_mode(p.get("run", ""), p.get("mode", ""))
return (0 if mode == "ping" else 1, int(p.get("count") or 0), p.get("run", ""))
def main() -> int:
folder = Path(sys.argv[1] if len(sys.argv) > 1 else ".")
thru: list[dict[str, str]] = []
lats: list[dict[str, str]] = []
for f in sorted(folder.glob("*.txt")):
text = f.read_text(encoding="utf-8", errors="replace")
lat = parse_lat(text)
if lat:
lat["run"] = f.stem
lats.append(lat)
continue
p = parse_bench(text)
if p.get("pub_msgs") or p.get("agg_msgs"):
p["run"] = f.stem
thru.append(p)
stamp = folder.name if re.fullmatch(r"\d{8}T\d{6}Z", folder.name) else ""
print("# NATS cluster message speed")
print()
if stamp:
print(f"Run **`{stamp}`** (UTC). ", end="")
print(
"Client: LXC **510** `verae-px-worker` (`10.10.10.20`), not a nats-* server. "
"Servers: `nats-a/b/c` on `10.10.10.2123` (`vmbr1` only)."
)
print()
print("Client URL:")
print()
print("```text")
print("nats://10.10.10.21:4222,nats://10.10.10.22:4222,nats://10.10.10.23:4222")
print("```")
print()
print("## Method")
print()
print("- **Core NATS** is fire-and-forget pub/sub (`nats bench`). No disk, no replica ack.")
print(
"- **JetStream** uses **file** storage and **replicas=3** (same as product streams). "
"The unique stream `benchstream` is deleted between JS loads."
)
print("- Throughput is **msgs/sec** from nats CLI **0.1.6** (`--no-progress --csv`). Its min/avg/max are publisher **rate spread**, not delay.")
print(
"- **Ping** delay: one publisher, sequential publish-then-wait. This is one-message round-trip through the cluster."
)
print(
"- **Flood** delay: N publishers dump the whole batch, then the subscriber drains. "
"This is **queueing under burst**, not wire RTT."
)
print("- Probe: `scripts/latency.mjs` (two connections, header timestamp).")
print()
print("## Throughput")
print()
print("| Run | Mode | Aggregate msgs/s | Pub msgs/s | Pub MB/s | Sub msgs/s | Sub MB/s |")
print("|-----|------|------------------|------------|----------|------------|----------|")
for p in sorted(thru, key=thru_sort):
print(
f"| `{p['run']}` | {p.get('mode', '')} | {fmt_int(p.get('agg_msgs'))} | "
f"{fmt_int(p.get('pub_msgs'))} | {p.get('pub_mb') or ''} | "
f"{fmt_int(p.get('sub_msgs'))} | {p.get('sub_mb') or ''} |"
)
print()
print("## Round-trip delay")
print()
print("| Run | Kind | Count | Pubs | Size | min | avg | p50 | p90 | p99 | max |")
print("|-----|------|-------|------|------|-----|-----|-----|-----|-----|-----|")
for p in sorted(lats, key=lat_sort):
kind = lat_mode(p.get("run", ""), p.get("mode", ""))
label = "ping (sequential RTT)" if kind == "ping" else "flood (burst queueing)"
print(
f"| `{p['run']}` | {label} | {p.get('count', '')} | {p.get('pubs', '')} | "
f"{p.get('size', '')} B | {p.get('min', '')} | {p.get('avg', '')} | "
f"{p.get('p50', '')} | {p.get('p90', '')} | {p.get('p99', '')} | {p.get('max', '')} |"
)
print()
print("## What the numbers mean")
print()
print(
"Product job/event/archive traffic is **JetStream r=3 file**. On this three-LXC stand that is about "
"**16k durable 128 B pubs/s** (about **13k** at 1 KiB). Core NATS is an upper bound for "
"non-durable fan-out: about **0.72.0M msgs/s** aggregate at 128 B, or **~630k msgs/s (~616 MB/s)** at 1 KiB with 4 publishers."
)
print()
print(
"A quiet request-reply is **~0.3 ms** average, **p99 < 1 ms**. Flood rows in the **150500 ms** band "
"are the subscriber catching up after a burst, which is what a job-events mailbox sees if publishers outrun consumers."
)
print()
print("Re-run on NS1: `bash scripts/bench.sh`. Raw logs/CSVs are under `results/<utc>/`.")
return 0
if __name__ == "__main__":
raise SystemExit(main())