|
2 | 2 | from contextlib import AbstractContextManager |
3 | 3 | from datetime import datetime, timezone |
4 | 4 | from importlib.metadata import version |
5 | | -from typing import Any, TypeVar |
| 5 | +from typing import Any, Generator, TypeVar |
| 6 | +import psutil |
6 | 7 |
|
7 | 8 | from packaging.version import Version, parse |
8 | 9 |
|
|
17 | 18 |
|
18 | 19 | from opentelemetry import context as context_api |
19 | 20 | from opentelemetry import trace |
20 | | -from opentelemetry.metrics import Meter, MeterProvider, get_meter |
| 21 | +from opentelemetry.metrics import Meter, MeterProvider, Observation, get_meter |
21 | 22 | from opentelemetry.propagate import extract, inject |
22 | 23 | from opentelemetry.semconv.trace import SpanAttributes |
23 | 24 | from opentelemetry.trace import Span, Tracer, TracerProvider |
@@ -205,6 +206,30 @@ def __init__( |
205 | 206 | unit="s", |
206 | 207 | description="Time the tasks waited before executing", |
207 | 208 | ) |
| 209 | + # current metrics to watch for in workers: CPU and memory utilization |
| 210 | + self._process = psutil.Process() |
| 211 | + # 6- CPU utilization |
| 212 | + self.worker_cpu_utilization = self._meter.create_observable_gauge( |
| 213 | + "worker_cpu_utilization", |
| 214 | + callbacks=[self._observe_cpu], |
| 215 | + unit="%", |
| 216 | + description="Worker CPU utilization percentage. Only for worker processes", |
| 217 | + ) |
| 218 | + # 7- Memory utilization |
| 219 | + self.worker_memory_utilization = self._meter.create_observable_gauge( |
| 220 | + "worker_memory_utilization", |
| 221 | + callbacks=[self._observe_memory], |
| 222 | + unit="By", |
| 223 | + description="Worker memory utilization in bytes. Only for worker processes", |
| 224 | + ) |
| 225 | + |
| 226 | + def _observe_memory(self, _options: Any) -> Generator[Observation, None, None]: |
| 227 | + if self.broker and self.broker.is_worker_process: |
| 228 | + yield Observation(self._process.memory_info().rss) |
| 229 | + |
| 230 | + def _observe_cpu(self, _options: Any) -> Generator[Observation, None, None]: |
| 231 | + if self.broker and self.broker.is_worker_process: |
| 232 | + yield Observation(self._process.cpu_percent()) |
208 | 233 |
|
209 | 234 | def pre_send(self, message: TaskiqMessage) -> TaskiqMessage: |
210 | 235 | """ |
|
0 commit comments