changed formula for waste_mem to account for default mem per CPU
This commit is contained in:
+22
-10
@@ -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)
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user