Coverage for python/lsst/analysis/ap/taskRuntimes.py: 8%
48 statements
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-10 06:35 -0300
« prev ^ index » next coverage.py v7.16.0, created at 2026-09-10 06:35 -0300
1# This file is part of analysis_ap.
2#
3# Developed for the LSST Data Management System.
4# This product includes software developed by the LSST Project
5# (https://www.lsst.org).
6# See the COPYRIGHT file at the top-level directory of this distribution
7# for details of code ownership.
8#
9# This program is free software: you can redistribute it and/or modify
10# it under the terms of the GNU General Public License as published by
11# the Free Software Foundation, either version 3 of the License, or
12# (at your option) any later version.
13#
14# This program is distributed in the hope that it will be useful,
15# but WITHOUT ANY WARRANTY; without even the implied warranty of
16# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
17# GNU General Public License for more details.
18#
19# You should have received a copy of the GNU General Public License
20# along with this program. If not, see <https://www.gnu.org/licenses/>.
22"""Collect per-quantum runtimes from the ``*_metadata`` datasets of a
23butler run collection.
25The single public entry point, `collect_task_runtimes`, walks every
26``<task>_metadata`` dataset, pulls the timing fields via
27`lsst.pipe.base.resource_usage.QuantumResourceUsage.from_task_metadata`,
28applies a per-task threshold (any-quantum-over-threshold keeps the
29whole task), and returns a tidy DataFrame, optionally with a box plot.
30"""
32from __future__ import annotations
34__all__ = ["collect_task_runtimes"]
36import pandas as pd
38from lsst.pipe.base.resource_usage import QuantumResourceUsage
41def collect_task_runtimes(butler, collections, threshold=1.0, *,
42 plot=False, ax=None):
43 """Per-task runtime and memory summary for a butler run collection.
45 Each ``<task>_metadata`` dataset under ``collections`` is loaded and
46 its timing fields extracted with
47 `~lsst.pipe.base.resource_usage.QuantumResourceUsage.from_task_metadata`.
48 Tasks whose every quantum runs faster than ``threshold`` seconds are
49 dropped; tasks with at least one quantum at or above ``threshold``
50 contribute all their quanta to the per-task summary statistics, so the
51 summary reflects cross-quantum variability rather than a single outlier.
53 Parameters
54 ----------
55 butler : `lsst.daf.butler.Butler`
56 Butler used to query the registry and load metadata datasets.
57 collections : `str` or iterable of `str`
58 Collections to query, typically a single run collection name.
59 threshold : `float`, optional
60 Minimum task duration (seconds) for inclusion. A task is kept iff
61 at least one of its quanta has ``total_time`` >= ``threshold``.
62 Default ``1.0``.
63 plot : `bool`, optional
64 If True, render a horizontal box plot of per-quantum
65 ``total_time`` per surviving task (tasks ordered by max
66 ``total_time`` descending) and return the ``(df, fig)`` pair
67 instead of just ``df``.
68 ax : `matplotlib.axes.Axes` or None
69 Axes to plot onto. Only used when ``plot=True``. If None, a new
70 figure and axes are created.
72 Returns
73 -------
74 df : `pandas.DataFrame`
75 One row per task, with columns:
77 - ``task``: pipeline task label
78 - ``n_quanta``: number of surviving quanta contributing to the row
79 - ``total_time_mean``, ``total_time_min``, ``total_time_max``,
80 ``total_time_std``: seconds
81 - ``memory_mean_<UNIT>``, ``memory_min_<UNIT>``,
82 ``memory_max_<UNIT>``, ``memory_std_<UNIT>``: where ``UNIT`` is
83 ``GB`` if any task's peak memory crosses 1 GB and ``MB``
84 otherwise. The unit is chosen once across the whole table so
85 the columns remain numerically comparable.
86 fig : `matplotlib.figure.Figure`
87 Only returned when ``plot=True``.
88 """
89 rows = []
90 for dataset_type in butler.registry.queryDatasetTypes("*_metadata"):
91 task = dataset_type.name[:-len("_metadata")]
92 for ref in butler.registry.queryDatasets(dataset_type, collections=collections):
93 metadata = butler.get(ref)
94 try:
95 usage = QuantumResourceUsage.from_task_metadata(metadata)
96 except KeyError:
97 # Quantum block exists but is missing one of the expected
98 # fields (e.g. an aborted run); skip rather than fail.
99 continue
100 if usage is None:
101 continue
102 rows.append({"task": task,
103 "total_time": usage.total_time,
104 "memory": usage.memory})
106 if not rows:
107 return (pd.DataFrame(), None) if plot else pd.DataFrame()
109 per_quantum = pd.DataFrame(rows)
110 per_task_max = per_quantum.groupby("task")["total_time"].max()
111 keep_tasks = per_task_max[per_task_max >= threshold].index
112 per_quantum = per_quantum[per_quantum["task"].isin(keep_tasks)]
114 summary = (per_quantum.groupby("task")
115 .agg(n_quanta=("total_time", "count"),
116 total_time_mean=("total_time", "mean"),
117 total_time_min=("total_time", "min"),
118 total_time_max=("total_time", "max"),
119 total_time_std=("total_time", "std"),
120 memory_mean=("memory", "mean"),
121 memory_min=("memory", "min"),
122 memory_max=("memory", "max"),
123 memory_std=("memory", "std"))
124 .reset_index()
125 .sort_values("total_time_max", ascending=False)
126 .reset_index(drop=True))
128 mem_cols = ["memory_mean", "memory_min", "memory_max", "memory_std"]
129 # Use GB if any task's peak memory crosses 1 GB, otherwise MB.
130 if summary["memory_max"].max() >= 1024**3:
131 mem_unit, mem_divisor = "GB", 1024**3
132 else:
133 mem_unit, mem_divisor = "MB", 1024**2
134 summary[mem_cols] = summary[mem_cols] / mem_divisor
135 summary = summary.rename(columns={c: f"{c}_{mem_unit}" for c in mem_cols})
137 if not plot:
138 return summary
140 import matplotlib.pyplot as plt
142 # Bottom-up order so the slowest task sits at the top of the horizontal
143 # plot.
144 order = per_task_max.loc[keep_tasks].sort_values(ascending=True).index.tolist()
145 data = [per_quantum.loc[per_quantum["task"] == t, "total_time"].to_numpy() for t in order]
146 if ax is None:
147 fig, ax = plt.subplots(figsize=(8, max(3, 0.3 * len(order))))
148 else:
149 fig = ax.figure
150 ax.boxplot(data, vert=False, tick_labels=order, showfliers=True)
151 ax.set_xlabel("total_time (s)")
152 if per_quantum["total_time"].max() > 100 * threshold:
153 ax.set_xscale("log")
154 ax.axvline(threshold, color="grey", linestyle=":", linewidth=1,
155 label=f"threshold = {threshold:g} s")
156 ax.set_title("Per-quantum total_time")
157 ax.grid(axis="x", linestyle=":", alpha=0.5)
158 ax.legend(loc="lower right", fontsize="small")
159 fig.tight_layout()
160 return summary, fig