Coverage for python/lsst/analysis/ap/taskRuntimes.py: 8%

48 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-23 11:04 +0000

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/>. 

21 

22"""Collect per-quantum runtimes from the ``*_metadata`` datasets of a 

23butler run collection. 

24 

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""" 

31 

32from __future__ import annotations 

33 

34__all__ = ["collect_task_runtimes"] 

35 

36import pandas as pd 

37 

38from lsst.pipe.base.resource_usage import QuantumResourceUsage 

39 

40 

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. 

44 

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. 

52 

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. 

71 

72 Returns 

73 ------- 

74 df : `pandas.DataFrame` 

75 One row per task, with columns: 

76 

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}) 

105 

106 if not rows: 

107 return (pd.DataFrame(), None) if plot else pd.DataFrame() 

108 

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)] 

113 

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)) 

127 

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}) 

136 

137 if not plot: 

138 return summary 

139 

140 import matplotlib.pyplot as plt 

141 

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