Разбивка по контейнерам (CPU/память/сеть) — то, что не смог дать
cAdvisor из-за несовместимости с containerd-снапшоттером сервера 1gb
(убран в 7f5c888). Свой минимальный сервис container-exporter (Python/
aiohttp/prometheus_client) опрашивает Docker Engine API параллельно в
фоновой задаче, не завязываясь на scrape-интервал Prometheus.
Доступ к докер-сокету изолирован через docker-socket-proxy: экспортеру
разрешены только GET /containers/json и /containers/*/stats, любые
изменяющие запросы блокируются на уровне прокси (POST=0) — полная
компрометация экспортера не даёт управлять Docker. Ни один из двух
сервисов не публикует портов наружу.
В host.json возвращены панели «Топ контейнеров по CPU/памяти» на новых
метриках vidconf_container_* (в 0.0.9 их убрали вместе с cAdvisor).
235 lines
12 KiB
Python
235 lines
12 KiB
Python
"""Экспортер метрик 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)
|