diff --git a/slurm-eff-tool.py b/slurm-eff-tool.py index 824c2e5..ecb9154 100755 --- a/slurm-eff-tool.py +++ b/slurm-eff-tool.py @@ -34,6 +34,8 @@ from dataclasses import dataclass, field, fields from pathlib import Path from typing import Any, Iterable +# Default GB per CPUs for the cluster +default_mempercpu_gb = 2.0 SACCT_FIELDS = [ "JobIDRaw", @@ -439,7 +441,8 @@ def is_top_level_row(row: dict[str, str]) -> bool: return "." not in jid and "_" not in jid -def build_job_records(rows: list[dict[str, str]], only_user: str | None = None) -> list[JobRecord]: +def build_job_records(rows: list[dict[str, str]], + only_user: str | None = None) -> list[JobRecord]: """Collapse sacct top-level and step rows into one JobRecord per base job.""" grouped: dict[str, list[dict[str, str]]] = defaultdict(list) for row in rows: @@ -563,6 +566,11 @@ def build_job_records(rows: list[dict[str, str]], only_user: str | None = None) def aggregate_records(records: list[JobRecord], args: argparse.Namespace) -> list[OutputRow]: """Aggregate records according to given grouping instructions.""" + + global default_mempercpu_gb + if args.deflt_mpcpu: + default_mempercpu_gb = float(args.deflt_mpcpu) + if args.aggr_regexp: compiled = [(pat, re.compile(pat)) for pat in args.aggr_regexp] buckets: dict[tuple[Any, ...], list[JobRecord]] = defaultdict(list) @@ -580,8 +588,8 @@ def aggregate_records(records: list[JobRecord], args: argparse.Namespace) -> lis if not matched: unmatched.append(rec) - out = [make_aggregate_row(v, username=k[1], jobname=k[0]) for k, v in buckets.items()] - out.extend(make_single_row(r) for r in unmatched) + out = [make_aggregate_row(v, username=k[1], jobname=k[0], ) for k, v in buckets.items()] + out.extend(make_single_row(r, dflt_mpcpu=default_mempercpu_gb) for r in unmatched) return out if args.aggr_user: @@ -590,18 +598,19 @@ def aggregate_records(records: list[JobRecord], args: argparse.Namespace) -> lis key = (rec.username, rec.reqtasks, rec.cpus, rec.nodes, rec.reqmem_gb, rec.reqwall_hours) buckets[key].append(rec) - return [make_aggregate_row(v, username=k[0], jobname=f"{common_prefix([r.jobname for r in v])}") for k, v in buckets.items()] + return [make_aggregate_row(v, username=k[0], jobname=f"{common_prefix([r.jobname for r in v])}", dflt_mpcpu=default_mempercpu_gb) for k, v in buckets.items()] - return [make_single_row(r) for r in records] + return [make_single_row(r, dflt_mpcpu=default_mempercpu_gb) for r in records] -def make_single_row(rec: JobRecord) -> OutputRow: +def make_single_row(rec: JobRecord, dflt_mpcpu: float) -> OutputRow: """Returns an OutputRow based on a single slurm job record.""" walltime = rec.elapsed_sec / 3600 waste_mem = None if rec.mem_eff is not None and rec.reqmem_gb is not None: - waste_mem = walltime * (100-rec.mem_eff)/100 * rec.reqmem_gb + waste_mem = max(0,walltime * (100-rec.mem_eff)/100 \ + * (rec.reqmem_gb - rec.cpus * dflt_mpcpu)) waste_cpu=None if rec.cpu_eff is not None: @@ -634,7 +643,8 @@ def make_single_row(rec: JobRecord) -> OutputRow: ) -def make_aggregate_row(records: list[JobRecord], username: str, jobname: str) -> OutputRow: +def make_aggregate_row(records: list[JobRecord], username: str, + jobname: str, dflt_mpcpu: float) -> 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] @@ -658,7 +668,8 @@ def make_aggregate_row(records: list[JobRecord], username: str, jobname: str) -> waste_mem = None if (memory_efficiency is not None) and (walltime is not None) \ and first.reqmem_gb is not None: - waste_mem = count * walltime * (100-memory_efficiency)/100 * first.reqmem_gb + waste_mem = max(0,count * walltime * (100-memory_efficiency)/100 \ + * (first.reqmem_gb - first.cpus * dflt_mpcpu)) cpu_efficiency = mean_or_none(cpu_eff_vals) @@ -840,6 +851,7 @@ def parse_args(argv: list[str]) -> argparse.Namespace: help="aggregate jobs matching regexp by regexp, CPUs, nodes, ReqMem, and timelimit; may be repeated", ) + p.add_argument("--deflt-mpcpu", help=f"Default memory to CPU ratio of cluster ({default_mempercpu_gb} GB/cpu)") p.add_argument("--sdev", action="store_true", help="after each efficiency average, add sdev, max, and min columns") p.add_argument("--json", action="store_true", help="emit JSON instead of an ASCII table") p.add_argument("-s", "--sort", help="comma-separated numeric sort columns or aliases; prefix with - for descending") @@ -883,7 +895,7 @@ def main(argv: list[str] | None = None) -> int: write_cache(args.output_cache, rows) records = build_job_records(rows, only_user=args.user) - out_rows = aggregate_records(records, args) + out_rows = aggregate_records(records, args) # , dflt_mpcpu=default_mempercpu_gb out_rows = sort_rows(out_rows, args.sort)