[[PageOutline]] = Job Checkpointing Examples = == Overview == This example implements a '''submitter/runner''' (also called cyclic checkpoint/restart) pattern. === Submitter -> Runner -> Application === A lightweight '''submitter''' job launches a '''runner''' job that in turn launches the '''application''' for a bounded time. After that time, the application writes a checkpoint, and exits cleanly. === Submitter self-resubmits before timing out === The submitter itself also runs for a bounded time while monitoring progress. Just before timing out, the submitter resubmits a new submitter job as needed until the workflow completes using the '''sbatch''' command rather than the SLURM '''requeue''' feature. For details on this checkpointing implementation, see [wiki://Workshops/JobCheckpointing/Examples#Checkpointworkflowimplementation Checkpoint workflow implementation] below. === Why do we use sbatch vs SLURM requeue? === ==== Unacceptable use of requeue ==== The SLURM resource manager provides an option to requeue a running job via the command '''scontrol requeue'''. (See '''man scontrol'''.) However, when a job is requeued, any previous record of elapsed time used by that same job is '''overwritten'''. Hence, the requeue option has the following 2 consequences. 1. You would '''not''' be able to track the total elapsed time used by your requeued job. 2. As a result, your job would be in violation of Tulane Information Technology's [https://it.tulane.edu/acceptable-use-policy Acceptable Use Policy], including your '''consent to monitor''' your use of Cypress. == Before you adopt checkpointing == '''Before you adopt a checkpointing workflow''', make a determination of whether your workflow consists - either entirely or in large part - of data elements or other components that can be processed '''entirely independently''' of each other - rather than each element depending on the result of the previous. This is also referred to as '''embarrassingly parallel'''. (See also [https://www.freecodecamp.org/news/embarrassingly-parallel-algorithms-explained-with-examples/ Embarrassingly Parallel Algorithms Explained].) In order to make the above determination, you may need to '''consult your software provider's user forum''', such as [https://github.com GitHub]. In case your workflow is '''embarrassingly parallel''', then please first refer to the material in [wiki:Workshops/JobParallelism Job Parallelism]. == Quick start == For a quick start using checkpointing, see [wiki:Workshops/JobCheckpointing/Examples#WorkingExamples Working Examples] below. == Background on checkpointing == === Use checkpointing for sequential workflows === This checkpointing implementation should be used for running '''sequential workflows''' on Cypress - where each component result of the workflow depends the previous component result in the workflow. (See above [wiki:Workshops/JobCheckpointing/Examples#Beforeyouadoptcheckpointing Before you adopt checkpointing].) === Checkpointing - an HPC workflow design pattern === This checkpointing implementation can be considered an implementation of an HPC workflow design pattern (see for example [https://github.com/deepakkum21/Books/blob/master/Design%20Patterns%20-%20Elements%20of%20Reusable%20Object%20Oriented%20Software%20-%20GOF.pdf Design Patterns]) that can be described by any of the following expressions. * submitter/runner pattern (Again see [https://github.com/deepakkum21/Books/blob/master/Design%20Patterns%20-%20Elements%20of%20Reusable%20Object%20Oriented%20Software%20-%20GOF.pdf Design Patterns].) * job chaining (See also [https://gjbex.github.io/Workflows-for-HPC/ Workflow for HPC].) * sequential workflows (See [https://docs.nersc.gov/jobs/best-practices/ NERSC - Best Practices for Jobs].) * restart workflows (See [https://support.ceci-hpc.be/doc/SubmittingJobs/WorkflowManagement/#introduction Wide,Deep, & Cyclic Workflows] and [https://nlsq.readthedocs.io/en/latest/tutorials/routine/three_workflows/hpc_workflow.html Checkpointing implemented in Python].) * checkpoint/restart job series (See [https://compendium.hpc.tu-dresden.de/jobs_and_resources/checkpoint_restart/ Checkpoint Motivation].) * using job dependencies to restart jobs that checkpoint their state (See [https://github.com/hongthai-lenguyen/Slurm_HPC_Guide/blob/main/slurm_guide.md#52-creating-pipelines-with-job-dependencies Creating Pipelines with (SLURM) Job Dependencies].) === How does the workflow determine when to stop? === ==== Limiting the work: using the sbatch --dependencies option ==== Note that in the context of the last reference in the above list - using job dependencies, when you use the SLURM '''sbatch --dependencies...''' option, the stopping point is not determined automatically. You may need '''additional effort''' to determine the precise number of '''sbatch --dependencies...''' commands to exiecute - in other words, to determine the length of the dependency chain - before starting the workflow. ==== Limiting the work: using the application's own work limit ==== This implementation differs from the above dependency chain, in that rather than using the '''sbatch --dependencies''' option, the stopping point is determined automatically during runtime. In this implementation, the workflow self restarts or resumes as needed and self terminates when the application's work limit is reached - in other words when the application's work is done. See for example the setting '''MAX_ITER''' in [wiki://Workshops/JobCheckpointing/Examples#checkpoint_submitter.sh checkpoint_submitter.sh] below. == Checkpoint workflow implementation == === Execution Model (Diagram) === {{{ +-----------+ +-----------+ +---------------+ | Submitter | | Runner | | Application | | (job) | | (job) | | (process) | +-----------+ +-----------+ +---------------+ | | | | if Application not done, | | | submit(sbatch) | | | ------------------------> | | | otherwise, exit | run | | | ------------------------> | | | | | ----+ | SIGTERM @ timeout | | | wait for self | ------------------------> | | | timeout | | ----+ | | | | | | | | | | handle SIGTERM | | | | | & write checkpoint | | | | <---+ | | | exit code | | | + <-------------------------+ | | Runner exit (progress / done / failure) | <---+ | | update workflow state | submit new Submitter | interpret Runner exit code + exit (Submitter state persisted e.g. to submitter_state.txt) }}} === Components explained === In this implementation, the checkpointing workflow consists of the following components. 1. '''the checkpoint submitter''' - to perform the following a. self restart as needed, and b. re-submit the runner once per self-restart 2. '''the checkpoint runner''' - to perform the following a. set its own application timer via the '''timeout''' command b. re-start the application as needed c. send the signal SIGTERM to the application at timeout 3. '''the checkpoint application''' - to perform the following a. compute b. catch the runner's signal SIGTERM c. "handle" the runner's signal SIGTERM by quickly saving the work and stop running (See [wiki:Workshops/JobCheckpointing/Examples#WorkingExamples Working Examples] for applications in BASH, Python and R.) === Checkpoint Submitter === Here is the job script for the '''checkpoint submitter''', a fully self-restarting job that controls the workflow by self-restarting as many times as needed and submitting one instance per each restart of the checkpoint runner job script. (See [wiki:Workshops/JobCheckpointing/Examples#CheckpointRunner Checkpoint Runner] further below.) ==== Overriding default settings ==== This checkpoint submitter uses settings - with default values - implemented using environment variables (see [wiki:cypress/WorkingOnClusters#EnvironmentVariables Environment Variables]), where the default values can be overridden by either of the following. 1. Override default settings by one or more '''export''' commands prior to the '''sbatch''' command {{{ # example: export all on one line, then submit [tulaneId@cypress1 ~]$export VAR1= VAR2= ... [tulaneId@cypress1 ~]$sbatch checkpoint_submitter.sh }}} {{{ # example: or export one line for each variable, then submit [tulaneId@cypress1 ~]$export VAR1= [tulaneId@cypress1 ~]$export VAR2= ... [tulaneId@cypress1 ~]$sbatch checkpoint_submitter.sh }}} 2. Alternatively, override settings by assigning values at the beginning, in front of the '''sbatch''' command {{{ # example: all on one line with submit [tulaneId@cypress1 ~]$VAR1= VAR2= ... sbatch checkpoint_submitter.sh }}} {{{ # example: on multiple lines using the line continuation character "\" with submit [tulaneId@cypress1 ~]$VAR1= \ VAR2= \ ... \ sbatch checkpoint_submitter.sh }}} ==== checkpoint_submitter.sh ==== {{{ #!/bin/bash ############################################################################### # CHECKPOINT SUBMITTER: self-restarting to re-submit runner # Controls repeated submissions of checkpoint_runner.sh ############################################################################### # User settings with defaults ############################################################################### # CKPT_PATH : path to checkpoint file # MAX_ITER : stop application after this many iterations total # RUNNER_SCRIPT : job script with checkpoint before end of time limit # SUBMITTER_MARGIN_SEC : time margin to restart before time limit # SUBMITTER_TIME_LIMIT : wall time <= 24 hours (format as --time in man sbatch) ############################################################################### # Edit settings below # Or pass values via... # sbatch --export=ALL,Var1=...,...,VarN=... checkpoint_submitter.sh # Or pass values via... # Var1=... VarN=... sbatch checkpoint_submitter.sh ############################################################################### CKPT_PATH="${CKPT_PATH:-state_iter.txt}" MAX_ITER="${MAX_ITER:-500}" RUNNER_SCRIPT="${RUNNER_SCRIPT:-checkpoint_runner.sh}" SUBMITTER_MARGIN_SEC="${SUBMITTER_MARGIN_SEC:-60}" SUBMITTER_TIME_LIMIT="${SUBMITTER_TIME_LIMIT:-3:30}" ############################################################################### # Slurm directives (keep SUBMITTER_TIME_LIMIT <= 24 hours) ############################################################################### #SBATCH --job-name=ckpt_submitter #SBATCH --output=log_submitter_%j.out #SBATCH --error=log_submitter_%j.err #SBATCH --time=24:00:00 #SBATCH --ntasks=1 #SBATCH --cpus-per-task=1 #SBATCH --mem=256M #SBATCH --partition=centos7 #SBATCH --qos=normal ############################################################################### # Exit on error ############################################################################### set -euo pipefail ############################################################################### # Helpers ############################################################################### shopt -s expand_aliases alias dtstamp='date +%Y%m%d-%H:%M:%S' info(){ echo "Info[$(dtstamp)]: $*"; } get_ckpt_iter() { [[ -f "${CKPT_PATH}" ]] || { echo 0; return; } tr -cd '0-9' < "${CKPT_PATH}" 2>/dev/null || echo 0 } submit_worker() { sbatch "${RUNNER_SCRIPT}" | awk '{print $4}' } to_seconds() { # Acceptable time formats include "minutes", "minutes:seconds", "hours:min- utes:seconds", # "days-hours", "days-hours:minutes" and "days-hours:minutes:seconds" local input="$1" # Handle days-hours format by replacing '-' with ':' to unify the delimiter # Then split components into an array IFS=':' read -r -a parts <<< "${input//-/:}" case ${#parts[@]} in 1) # Format: minutes seconds=$((parts[0] * 60)) ;; 2) # Formats: minutes:seconds OR days-hours if [[ $input == *"-"* ]]; then seconds=$((parts[0] * 86400 + parts[1] * 3600)) else seconds=$((parts[0] * 60 + parts[1])) fi ;; 3) # Formats: hours:minutes:seconds OR days-hours:minutes if [[ $input == *"-"* ]]; then seconds=$((parts[0] * 86400 + parts[1] * 3600 + parts[2] * 60)) else seconds=$((parts[0] * 3600 + parts[1] * 60 + parts[2])) fi ;; 4) # Format: days-hours:minutes:seconds seconds=$((parts[0] * 86400 + parts[1] * 3600 + parts[2] * 60 + parts[3])) ;; *) echo "Unsupported format" exit 1 ;; esac echo "$seconds" } should_restart() { now=$(date +%s) elapsed=$(( now - restart_start_sec )) remaining=$(( restart_limit_sec - elapsed )) echo -ne "Time before restart: $remaining (sec)\r" >&2 (( remaining < SUBMITTER_MARGIN_SEC )) } ############################################################################### # Log settings ############################################################################### info "Start on $(hostname); JOBID=${SLURM_JOB_ID:-unknown}" info "Submitter settings:" info "CKPT_PATH=${CKPT_PATH}" info "MAX_ITER=${MAX_ITER}" info "RUNNER_SCRIPT=${RUNNER_SCRIPT}" info "SUBMITTER_MARGIN_SEC=${SUBMITTER_MARGIN_SEC}" info "SUBMITTER_TIME_LIMIT=${SUBMITTER_TIME_LIMIT}" ############################################################################### # Restore state if present ############################################################################### state_file="submitter_state.txt" submitter_cycle=0 if [[ -f "${state_file}" ]]; then source "${state_file}" info "Loaded state: submitter_cycle=${submitter_cycle}" else info "No previous state found; starting fresh." fi ############################################################################### # Track time for restart ############################################################################### restart_start_sec=$(date +%s) restart_limit_sec=$(to_seconds "${SUBMITTER_TIME_LIMIT}") ############################################################################### # Main Loop ############################################################################### info "Submitter cycle: ${submitter_cycle}" module load anaconda3/2023.07 # using python3 while true; do current_iter=$(get_ckpt_iter) info "Current checkpoint iteration: ${current_iter}" if (( current_iter >= MAX_ITER )); then info "Reached MAX_ITER=${MAX_ITER}; running final analytics." # List of analytics scripts to check analytics_scripts=( "aggregate_usage.py" "plot_progress.py" ) # Run resource usage aggregation if command -v python3 >/dev/null 2>&1; then # Loop over each script and run only if present for script in "${analytics_scripts[@]}"; do if [[ -f "${script}" ]]; then info "Running analytics ${script} ..." python3 "${script}" else info "WARNING: ${script} not found; skipping." fi done else info "WARNING: python3 not found; skipping analytics." fi info "Workflow complete. Exiting submitter." exit 0 fi JOBID=$(submit_worker) info "Submitted worker: ${JOBID}" # Log worker JobIDs for later aggregation id_log="worker_ids.txt" echo "${JOBID}" >> "${id_log}" # Wait for worker to finish info "Waiting on first of worker ${JOBID} to finish or self-restart." while squeue -j "${JOBID}" >/dev/null 2>&1; do sleep 10 if should_restart; then info "Submitter approaching timeout -> restarting self" # Save state echo "submitter_cycle=$((submitter_cycle+1))" > "${state_file}" # Submit a new submitter instance sbatch checkpoint_submitter.sh exit 0 fi done # Read worker exit code rc=$(sacct -j "${JOBID}" --format=ExitCode -P | awk -F: 'NR==2{print $1}') info "Worker finished with code ${rc}" case "${rc}" in 0) info "Worker reports completed workflow." exit 0 ;; 99|124) info "Checkpoint advance -> continuing workflow." ;; *) info "Unexpected worker failure (code ${rc}). Stopping." exit "${rc}" ;; esac # Loop continues to submit next worker done }}} === Checkpoint Runner === Here is the job script for the '''checkpoint runner''' to be used with the following. * for jobs submitted by [wiki:Workshops/JobCheckpointing/Examples#CheckpointSubmitter Checkpoint Submitter] above * for running any of the following, minimal working examples of checkpointed applications in BASH, Python, or R. * For BASH, see [wiki:Workshops/JobCheckpointing/Examples#BASHCheckpointingExample BASH Checkpointint Example] * For Python, see [wiki:Workshops/JobCheckpointing/Examples#PythonCheckpointingExample Python Checkpointint Example] * For R, see [wiki:Workshops/JobCheckpointing/Examples#RCheckpointingExample R Checkpointint Example] ==== checkpoint_runner.sh ==== {{{ #!/bin/bash ############################################################################### # CHECKPOINT RUNNER: run application w/ checkpoint + controlled stop # Use with checkpoint_submitter.sh to re-submit as needed ############################################################################### # User settings with defaults ############################################################################### # APP_CMD : command to run (string) - must handle signal SIGTERM # CHECKPOINT_EVERY : checkpoint after every this many application iterations # CKPT_PATH : path to checkpoint file # LAUNCH_MODE : run app command directly or via srun # MAX_ITER : stop application after this many iterations total # MODULE_LIST : modules loaded before running app command # RUNNER_MARGIN_SEC: seconds before wall time to checkpoint (for timeout) # RUNNER_TIME_LIMIT: wall time <= 24 hours (format as --time in man sbatch) # SRUN_ARGS : arguments to srun, if needed ############################################################################### # Edit settings below # Or pass values via... # sbatch --export=ALL,Var1=...,...,VarN=... checkpoint_submitter.sh # Or pass values via... # Var1=... VarN=... sbatch checkpoint_submitter.sh ############################################################################### APP_CMD="${APP_CMD:-python3 checkpoint_signal_iter.py}" CHECKPOINT_EVERY="${CHECKPOINT_EVERY:-20}" CKPT_PATH="${CKPT_PATH:-state_iter.txt}" LAUNCH_MODE="${LAUNCH_MODE:-direct}" MAX_ITER="${MAX_ITER:-500}" MODULE_LIST="${MODULE_LIST:-anaconda3/2023.07}" RUNNER_MARGIN_SEC="${RUNNER_MARGIN_SEC:-60}" RUNNER_TIME_LIMIT="${RUNNER_TIME_LIMIT:-00:03:00}" SRUN_ARGS="${SRUN_ARGS:--n 1}" ############################################################################### # Slurm directives (keep RUNNER_TIME_LIMIT <= 24 hours) ############################################################################### #SBATCH --job-name=ckpt_runner_demo #SBATCH --output=log_runner_%j.out #SBATCH --error=log_runner_%j.err #SBATCH --time=00:03:00 #SBATCH --ntasks=1 #SBATCH --cpus-per-task=1 #SBATCH --mem=512M #SBATCH --partition=centos7 #SBATCH --qos=normal ############################################################################### # Exit on error ############################################################################### set -euo pipefail ############################################################################### # Helpers ############################################################################### shopt -s expand_aliases alias dtstamp="date +%Y%m%d-%H:%M:%S" info(){ echo "Info[$(dtstamp)]: $*"; } get_ckpt_iter() { [[ -f "${CKPT_PATH}" ]] || { echo 0; return; } tr -cd '0-9' < "${CKPT_PATH}" 2>/dev/null || echo 0 } to_seconds() { # Acceptable time formats include "minutes", "minutes:seconds", "hours:min- utes:seconds", # "days-hours", "days-hours:minutes" and "days-hours:minutes:seconds" local input="$1" # Handle days-hours format by replacing '-' with ':' to unify the delimiter # Then split components into an array IFS=':' read -r -a parts <<< "${input//-/:}" case ${#parts[@]} in 1) # Format: minutes seconds=$((parts[0] * 60)) ;; 2) # Formats: minutes:seconds OR days-hours if [[ $input == *"-"* ]]; then seconds=$((parts[0] * 86400 + parts[1] * 3600)) else seconds=$((parts[0] * 60 + parts[1])) fi ;; 3) # Formats: hours:minutes:seconds OR days-hours:minutes if [[ $input == *"-"* ]]; then seconds=$((parts[0] * 86400 + parts[1] * 3600 + parts[2] * 60)) else seconds=$((parts[0] * 3600 + parts[1] * 60 + parts[2])) fi ;; 4) # Format: days-hours:minutes:seconds seconds=$((parts[0] * 86400 + parts[1] * 3600 + parts[2] * 60 + parts[3])) ;; *) echo "Unsupported format" exit 1 ;; esac echo "$seconds" } ############################################################################### # Log settings ############################################################################### info "Start on $(hostname); JOBID=${SLURM_JOB_ID:-unknown}" info "Runner settings:" info "APP_CMD=${APP_CMD}" info "CHECKPOINT_EVERY=${CHECKPOINT_EVERY}" info "CKPT_PATH=${CKPT_PATH}" info "LAUNCH_MODE=${LAUNCH_MODE}" info "MAX_ITER=${MAX_ITER}" info "MODULE_LIST=${MODULE_LIST}" info "RUNNER_MARGIN_SEC=${RUNNER_MARGIN_SEC}" info "RUNNER_TIME_LIMIT=${RUNNER_TIME_LIMIT}" info "SRUN_ARGS=${SRUN_ARGS}" ############################################################################### # Load requested modules ############################################################################### info "Loading modules." module load "${MODULE_LIST}" ############################################################################### # Timeout window ############################################################################### TOTAL_SEC=$(to_seconds "${RUNNER_TIME_LIMIT}") RUN_WINDOW_SEC=$(( TOTAL_SEC - RUNNER_MARGIN_SEC )) if (( RUN_WINDOW_SEC <= 0 )); then info "WARNING: RUN_WINDOW_SEC <= 0; using 1s." RUN_WINDOW_SEC=1 fi before_iter=$(get_ckpt_iter) set +e if [[ "${LAUNCH_MODE}" == "srun" ]]; then # timeout -> srun -> bash -lc "APP_CMD" # On expiry, timeout sends TERM to srun; srun forwards signals to its job step tasks. timeout "${RUN_WINDOW_SEC}s" srun ${SRUN_ARGS} bash -lc "${APP_CMD}" else # direct mode: run in the batch shell timeout "${RUN_WINDOW_SEC}s" bash -lc "${APP_CMD}" fi rc=$? set -e info "Program exit code (from timeout wrapper): ${rc}" ############################################################################### # Interpret exit-code from timeout # 0 = completed all iterations (no further submissions necessary) # 99 = app voluntarily checkpointed early -> submitter should continue # 124 = timeout before walltime expiry -> checkpoint likely written -> continue ############################################################################### if [[ ${rc} -eq 0 ]]; then info "Application finished all iterations." exit 0 elif [[ ${rc} -eq 99 ]]; then info "Application exited after writing checkpoint." exit 99 elif [[ ${rc} -eq 124 ]]; then after_iter=$(get_ckpt_iter) if (( after_iter > before_iter )); then info "Timeout + checkpoint advance detected (${before_iter}->${after_iter})." exit 124 else info "Timeout but checkpoint did not advance. Marking failure." exit 1 fi else info "Unexpected exit code: ${rc}" exit "${rc}" fi }}} === Post-workflow analysis in Python === ==== Plotting execution time and memory usage ==== Once the last in the series of [wiki://Workshops/JobCheckpointing/Examples#checkpoint_submitter.sh checkpoint_submitter.sh] jobs has detected the end of the computational workflow, it will call on the [wiki:///Workshops/JobCheckpointing/Examples#Pythonscripts Python scripts] listed below to perform a post-worflow analysis in order to plot execution time and memory usage (see output fields '''Elapsed''' and '''MaxRSS''' in '''man sacct''') and identify possible anomalies from among all of the [wiki://Workshops/JobCheckpointing/Examples#checkpoint_runner.sh checkpoint_runner] jobs that have participated in the workflow. The types of anomalies determined, if any are the following. * | Z | > 3.5 for single job execution time ('''Elapsed''' in the above), where {{{ Z = (value - mean) / (standard deviation) }}} * modZ > 3.5 for maximum memory usage by a single job ('''MaxRSS''' in the above), where {{{ modZ = 0.6745 * | value - median | / (Median Absolute Deviation) }}} See also [https://www.itl.nist.gov/div898/handbook/eda/section3/eda35h.htm NIST - Detection of Outliers]. ==== Python scripts ==== Here are two Python scripts used by the checkpoint submitter to analyze and plot usage and annotate anomalies in execution time and memory usage by the checkpoint runner jobs. ==== Aggregate usage ==== Here is the listing for the file '''aggregate_usage.py''' used in [wiki://Workshops/JobCheckpointing/Examples#checkpoint_submitter.sh checkpoint_submitter.sh]. ===== aggregate_usage.py ===== {{{ #!/usr/bin/env python3 """ aggregate_usage.py (SLURM 14.x–compatible) ------------------------------------------------------------- Combines parent and .batch accounting rows so MaxRSS is preserved. SLURM 14.x specifics: - Only .batch steps contain correct MaxRSS - Parent jobs often have empty MaxRSS - ElapsedRaw is not available - Elapsed may be MM, MM:SS, HH:MM:SS, D-HH, D-HH:MM, D-HH:MM:SS Outputs worker_usage.csv: JobID, ElapsedSeconds, TotalCPU, MaxRSS_bytes, MaxRSS_MiB, ReqMem, NodeList, State """ import subprocess import pandas as pd import sys import re from io import StringIO IDFILE = "worker_ids.txt" OUTFILE = "worker_usage.csv" # --------------------------------------------------------- # Utilities # --------------------------------------------------------- def run_cmd(cmd): return subprocess.check_output(cmd, shell=True, text=True) def parse_elapsed_to_seconds(elapsed): if not isinstance(elapsed, str): return None s = elapsed.strip() days = 0 # Day prefix if "-" in s: d, rest = s.split("-", 1) try: days = int(d) except: return None else: rest = s parts = rest.split(":") if len(parts) == 1: # MM or D-HH try: x = int(parts[0]) except: return None if days > 0: hours, minutes, seconds = x, 0, 0 else: hours, minutes, seconds = 0, x, 0 elif len(parts) == 2: # MM:SS or D-HH:MM try: a = int(parts[0]) b = int(parts[1]) except: return None if days > 0: hours, minutes, seconds = a, b, 0 else: hours, minutes, seconds = 0, a, b elif len(parts) == 3: # HH:MM:SS or D-HH:MM:SS try: hours = int(parts[0]) minutes = int(parts[1]) seconds = int(parts[2]) except: return None else: return None return days*86400 + hours*3600 + minutes*60 + seconds def parse_rss_to_bytes(val): if val is None or pd.isna(val): return None s = str(val).strip().rstrip("nN") try: if s.endswith(("K","k")): return float(s[:-1]) * 1e3 if s.endswith(("M","m")): return float(s[:-1]) * 1e6 if s.endswith(("G","g")): return float(s[:-1]) * 1e9 if re.fullmatch(r"\d+(\.\d+)?", s): return float(s) except: return None return None # --------------------------------------------------------- # Main # --------------------------------------------------------- def main(): try: with open(IDFILE) as f: jobids = [x.strip() for x in f if x.strip()] except: print(f"ERROR: {IDFILE} not found.") sys.exit(1) if not jobids: print("No job IDs to process.") sys.exit(0) fields = [ "JobID", "Elapsed", "TotalCPU", "MaxRSS", "ReqMem", "NodeList", "State" ] cmd = f"sacct -j {','.join(jobids)} --format={','.join(fields)} -P" raw = run_cmd(cmd) df = pd.read_csv(StringIO(raw), sep="|") df.dropna(how="all", inplace=True) # Split BaseID with explicit keywords (older pandas compatibility) df["BaseID"] = ( df["JobID"].astype(str) .str.split(pat=".", n=1, regex=False) .str[0] ) df["IsBatch"] = df["JobID"].astype(str).str.endswith(".batch") parent = df[~df["IsBatch"]].copy() batch = df[df["IsBatch"]].copy() merged = pd.merge( parent, batch[["BaseID","MaxRSS"]], on="BaseID", how="left", suffixes=("", "_batch") ) def choose(x): if pd.notna(x["MaxRSS"]): return x["MaxRSS"] return x["MaxRSS_batch"] merged["MaxRSS"] = merged.apply(choose, axis=1) merged["MaxRSS_bytes"] = merged["MaxRSS"].apply(parse_rss_to_bytes) merged["MaxRSS_MiB"] = merged["MaxRSS_bytes"] / (1024.0**2) merged["ElapsedSeconds"] = merged["Elapsed"].apply(parse_elapsed_to_seconds) out = merged[[ "BaseID","ElapsedSeconds","TotalCPU","MaxRSS","MaxRSS_bytes", "MaxRSS_MiB","ReqMem","NodeList","State" ]].rename(columns={"BaseID": "JobID"}) out.to_csv(OUTFILE, index=False) print(f"Wrote {OUTFILE}") print(out.head()) if __name__ == "__main__": main() }}} ==== Plot progress ==== Here is the listing for the file '''plot_progress.py''' used in [wiki://Workshops/JobCheckpointing/Examples#checkpoint_submitter.sh checkpoint_submitter.sh] for generating .png plot files for run time and memory usage by the checkpoint runner jobs. ===== plot_progress.py ===== {{{ #!/usr/bin/env python3 """ plot_progress.py Produces: - progress_with_anomalies.png (dual-axis: checkpoints + MaxRSS MiB) - resource_trend.png (dual-axis: ElapsedSeconds + MaxRSS MiB) - anomaly_details.csv (per-job anomalies) - anomaly_stats_table.png (summary table) """ import pandas as pd import matplotlib.pyplot as plt import datetime as dt import numpy as np import glob import re import os USAGE_FILE = "worker_usage.csv" LOG_GLOB = "log_runner_*.out" # --------------------------------------------------------- # Utilities # --------------------------------------------------------- def parse_worker_logs(): rows = [] for fp in glob.glob(LOG_GLOB): m = re.search(r"log_runner_(\d+)\.out", os.path.basename(fp)) if not m: continue jid = int(m.group(1)) with open(fp, "r", errors="ignore") as f: for line in f: if "Info[" in line: try: ts = line.split("Info[")[1].split("]")[0] t = dt.datetime.strptime(ts,"%Y%m%d-%H:%M:%S") rows.append((jid,t)) break except: pass return pd.DataFrame(rows, columns=["JobID","StartTime"]) if rows else \ pd.DataFrame(columns=["JobID","StartTime"]) def numeric_jobid(s): m = re.match(r"^(\d+)", str(s)) return int(m.group(1)) if m else None def time_anomalies(elapsed): s = elapsed.astype(float) if s.std(ddof=0)==0: return pd.Series(False,index=s.index), pd.Series(0.0,index=s.index) z = (s-s.mean())/s.std(ddof=0) return (z.abs()>3.5), z def mem_anomalies(mem_mib, window=10, thresh=3.5): s = mem_mib.astype(float) med = s.rolling(window,min_periods=3).median() absd = (s-med).abs() mad = absd.rolling(window,min_periods=3).median() modz = 0.6745 * absd / mad.replace(0,np.nan) modz = modz.fillna(0) return (modz>thresh), modz def save_summary(summary, details, fname="anomaly_stats_table.png"): fig, ax = plt.subplots(figsize=(11,6)) ax.axis("off") rows = [[k,str(v)] for k,v in summary.items()] stab = ax.table( cellText=rows, colLabels=["Metric","Value"], bbox=[0,0.55,0.45,0.4] ) stab.auto_set_font_size(False) stab.set_fontsize(11) anom = details[details["AnomalyType"]!="none"] if len(anom)>12: anom = anom.head(12) if not anom.empty: atab = ax.table( cellText=anom.astype(str).values, colLabels=list(anom.columns), bbox=[0,0.05,1,0.45] ) atab.auto_set_font_size(False) atab.set_fontsize(9) fig.tight_layout() fig.savefig(fname,dpi=150) plt.close(fig) print(f"Saved {fname}") # --------------------------------------------------------- # Main # --------------------------------------------------------- def main(): df = pd.read_csv(USAGE_FILE) df = df[df["ElapsedSeconds"].notnull()].reset_index(drop=True) if "MaxRSS_MiB" not in df.columns: print("ERROR: worker_usage.csv missing MaxRSS_MiB.") return df["JobID"] = df["JobID"].apply(numeric_jobid) df = df[df["JobID"].notnull()].reset_index(drop=True) logs = parse_worker_logs() merged = pd.merge(df, logs, on="JobID", how="inner") if merged.empty: print("ERROR: No matching logs.") return merged = merged.sort_values("StartTime").reset_index(drop=True) times = merged["StartTime"].tolist() idxs = np.arange(len(merged)) elapsed = merged["ElapsedSeconds"].astype(float) mem_mib = merged["MaxRSS_MiB"].astype(float) # Detect anomalies mtime, ztime = time_anomalies(elapsed) mmem, zmem = mem_anomalies(mem_mib) both = mtime & mmem tonly = mtime & ~mmem monly = mmem & ~mtime # ----------------------------------------------------- # 1. progress_with_anomalies.png # ----------------------------------------------------- fig, ax1 = plt.subplots(figsize=(12,6)) # Checkpoint curve with dots ax1.plot(times, idxs, 'o-', color="gray", label="Checkpoint Progress") ax1.set_xlabel("Wall Time") ax1.set_ylabel("Checkpoint Number",color="gray") # Anomaly markers for i,t in enumerate(times): if both[i]: ax1.plot(t, idxs[i], 'o', color="purple", markersize=10, label="Both anomalies" if i==0 else "") elif tonly[i]: ax1.plot(t, idxs[i], 'o', color="orange", markersize=9, label="Time anomaly" if i==0 else "") elif monly[i]: ax1.plot(t, idxs[i], 'o', color="blue", markersize=9, label="Memory anomaly" if i==0 else "") # MaxRSS dual-axis with dots ax2 = ax1.twinx() ax2.plot(times, mem_mib.values, 'bo-', markersize=5, alpha=0.70, label="MaxRSS (MiB)") ax2.set_ylabel("MaxRSS (MiB)", color="blue") # Highlight memory anomalies on memory curve for i,t in enumerate(times): if monly[i]: ax2.plot(t, mem_mib.values[i], 'bo', markersize=8) elif both[i]: ax2.plot(t, mem_mib.values[i], 'mo', markersize=10) h1,l1 = ax1.get_legend_handles_labels() h2,l2 = ax2.get_legend_handles_labels() ax1.legend(h1+h2, l1+l2, loc="upper left") fig.tight_layout() fig.savefig("progress_with_anomalies.png") plt.close(fig) print("Saved progress_with_anomalies.png") # ----------------------------------------------------- # 2. resource_trend.png (Elapsed vs MaxRSS dual-axis) # ----------------------------------------------------- fig2, axL = plt.subplots(figsize=(12,6)) axL.plot(elapsed.values, 'ko-', markersize=5, label="Elapsed (sec)") axL.set_xlabel("Worker Job Index (time-ordered)") axL.set_ylabel("Elapsed Seconds",color="black") axR = axL.twinx() axR.plot(mem_mib.values, 'bo-', markersize=5, alpha=0.7, label="MaxRSS (MiB)") axR.set_ylabel("MaxRSS (MiB)",color="blue") hL,lL = axL.get_legend_handles_labels() hR,lR = axR.get_legend_handles_labels() axL.legend(hL+hR, lL+lR, loc="upper left") fig2.tight_layout() fig2.savefig("resource_trend.png") plt.close(fig2) print("Saved resource_trend.png") # ----------------------------------------------------- # 3. anomaly_details + summary table # ----------------------------------------------------- details = pd.DataFrame({ "JobID": merged["JobID"], "StartTime": merged["StartTime"], "ElapsedSeconds": elapsed, "MaxRSS_MiB": mem_mib, "TimeZ": ztime, "MemModZ": zmem }) atype = np.where(both, "both", np.where(tonly, "time", np.where(monly,"memory","none"))) details["AnomalyType"] = atype summary = { "Total jobs": len(details), "Anomalies": int((atype!="none").sum()), "Time only": int((atype=="time").sum()), "Memory only": int((atype=="memory").sum()), "Both": int((atype=="both").sum()) } details.to_csv("anomaly_details.csv", index=False) print("Saved anomaly_details.csv") save_summary(summary, details) if __name__ == "__main__": main() }}} == Working Examples == === Make copies of the BASH scripts checkpoint_submitter.sh and checkpoint_runner.sh All of the examples linked further below make use of the following two (2) BASH job scripts mentioned above. * [wiki://Workshops/JobCheckpointing/Examples#checkpoint_submitter.sh checkpoint_submitter.sh] * [wiki://Workshops/JobCheckpointing/Examples#checkpoint_runner.sh checkpoint_runner.sh] === Make copies of the Python scripts aggregate_usage.py and plot_progress.py Also, the '''checkpoint_submitter.sh''' script calls out the following two (2) Python scripts for post-workflow analysis of execution time and memory usage by the checkpoint runner jobs. * [wiki://Workshops/JobCheckpointing/Examples#aggregate_usage.py aggregate_usage.py] * [wiki://Workshops/JobCheckpointing/Examples#plot_progress.py plot_progress.py] === BASH Checkpointing Example === See [wiki:Workshops/JobCheckpointing/Examples/BASH BASH Checkpointing Example] === Python Checkpointing Example === See [wiki:Workshops/JobCheckpointing/Examples/Python Python Checkpointing Example] === R Checkpointing Example === See [wiki:Workshops/JobCheckpointing/Examples/R R Checkpointing Example]