From d37afc0efc8224976cae4e0ef6e578e28dc8c864 Mon Sep 17 00:00:00 2001 From: Derek Feichtinger Date: Tue, 14 Jul 2026 09:55:48 +0200 Subject: [PATCH] code for date based histograms with nice bin widths --- slurm-eff-tool.py | 175 +++++++++++++++++++++++++++++++++++++++++++--- 1 file changed, 167 insertions(+), 8 deletions(-) diff --git a/slurm-eff-tool.py b/slurm-eff-tool.py index 7669970..4e69e11 100755 --- a/slurm-eff-tool.py +++ b/slurm-eff-tool.py @@ -49,6 +49,8 @@ import sys import signal import tempfile import time +from datetime import datetime +from zoneinfo import ZoneInfo from collections import defaultdict from dataclasses import dataclass, field, fields, asdict from pathlib import Path @@ -66,6 +68,7 @@ debug = 0 # Default GB per CPUs for the cluster default_mempercpu_gb = 2.2 +SLURM_TIME_FORMAT = "%Y-%m-%dT%H:%M:%S" SACCT_DELIM = "\x1f" # ASCII Unit Separator SACCT_FIELDS = [ @@ -78,6 +81,7 @@ SACCT_FIELDS = [ "AllocCPUS", "NNodes", "ReqMem", + "Start", "ElapsedRaw", "TimelimitRaw", "CPUTimeRAW", @@ -106,6 +110,7 @@ ALL_COLUMNS = [ "MaxRSS_max", "MemPerCPU", "Mem_Eff", + "Start", "Planned_Time", "Elig_Qtime", "waste_Mem", @@ -257,6 +262,7 @@ HISTO_COLUMNS = { "Elig_Qtime", "Mem_Eff", "Planned_Time", + "Start", "Time_Eff", "UsedMem", "Walltime", @@ -431,6 +437,7 @@ class JobRecord: mem_per_cpu_gb: float | None mem_alloc_tres: float | None mem_used_tres: float | None + start: float | None planned_sec: float reqwall_hours: float | None reqmem_bytes_total: float | None @@ -509,6 +516,7 @@ class OutputRow: AllocMem: float | None UsedMem: float | None MemPerCPU: float | None + Start: float | None ReqWalltime: float | None Planned_Time: float | None Elig_Qtime: float | None @@ -539,6 +547,7 @@ class OutputRow: "UsedMem": self.UsedMem, "MaxRSS_max": self.maxrss_max, "MemPerCPU": self.MemPerCPU, + "Start": self.Start, "ReqWalltime": self.ReqWalltime, "Walltime": self.walltime, "Walltime_max": self.walltime_max, @@ -567,7 +576,7 @@ class OutputRow: def die(msg: str, code: int = 2) -> None: - print(f"error: {msg}", file=sys.stderr) + print(f"ERROR: {msg}", file=sys.stderr) raise SystemExit(code) @@ -597,6 +606,21 @@ def slurm_duration_to_seconds(value: str) -> float: return 0.0 return 0.0 +def slurm_date_to_epoch( + value: str, + # *, + # timezone: ZoneInfo, +) -> float: + """ + Convert a Slurm timestamp to Unix epoch seconds. + + The sacct timestamp has no timezone offset, so `timezone` must describe + the timezone used by the Slurm cluster. + """ + # Since Slurm does not add timezone info, we consider all values as UTC + parsed = datetime.strptime(value, SLURM_TIME_FORMAT) + # localized = parsed.replace(tzinfo=timezone) + return parsed.timestamp() def format_seconds(seconds: int | float | None) -> str: if seconds is None or seconds <= 0: @@ -782,7 +806,11 @@ def read_binary_cache(filename: str) -> dict[str, str or list[JobRecord]]: if result['slurm_eff_version'] != VERSION: sys.stderr.write(f"WARNING: cache was written by version {result['slurm_eff_version']}, but we are running version {VERSION}\n") - result['records'] = [JobRecord(**d) for d in blob["records"]] + try: + result['records'] = [JobRecord(**d) for d in blob["records"]] + except TypeError as e: + sys.stderr.write(f"ERROR reading records from {filename}: {e}\n") + sys.exit(1) return result def show_cache_info(filename: str): @@ -825,7 +853,7 @@ def build_job_records(rows: list[dict[str, str]]) -> list[JobRecord]: if state.startswith("CANCELLED by"): state = "CANCELLED" stateid = state_mappings[state] - + # seff-style practical peak RSS: maximum MaxRSS across non-extern job steps. # This is intentionally not a sum over all ranks/tasks. step_rows = [ @@ -869,6 +897,12 @@ def build_job_records(rows: list[dict[str, str]]) -> list[JobRecord]: planned_sec = float(top.get("PlannedCPURAW") or 0) elig_qtime_sec = planned_sec / alloc_cpus + # Start time by default is in YYYY-MM-DDTHH:MM:SS format. + start = None + startstr = top.get("Start") or None + if startstr: + start = slurm_date_to_epoch(startstr) + # sacct TimelimitRaw is minutes, while ElapsedRaw and CPUTimeRAW are seconds. timelimit_raw_minutes = float(top.get("TimelimitRaw") or 0) timelimit_seconds = int(round(timelimit_raw_minutes * 60)) @@ -930,6 +964,7 @@ def build_job_records(rows: list[dict[str, str]]) -> list[JobRecord]: planned_sec=planned_sec, reqwall_hours=reqwall_hours, reqmem_bytes_total=reqmem_total, + start=start, elapsed_sec=elapsed, elig_qtime_sec=elig_qtime_sec, timelimit_sec=timelimit_seconds, @@ -1078,6 +1113,8 @@ def make_single_row(rec: JobRecord, dflt_mpcpu: float, vectors[colname] = [rec.elapsed_sec/3600] elif colname == "usedmem" and rec.mem_used_tres is not None: vectors[colname] = [rec.mem_used_tres] + elif colname == "start" and rec.start is not None: + vectors[colname] = [rec.start] elif colname == "elig_qtime": vectors[colname] = [rec.elig_qtime_sec/3600] elif colname == "planned_time": @@ -1100,6 +1137,7 @@ def make_single_row(rec: JobRecord, dflt_mpcpu: float, AllocMem=rec.mem_alloc_tres, UsedMem=rec.mem_used_tres, MemPerCPU=rec.mem_per_cpu_gb, + Start=rec.start, ReqWalltime=rec.reqwall_hours, Planned_Time=rec.planned_sec/3600, Elig_Qtime=rec.elig_qtime_sec/3600, @@ -1128,7 +1166,10 @@ def make_aggregate_row(records: list[JobRecord], username: str, partition: str, mem_eff_vals = [r.mem_eff for r in records if r.mem_eff is not None] time_eff_vals = [r.time_eff for r in records if r.time_eff is not None] - walltime=avg_function([r.elapsed_sec for r in records if r.elapsed_sec is not None]) + # an average start date is not very useful, but we do it for consistency + start = avg_function([r.start for r in records if r.start is not None]) + + walltime = avg_function([r.elapsed_sec for r in records if r.elapsed_sec is not None]) if walltime is not None: walltime /= 3600 walltime_max=max([r.elapsed_sec for r in records @@ -1189,6 +1230,9 @@ def make_aggregate_row(records: list[JobRecord], username: str, partition: str, elif colname == "usedmem": vectors[colname] = [r.mem_used_tres for r in records if r.mem_used_tres is not None] + elif colname == "start": + vectors[colname] = [r.start for r in records + if r.start 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] @@ -1213,6 +1257,7 @@ def make_aggregate_row(records: list[JobRecord], username: str, partition: str, AllocMem=alloc_mem, UsedMem=used_mem, MemPerCPU=first.mem_per_cpu_gb, + Start=start, ReqWalltime=first.reqwall_hours, Planned_Time=planned_time_h, Elig_Qtime=elig_qtime_h, @@ -1297,6 +1342,8 @@ def format_value(value: Any, column: str | None = None) -> str: if column in ["ReqWalltime", "Walltime", "Walltime_max", "Elig_Qtime", "Planned_Time"]: return f"{value:.2f}h" + if column in ["Start"]: + return datetime.fromtimestamp(value).strftime("%Y-%m-%dT%H:%M:%S") if isinstance(value, float): return f"{value:.2f}" @@ -1339,7 +1386,6 @@ from math import ceil, exp, floor, log, log10 Metricnumber = Union[int, float] - def percentile(sorted_values: Sequence[float], p: float) -> float: """Return percentile p in [0, 100] using linear interpolation.""" if not sorted_values: @@ -1440,6 +1486,92 @@ def shared_bins( ] +NICE_TIME_BIN_WIDTHS: tuple[int, ...] = ( + 60, # 1 minute + 5 * 60, # 5 minutes + 15 * 60, # 15 minutes + 30 * 60, # 30 minutes + 60 * 60, # 1 hour + 2 * 60 * 60, # 2 hours + 3 * 60 * 60, # 3 hours + 6 * 60 * 60, # 6 hours + 12 * 60 * 60, # 12 hours + 24 * 60 * 60, # 1 day + 2 * 24 * 60 * 60, # 2 days + 7 * 24 * 60 * 60, # 1 week + 14 * 24 * 60 * 60, # 2 weeks + 30 * 24 * 60 * 60, # approximately 1 month + 90 * 24 * 60 * 60, # approximately 3 months + 365 * 24 * 60 * 60, # approximately 1 year +) + + +def nice_time_bin_width( + proposed_width: float, + *, + candidates: tuple[int, ...] = NICE_TIME_BIN_WIDTHS, +) -> int: + """ + Round a proposed width in seconds up to a calendar-friendly width. + """ + if proposed_width <= 0: + raise ValueError("proposed_width must be positive") + + for candidate in candidates: + if candidate >= proposed_width: + return candidate + + # For very large ranges, round upward to whole years. + year = 365 * 24 * 60 * 60 + return int(-(-proposed_width // year) * year) + +def time_bin_width(values: list[float]) -> int: + """ + Find appropriate time bin width from first getting an + estimate by freedaman-diaconis algorithm, then adapt + to nice bin widths. + """ + width = freedman_diaconis_bin_width(values) + logging.debug(f"freedman-diaconis bin width candidate: {width}") + + if width is None: + # Sensible fallback for degenerate or very small data. + return 24 * 60 * 60 + + return nice_time_bin_width(width) + + +def aligned_time_bin_edges( + values: list[float], + width: int, +) -> list[float]: + """ + Generate shared, epoch-aligned time-bin edges. + + Bins are aligned to multiples of `width` relative to the Unix epoch. + This gives natural alignment for widths such as hours and days in UTC. + """ + if not values: + raise ValueError("no values") + + if width <= 0: + raise ValueError("width must be positive") + + start = floor(min(values) / width) * width + stop = ceil(max(values) / width) * width + + # Ensure that the maximum value is enclosed if it lies exactly on `stop`. + if stop <= max(values): + stop += width + + number_of_bins = round((stop - start) / width) + + return [ + float(start + i * width) + for i in range(number_of_bins + 1) + ] + + def histogram_counts( values: Iterable[Metricnumber], edges: Sequence[float], @@ -1469,7 +1601,7 @@ def histogram_counts( return counts -def histogram_table( +def histogram_metric_table( distributions: Sequence[Sequence[Metricnumber]], *, min_bins: int = 8, @@ -1493,12 +1625,27 @@ def histogram_table( return edges, counts +def histogram_date_table(distributions: Sequence[Sequence[Metricnumber]]): + """ + Return shared bin edges and bin counts for each date based distribution. + """ + all_start_times = [timestamp for distr in distributions + for timestamp in distr] + width = time_bin_width(all_start_times) + logging.debug(f"adapting to nice date bin width: {width}") + edges = aligned_time_bin_edges(all_start_times, width) + + counts = [histogram_counts(distribution, edges) + for distribution in distributions] + return edges, counts + def histo_graphs(out_rows: list[OutputRow], args: argparse.Namespace): # dumb term defaults width = 120 height = 30 - outfileroot=None + outfileroot = None term = args.hterm + xax_format = "" plotstyle = "set style fill transparent solid 0.35 noborder" linestyle = "with steps lw 2" @@ -1524,9 +1671,16 @@ def histo_graphs(out_rows: list[OutputRow], args: argparse.Namespace): colnames = args.histo.lower().split(',') for metric in colnames: + # TODO: implement the nicer binning for date based x axis (e.g. for Start time) # create bins try: - edges, counts = histogram_table([getattr(row,"vectors")[metric] for row in out_rows]) + if metric in ["start"]: + edges, counts = histogram_date_table( + [getattr(row,"vectors")[metric] for row in out_rows]) + pass + else: + edges, counts = histogram_metric_table( + [getattr(row,"vectors")[metric] for row in out_rows]) except ValueError as e: sys.stderr.write(f"WARNING: failed to produce histogram for {metric}: {e}\n") return @@ -1560,8 +1714,13 @@ def histo_graphs(out_rows: list[OutputRow], args: argparse.Namespace): logy = "" if args.hlogy: logy = "set logscale y\n" + if metric == "start": + xax_format = (f'set xdata time\n' + f'set timefmt "%s"\n' + f'set format x "%m-%d\\n%H:%M"\n') gnuplot_cmd = (f'set title "{metric} distribution" noenhanced\n' + f'{xax_format}' f'set xlabel "{metric}" noenhanced\n' f'set ylabel "jobs"\n' f'{plotstyle}\n'