|
| 1 | +"""ProgressBar shows the how many dask tasks finished/remains using tqdm.""" |
| 2 | + |
| 3 | +from typing import Any, Optional, Dict, Tuple, Union |
| 4 | +from time import time |
| 5 | + |
| 6 | +from dask.callbacks import Callback |
| 7 | + |
| 8 | +from .utils import is_notebook |
| 9 | + |
| 10 | +if is_notebook(): |
| 11 | + from tqdm.notebook import tqdm |
| 12 | +else: |
| 13 | + from tqdm import tqdm |
| 14 | + |
| 15 | +# pylint: disable=method-hidden,too-many-instance-attributes |
| 16 | +class ProgressBar(Callback): # type: ignore |
| 17 | + """A progress bar for DataPrep.EDA. |
| 18 | +
|
| 19 | + Parameters |
| 20 | + ---------- |
| 21 | + minimum : int, optional |
| 22 | + Minimum time threshold in seconds before displaying a progress bar. |
| 23 | + Default is 0 (always display) |
| 24 | + _min_tasks : int, optional |
| 25 | + Minimum graph size to show a progress bar, default is 5 |
| 26 | + width : int, optional |
| 27 | + Width of the bar. None means auto width. |
| 28 | + interval : float, optional |
| 29 | + Update resolution in seconds, default is 0.1 seconds |
| 30 | + """ |
| 31 | + |
| 32 | + _minimum: float = 0 |
| 33 | + _min_tasks: int = 5 |
| 34 | + _width: Optional[int] = None |
| 35 | + _interval: float = 0.1 |
| 36 | + _last_duration: float = 0 |
| 37 | + _pbar: Optional[tqdm] = None |
| 38 | + _state: Optional[Dict[str, Any]] = None |
| 39 | + _started: Optional[float] = None |
| 40 | + _last_task: Optional[str] = None # in case we initialize the pbar in _finish |
| 41 | + |
| 42 | + def __init__( |
| 43 | + self, |
| 44 | + minimum: float = 0, |
| 45 | + min_tasks: int = 5, |
| 46 | + width: Optional[int] = None, |
| 47 | + interval: float = 0.1, |
| 48 | + ) -> None: |
| 49 | + super().__init__() |
| 50 | + self._minimum = minimum |
| 51 | + self._min_tasks = min_tasks |
| 52 | + self._width = width |
| 53 | + self._interval = interval |
| 54 | + |
| 55 | + def _start(self, _dsk: Any) -> None: |
| 56 | + """A hook to start this callback.""" |
| 57 | + |
| 58 | + def _start_state(self, _dsk: Any, state: Dict[str, Any]) -> None: |
| 59 | + """A hook called before every task gets executed.""" |
| 60 | + self._started = time() |
| 61 | + self._state = state |
| 62 | + _, ntasks = self._count_tasks() |
| 63 | + |
| 64 | + if ntasks > self._min_tasks: |
| 65 | + self._init_bar() |
| 66 | + |
| 67 | + def _pretask( |
| 68 | + self, key: Union[str, Tuple[str, ...]], _dsk: Any, _state: Dict[str, Any] |
| 69 | + ) -> None: |
| 70 | + """A hook called before one task gets executed.""" |
| 71 | + if self._started is None: |
| 72 | + raise ValueError("ProgressBar not started properly") |
| 73 | + |
| 74 | + if self._pbar is None and time() - self._started > self._minimum: |
| 75 | + self._init_bar() |
| 76 | + |
| 77 | + if isinstance(key, tuple): |
| 78 | + key = key[0] |
| 79 | + |
| 80 | + if self._pbar is not None: |
| 81 | + self._pbar.set_description(f"Computing {key}") |
| 82 | + else: |
| 83 | + self._last_task = key |
| 84 | + |
| 85 | + def _posttask( |
| 86 | + self, |
| 87 | + _key: str, |
| 88 | + _result: Any, |
| 89 | + _dsk: Any, |
| 90 | + _state: Dict[str, Any], |
| 91 | + _worker_id: Any, |
| 92 | + ) -> None: |
| 93 | + """A hook called after one task gets executed.""" |
| 94 | + |
| 95 | + if self._pbar is not None: |
| 96 | + self._update_bar() |
| 97 | + |
| 98 | + def _finish(self, _dsk: Any, _state: Dict[str, Any], _errored: bool) -> None: |
| 99 | + """A hook called after all tasks get executed.""" |
| 100 | + if self._started is None: |
| 101 | + raise ValueError("ProgressBar not started properly") |
| 102 | + |
| 103 | + if self._pbar is None and time() - self._started > self._minimum: |
| 104 | + self._init_bar() |
| 105 | + |
| 106 | + if self._pbar is not None: |
| 107 | + self._update_bar() |
| 108 | + self._pbar.close() |
| 109 | + |
| 110 | + self._state = None |
| 111 | + self._started = None |
| 112 | + self._pbar = None |
| 113 | + |
| 114 | + def _update_bar(self) -> None: |
| 115 | + if self._pbar is None: |
| 116 | + return |
| 117 | + ndone, _ = self._count_tasks() |
| 118 | + |
| 119 | + self._pbar.update(max(0, ndone - self._pbar.n)) |
| 120 | + |
| 121 | + def _init_bar(self) -> None: |
| 122 | + if self._pbar is not None: |
| 123 | + raise ValueError("ProgressBar already initialized.") |
| 124 | + ndone, ntasks = self._count_tasks() |
| 125 | + |
| 126 | + if self._last_task is not None: |
| 127 | + desc = f"Computing {self._last_task}" |
| 128 | + else: |
| 129 | + desc = "" |
| 130 | + |
| 131 | + if self._width is None: |
| 132 | + self._pbar = tqdm( |
| 133 | + total=ntasks, |
| 134 | + dynamic_ncols=True, |
| 135 | + mininterval=self._interval, |
| 136 | + initial=ndone, |
| 137 | + desc=desc, |
| 138 | + ) |
| 139 | + else: |
| 140 | + self._pbar = tqdm( |
| 141 | + total=ntasks, |
| 142 | + ncols=self._width, |
| 143 | + mininterval=self._interval, |
| 144 | + initial=ndone, |
| 145 | + desc=desc, |
| 146 | + ) |
| 147 | + |
| 148 | + self._pbar.start_t = self._started |
| 149 | + self._pbar.refresh() |
| 150 | + |
| 151 | + def _count_tasks(self) -> Tuple[int, int]: |
| 152 | + if self._state is None: |
| 153 | + raise ValueError("ProgressBar not started properly") |
| 154 | + |
| 155 | + state = self._state |
| 156 | + ndone = len(state["finished"]) |
| 157 | + ntasks = sum(len(state[k]) for k in ["ready", "waiting", "running"]) + ndone |
| 158 | + |
| 159 | + return ndone, ntasks |
0 commit comments