built in vector storing machinery for various metrics
This commit is contained in:
+75
-11
@@ -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
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user