From 9aaa17fa76cca73ddec1839a1f29bf852581922b Mon Sep 17 00:00:00 2001 From: Derek Feichtinger Date: Mon, 6 Jul 2026 17:12:51 +0200 Subject: [PATCH] built in vector storing machinery for various metrics --- slurm-eff-tool.py | 86 +++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 75 insertions(+), 11 deletions(-) diff --git a/slurm-eff-tool.py b/slurm-eff-tool.py index e7697a2..b49c0e8 100755 --- a/slurm-eff-tool.py +++ b/slurm-eff-tool.py @@ -510,7 +510,8 @@ class OutputRow: "jobname": self.jobname, "waste_CPU": self.waste_CPU, "waste_Mem": self.waste_Mem, - "waste_total": self.waste_total + "waste_total": self.waste_total, + "vectors": self.vectors } if sdev: for col, vals in [ @@ -909,11 +910,12 @@ def filter_records(records: list[JobRecord], def aggregate_records(records: list[JobRecord], args: argparse.Namespace) -> list[OutputRow]: """Aggregate records according to given grouping instructions.""" - global default_mempercpu_gb if args.dflt_mpcpu: default_mempercpu_gb = float(args.dflt_mpcpu) + storvecs=args.histo.lower().split(",") + if args.aggr_regexp: compiled = [(pat, re.compile(pat)) for pat in args.aggr_regexp] buckets: dict[tuple[Any, ...], list[JobRecord]] = defaultdict(list) @@ -940,9 +942,11 @@ def aggregate_records(records: list[JobRecord], args: argparse.Namespace) -> lis out = [make_aggregate_row(v, username=k[1], jobname=k[0], \ partition=k[2], - dflt_mpcpu=default_mempercpu_gb) \ + dflt_mpcpu=default_mempercpu_gb, + storvecs=storvecs) \ for k, v in buckets.items()] - out.extend(make_single_row(r, dflt_mpcpu=default_mempercpu_gb) \ + out.extend(make_single_row(r, dflt_mpcpu=default_mempercpu_gb, + storvecs=storvecs) \ for r in unmatched) return out @@ -958,12 +962,17 @@ def aggregate_records(records: list[JobRecord], args: argparse.Namespace) -> lis rec.reqmem_gb, rec.reqwall_hours) buckets[key].append(rec) - return [make_aggregate_row(v, username=k[0],partition=k[1],jobname=f"{common_prefix([r.jobname for r in v])}", dflt_mpcpu=default_mempercpu_gb) for k, v in buckets.items()] + return [make_aggregate_row(v, username=k[0],partition=k[1], \ + jobname=f"{common_prefix([r.jobname for r in v])}", + dflt_mpcpu=default_mempercpu_gb, storvecs=storvecs) + for k, v in buckets.items()] - return [make_single_row(r, dflt_mpcpu=default_mempercpu_gb) for r in records] + return [make_single_row(r, dflt_mpcpu=default_mempercpu_gb, + storvecs=storvecs) for r in records] -def make_single_row(rec: JobRecord, dflt_mpcpu: float) -> OutputRow: +def make_single_row(rec: JobRecord, dflt_mpcpu: float, + storvecs: list[str] | None = None) -> OutputRow: """Returns an OutputRow based on a single slurm job record.""" walltime = rec.elapsed_sec / 3600 @@ -987,6 +996,29 @@ def make_single_row(rec: JobRecord, dflt_mpcpu: float) -> OutputRow: jobid = rec.jobidraw if jobid == "": jobid = rec.jobid + + vectors: dict[str, list[int | float]] = { + "cpu_eff": [rec.cpu_eff] if rec.cpu_eff is not None else [], + "mem_eff": [rec.mem_eff] if rec.mem_eff is not None else [], + "time_eff": [rec.time_eff] if rec.time_eff is not None else [] + } + + if storvecs: + for colname in storvecs: + vectors[colname] = [] + if colname == "walltime": + vectors[colname] = [rec.elapsed_sec/3600] + elif colname == "maxrss" and rec.maxrss_bytes is not None: + vectors[colname] = [rec.maxrss_bytes / (1024**3)] + elif colname == "used_mem" and rec.mem_used_tres is not None: + vectors[colname] = [rec.mem_used_tres] + elif colname == "elig_qtime": + vectors[colname] = [rec.elig_qtime_sec/3600] + elif colname == "planned_time": + vectors[colname] = [rec.planned_sec/3600] + else: + sys.stderr.write(f"ERROR: Cannot gather vector information for metric {colname}\n") + sys.exit(1) return OutputRow( username=rec.username, @@ -1014,14 +1046,13 @@ def make_single_row(rec: JobRecord, dflt_mpcpu: float) -> OutputRow: waste_Mem=waste_mem, waste_CPU=waste_cpu, waste_total=waste_total, - vectors = {"cpu_eff": [rec.cpu_eff] if rec.cpu_eff is not None else [], - "mem_eff": [rec.mem_eff] if rec.mem_eff is not None else [], - "time_eff": [rec.time_eff] if rec.time_eff is not None else []}, + vectors = vectors ) def make_aggregate_row(records: list[JobRecord], username: str, partition: str, - jobname: str, dflt_mpcpu: float) -> OutputRow: + jobname: str, dflt_mpcpu: float, + storvecs: list[str] | None = None) -> OutputRow: """Returns an OutputRow based on the given list of job records.""" first = records[0] cpu_eff_vals = [r.cpu_eff for r in records if r.cpu_eff is not None] @@ -1077,6 +1108,29 @@ def make_aggregate_row(records: list[JobRecord], username: str, partition: str, "mem_eff": [r.mem_eff for r in records if r.mem_eff is not None], "time_eff": [r.time_eff for r in records if r.time_eff is not None], } + if storvecs: + for colname in storvecs: + if colname == "walltime": + vectors[colname] = [r.elapsed_sec/3600 for r in records + if r.elapsed_sec is not None] + elif colname == "maxrss": + vectors[colname] = [r.maxrss_bytes / (1024**3) for r in records + if r.maxrss_bytes is not None] + elif colname == "used_mem": + vectors[colname] = [r.mem_used_tres for r in records + if r.mem_used_tres is not None] + elif colname == "elig_qtime": + vectors[colname] = [r.elig_qtime_sec/3600 for r in records + if r.elig_qtime_sec is not None] + elif colname == "planned_time": + vectors[colname] = [r.planned_sec/3600 for r in records + if r.planned_sec is not None] + else: + sys.stderr.write(f"ERROR: Cannot gather vector information for metric {colname}\n") + sys.exit(1) + #planned_time, elig_qtime, + # vectors[colname] = [getattr(r,colname) for r in records + # if getattr(r,colname) is not None] return OutputRow( @@ -1259,6 +1313,7 @@ def parse_args(argv: list[str]) -> argparse.Namespace: p.add_argument("-O", "--output-raw", help="write raw sacct output cache to this file (text format, large).") p.add_argument("-F", "--from-raw", help="read raw sacct output cache from this file.") p.add_argument("-i", "--info", help="show metadata information for the given binary cache file") + p.add_argument("-H", "--histo", help="print histogram of the given metric.", default=None) p.add_argument("--dflt-mpcpu", help=f"Default memory/CPU ratio of cluster [{default_mempercpu_gb} GB/cpu]. Used in memory waste calculation.") p.add_argument("--sdev", action="store_true", help="after each efficiency average, add sdev, max, and min columns") @@ -1338,6 +1393,7 @@ def main(argv: list[str] | None = None) -> int: if args.write_binary_cache: write_binary_cache(records, args.write_binary_cache) + # Aggregate out_rows = aggregate_records(records, args) # , dflt_mpcpu=default_mempercpu_gb if args.expr: expr = SafeExpression(args.expr) @@ -1346,11 +1402,19 @@ def main(argv: list[str] | None = None) -> int: out_rows = sort_rows(out_rows, args.sort) + # Output table if args.json: print(json.dumps([r.as_dict(sdev=args.sdev) for r in out_rows], indent=2)) else: print_table(out_rows, output_columns, args.sdev) + if args.histo: + for row in out_rows: + rowdict = row.as_dict() + for colname, vals in rowdict["vectors"].items(): + print(f'DEBUG #### {colname}: {vals}') + #print(f'{rowdict["username"]: {rowdict["vectors"][args.histo]}}') + return 0