"""Экспортер метрик Prometheus по контейнерам поверх Docker Engine API. Отдельный минимальный сервис (не часть backend — backend не должен иметь доступ к докер-сокету, см. `docs/deploy/monitoring.md`). Опрашивает Docker Engine API ЧЕРЕЗ `docker-socket-proxy` (`deploy/docker-compose.yml`) — самому экспортеру `/var/run/docker.sock` не примонтирован, только доступ по внутренней сети compose к прокси, который разрешает исключительно `GET /containers/json` и `GET /containers/{id}/stats` (только чтение, без прав управлять Docker — см. обоснование в `docs/deploy/monitoring.md`). Стек — `aiohttp` (не `httpx`/FastAPI, как у backend): это отдельный процесс вне зависимостей backend, а `aiohttp` даёт async HTTP-клиент и сервер в одной небольшой библиотеке, без лишних зависимостей ради пяти строк веб-сервера. Метрики (КОНТРАКТ с `deploy/monitoring/prometheus.yml` и `deploy/monitoring/grafana/dashboards/host.json` — при переименовании поправить оба файла), лейбл везде — `name` (имя контейнера): - `vidconf_container_cpu_percent` — доля CPU, использованная контейнером за последний интервал опроса, тот же расчёт, что у `docker stats` (0-100 на одно ядро, суммарно может быть больше 100 при использовании нескольких ядер). - `vidconf_container_memory_usage_bytes` / `_limit_bytes` — из `stats.memory_stats`, `usage` очищен от переиспользуемого page cache (`inactive_file`) тем же приёмом, что cAdvisor использовал для `container_memory_working_set_bytes` — именно это число cgroup учитывает при OOM. Без явного лимита памяти у сервиса в docker-compose `limit` равен общему объёму RAM хоста — это не баг. - `vidconf_container_network_rx_bytes` / `_tx_bytes` — суммарно по всем сетевым интерфейсам контейнера, счётчик со старта контейнера (обнуляется при пересоздании) — для скорости передачи взять `rate()`. Опрос контейнеров идёт ПАРАЛЛЕЛЬНО (`asyncio.gather`) в фоновой задаче со своим интервалом (`POLL_INTERVAL_S`), а не синхронно на каждый scrape Prometheus: `GET /containers/{id}/stats?stream=false` не мгновенен (Docker считает CPU-дельту по двум внутренним замерам с разницей ~1с — проверено на реальном стенде), и линейный опрос дюжины контейнеров легко вышел бы за scrape-интервал Prometheus (15с). Параллельный опрос укладывается в ~1-2с независимо от числа контейнеров (проверено на 3 контейнерах, 4 запроса на каждый параллельно — 2с суммарно). Фоновая задача пишет значения прямо в метрики `prometheus_client` (они и есть кэш), обработчик `/metrics` не трогает Docker API вообще — стабильная латентность ответа независимо от состояния Docker API. """ import asyncio import logging import os import time from typing import Any import aiohttp from aiohttp import web from prometheus_client import CONTENT_TYPE_LATEST, Gauge, generate_latest logging.basicConfig(level=logging.INFO) logger = logging.getLogger("container_exporter") DOCKER_API_BASE = os.environ.get("DOCKER_API_BASE", "http://docker-socket-proxy:2375") POLL_INTERVAL_S = float(os.environ.get("POLL_INTERVAL_S", "15")) REQUEST_TIMEOUT_S = float(os.environ.get("REQUEST_TIMEOUT_S", "10")) LISTEN_PORT = int(os.environ.get("PORT", "9419")) CPU_PERCENT = Gauge( "vidconf_container_cpu_percent", "Использование CPU контейнером, % от одного ядра (как `docker stats`)", labelnames=("name",), ) MEMORY_USAGE_BYTES = Gauge( "vidconf_container_memory_usage_bytes", "Память контейнера без переиспользуемого page cache (working set)", labelnames=("name",), ) MEMORY_LIMIT_BYTES = Gauge( "vidconf_container_memory_limit_bytes", "Лимит памяти контейнера (= RAM хоста, если лимит не задан явно)", labelnames=("name",), ) NETWORK_RX_BYTES = Gauge( "vidconf_container_network_rx_bytes", "Байт получено контейнером по сети, суммарно по интерфейсам, с момента старта контейнера", labelnames=("name",), ) NETWORK_TX_BYTES = Gauge( "vidconf_container_network_tx_bytes", "Байт отправлено контейнером по сети, суммарно по интерфейсам, с момента старта контейнера", labelnames=("name",), ) class PollState: """Флаг успешности последнего опроса — источник `/health` для compose.""" last_poll_ok: bool = False def _container_name(container: dict[str, Any]) -> str: """Имя контейнера без ведущего `/` (Docker API отдаёт список имён с ним).""" names = container.get("Names") or [] if names: return str(names[0]).lstrip("/") return str(container["Id"])[:12] def _cpu_percent(stats: dict[str, Any]) -> float | None: """Доля CPU за интервал между `precpu_stats` и `cpu_stats` — расчёт `docker stats` (CLI). Оба поля Docker кладёт в один ответ `stats?stream=false` сам (снимает дважды с разницей ~1с внутри демона) — считать дельту вручную между двумя нашими запросами не нужно. """ cpu = stats.get("cpu_stats") or {} precpu = stats.get("precpu_stats") or {} cpu_total: int | None = (cpu.get("cpu_usage") or {}).get("total_usage") precpu_total: int | None = (precpu.get("cpu_usage") or {}).get("total_usage") system: int | None = cpu.get("system_cpu_usage") presystem: int | None = precpu.get("system_cpu_usage") # Первый снимок контейнера (только что создан) — precpu_stats пуст. if cpu_total is None or precpu_total is None or system is None or presystem is None: return None cpu_delta = cpu_total - precpu_total system_delta = system - presystem if system_delta <= 0 or cpu_delta < 0: return None percpu_usage = (cpu.get("cpu_usage") or {}).get("percpu_usage") or [] online_cpus = cpu.get("online_cpus") or len(percpu_usage) or 1 return float((cpu_delta / system_delta) * online_cpus * 100.0) def _apply_stats(name: str, stats: dict[str, Any]) -> None: cpu_percent = _cpu_percent(stats) if cpu_percent is not None: CPU_PERCENT.labels(name=name).set(cpu_percent) memory = stats.get("memory_stats") or {} usage = memory.get("usage") limit = memory.get("limit") cache = (memory.get("stats") or {}).get("inactive_file", 0) if usage is not None: MEMORY_USAGE_BYTES.labels(name=name).set(max(usage - cache, 0)) if limit is not None: MEMORY_LIMIT_BYTES.labels(name=name).set(limit) networks: dict[str, Any] = stats.get("networks") or {} if networks: rx = sum(net.get("rx_bytes", 0) for net in networks.values()) tx = sum(net.get("tx_bytes", 0) for net in networks.values()) NETWORK_RX_BYTES.labels(name=name).set(rx) NETWORK_TX_BYTES.labels(name=name).set(tx) async def _fetch_and_apply(session: aiohttp.ClientSession, container: dict[str, Any]) -> None: name = _container_name(container) container_id = container["Id"] timeout = aiohttp.ClientTimeout(total=REQUEST_TIMEOUT_S) try: async with session.get( f"{DOCKER_API_BASE}/containers/{container_id}/stats", params={"stream": "false"}, timeout=timeout, ) as response: response.raise_for_status() stats = await response.json(content_type=None) except Exception: logger.warning("не удалось получить stats контейнера %s", name, exc_info=True) return _apply_stats(name, stats) async def _poll_once(session: aiohttp.ClientSession) -> None: timeout = aiohttp.ClientTimeout(total=REQUEST_TIMEOUT_S) async with session.get(f"{DOCKER_API_BASE}/containers/json", timeout=timeout) as response: response.raise_for_status() containers = await response.json(content_type=None) # Параллельно — см. докстринг модуля про латентность stats-запроса. await asyncio.gather(*(_fetch_and_apply(session, c) for c in containers)) async def _poll_loop(state: PollState) -> None: async with aiohttp.ClientSession() as session: while True: start = time.monotonic() try: await _poll_once(session) state.last_poll_ok = True except Exception: logger.exception("опрос Docker API завершился ошибкой") state.last_poll_ok = False elapsed = time.monotonic() - start await asyncio.sleep(max(POLL_INTERVAL_S - elapsed, 1.0)) # `aiohttp.web.Response(content_type=...)` не принимает `charset` внутри # строки типа (падает `ValueError`), а `CONTENT_TYPE_LATEST` из # `prometheus_client` — "text/plain; version=0.0.4; charset=utf-8" целиком; # разбираем один раз на media type (без `charset=...`) и сам charset. _METRICS_MEDIA_TYPE = "; ".join( part for part in CONTENT_TYPE_LATEST.split("; ") if not part.startswith("charset=") ) async def metrics_handler(_request: web.Request) -> web.Response: return web.Response(body=generate_latest(), content_type=_METRICS_MEDIA_TYPE, charset="utf-8") def _make_health_handler(state: PollState) -> Any: async def health_handler(_request: web.Request) -> web.Response: if state.last_poll_ok: return web.Response(text="ok") return web.Response(text="опрос Docker API ещё не выполнялся успешно", status=503) return health_handler def create_app() -> web.Application: app = web.Application() state = PollState() app.router.add_get("/metrics", metrics_handler) app.router.add_get("/health", _make_health_handler(state)) async def on_startup(_app: web.Application) -> None: app["poll_task"] = asyncio.create_task(_poll_loop(state)) async def on_cleanup(_app: web.Application) -> None: app["poll_task"].cancel() app.on_startup.append(on_startup) app.on_cleanup.append(on_cleanup) return app if __name__ == "__main__": web.run_app(create_app(), host="0.0.0.0", port=LISTEN_PORT)