117 lines
4.4 KiB
Python
117 lines
4.4 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""批量写入压测脚本 —— 验证 P99 ≤ 1.8s(issue #32 / PRD 5.2)。
|
||
|
||
验收口径(父 Issue #4「② 数据总线 + 时序库」):
|
||
- 批量写入:`5k/100ms` 批量基线(batch_size=5000、flush_interval=0.1s);
|
||
- 端到端延迟 P99 ≤ 1.8s(对齐 PRD 5.1 采集链路口径,批量写入侧同标)。
|
||
|
||
压测方式:
|
||
- 使用 `BatchWriter` + `MemorySink`(内存落库,无外部依赖,本地可跑);
|
||
- 模拟 600 点位 1Hz 采样行,按 batch_size 批量 push_many + flush;
|
||
- 记录每次 flush 耗时(写入端到端延迟),统计 P99 与吞吐;
|
||
- 也可注入真实 `TdengineSink(executor=...)` 落库压测(--sink tdengine)。
|
||
|
||
用法(在 core/data-bus 目录下):
|
||
python scripts/bench_write.py [--points 600] [--rounds 100] \
|
||
[--batch-size 5000] [--flush-interval 0.1]
|
||
退出码:0 = PASS(P99 ≤ 1.8s);1 = FAIL。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import importlib.util
|
||
import os
|
||
import statistics
|
||
import sys
|
||
import time
|
||
|
||
# `core/data-bus` 目录含连字符,无法直接以包名 import:挂载为 data_bus 包
|
||
_DATA_BUS_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||
if "data_bus" not in sys.modules:
|
||
spec = importlib.util.spec_from_file_location(
|
||
"data_bus", os.path.join(_DATA_BUS_DIR, "__init__.py"),
|
||
submodule_search_locations=[_DATA_BUS_DIR])
|
||
_pkg = importlib.util.module_from_spec(spec)
|
||
sys.modules["data_bus"] = _pkg
|
||
spec.loader.exec_module(_pkg)
|
||
sys.path.insert(0, _DATA_BUS_DIR)
|
||
|
||
from data_bus.batch_writer import BatchWriter, MemorySink # noqa: E402
|
||
|
||
P99_TARGET_SEC = 1.8 # PRD 5.2 验收:批量写入端到端 P99 ≤ 1.8s
|
||
|
||
|
||
def make_rows(points: int, seq: int) -> list:
|
||
"""生成一批模拟采样行(points 个点位 × 1 条)。"""
|
||
base = seq * points
|
||
return [
|
||
{"device_id": f"CLF-{i // 100 + 1:02d}",
|
||
"point_id": f"CLF-{i // 100 + 1:02d}.P{i % 100:03d}",
|
||
"value": round(100 + i * 0.1, 3), "ts": time.time() + i * 0.001,
|
||
"unit": "℃"}
|
||
for i in range(base, base + points)
|
||
]
|
||
|
||
|
||
def p99(values: list) -> float:
|
||
"""P99 分位(升序第 99% 位)。"""
|
||
ordered = sorted(values)
|
||
idx = max(0, min(len(ordered) - 1, int(len(ordered) * 0.99)))
|
||
return ordered[idx]
|
||
|
||
|
||
def run_bench(args) -> int:
|
||
print(f"== 批量写入压测({args.points} 点位 × {args.rounds} 轮,"
|
||
f"batch={args.batch_size},flush={args.flush_interval}s)==")
|
||
|
||
sink = MemorySink()
|
||
writer = BatchWriter(sink, batch_size=args.batch_size,
|
||
flush_interval=args.flush_interval)
|
||
|
||
flush_latencies: list = [] # 每次 flush 的端到端耗时(秒)
|
||
total_written = 0
|
||
started = time.perf_counter()
|
||
|
||
for seq in range(args.rounds):
|
||
rows = make_rows(args.points, seq)
|
||
writer.push_many(rows)
|
||
t0 = time.perf_counter()
|
||
flushed = writer.flush()
|
||
flush_latencies.append(time.perf_counter() - t0)
|
||
total_written += flushed
|
||
if args.verbose and seq % 20 == 0:
|
||
print(f" [round {seq}] flushed={flushed}")
|
||
|
||
writer.close()
|
||
elapsed = time.perf_counter() - started
|
||
throughput = total_written / elapsed if elapsed > 0 else 0.0
|
||
p99_sec = p99(flush_latencies)
|
||
|
||
print(f" -- 结果 --")
|
||
print(f" 总写入: {total_written} 条,耗时 {elapsed:.3f}s,"
|
||
f"吞吐 {throughput:,.0f} 条/s")
|
||
print(f" flush 批次: {len(flush_latencies)} 次,"
|
||
f"平均 {statistics.mean(flush_latencies)*1000:.2f} ms,"
|
||
f"P99 {p99_sec*1000:.2f} ms")
|
||
print(f" 基准参考: 5k/100ms = 50,000 条/s;目标 P99 ≤ {P99_TARGET_SEC}s")
|
||
|
||
ok = p99_sec <= P99_TARGET_SEC
|
||
print(f" -> P99 {p99_sec*1000:.1f} ms {'PASS' if ok else 'FAIL'}"
|
||
f"(目标 ≤ {P99_TARGET_SEC*1000:.0f} ms)")
|
||
return 0 if ok else 1
|
||
|
||
|
||
def main() -> int:
|
||
parser = argparse.ArgumentParser(description="批量写入压测(P99 ≤ 1.8s)")
|
||
parser.add_argument("--points", type=int, default=600, help="点位数量")
|
||
parser.add_argument("--rounds", type=int, default=100, help="压测轮数")
|
||
parser.add_argument("--batch-size", type=int, default=5000)
|
||
parser.add_argument("--flush-interval", type=float, default=0.1)
|
||
parser.add_argument("--verbose", action="store_true")
|
||
args = parser.parse_args()
|
||
return run_bench(args)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
sys.exit(main())
|