FastAPI для MLOps: Проектирование масштабируемых систем инференса

Интенсивный курс по созданию высокопроизводительных API для машинного обучения. Вы научитесь оптимизировать работу с тяжелыми моделями, внедрять очереди задач и настраивать профессиональный мониторинг для production-сред.

Специфика FastAPI в MLOps и жизненный цикл приложения lifespan

Специфика FastAPI в MLOps и жизненный цикл приложения lifespan

Представьте, что вы разворачиваете нейросеть весом 5 ГБ. Если ваш API будет загружать её с диска в оперативную память при каждом входящем HTTP-запросе, время ответа (latency) составит десятки секунд, а на третьем параллельном запросе сервер упадет с ошибкой Out Of Memory. Пропускная способность системы описывается простой зависимостью: Throughput=ConcurrencyLatencyThroughput = \frac{Concurrency}{Latency}. Если LatencyLatency огромно из-за постоянной загрузки весов, ваш ThroughputThroughput стремится к нулю.

Главный парадокс интеграции машинного обучения в веб заключается в конфликте парадигм. Традиционный веб тяготеет к архитектуре Stateless (без сохранения состояния): запрос пришел, сервер сходил в базу данных, отдал ответ и «забыл» о клиенте. Инференс ML-моделей — это жесткий Stateful: модель должна быть загружена в память (RAM или VRAM видеокарты) до того, как придут пользователи, и постоянно находиться там.

FastAPI стал стандартом де-факто в MLOps именно потому, что он элегантно решает эту проблему, предоставляя разработчику полный контроль над состоянием приложения с помощью механизма lifespan.

Эволюция управления состоянием: почему lifespan?

В ранних версиях FastAPI (и Starlette, на котором он основан) для загрузки моделей использовались декораторы событий: @app.on_event("startup") и @app.on_event("shutdown"). Выглядело это так: мы объявляли глобальную переменную, в startup загружали в нее модель, а в shutdown очищали.

На технических собеседованиях часто спрашивают, почему этот подход признан устаревшим (deprecated). У него три фатальных недостатка:

  1. Глобальные переменные: Они усложняют тестирование. Если вы хотите запустить тесты параллельно с разными мок-моделями, глобальное состояние приведет к гонке данных.
  2. Отсутствие гарантий: Если приложение падает во время работы, событие shutdown может не вызваться, оставляя «висящие» процессы и занятую память на GPU.
  3. Разрыв контекста: Логика инициализации и очистки разнесена по разным функциям, хотя семантически это один процесс.

Современный стандарт — использование асинхронного контекстного менеджера lifespan.

Анатомия lifespan в FastAPI

Контекстный менеджер lifespan оборачивает всю жизнь вашего приложения. Всё, что написано до ключевого слова yield, выполняется при старте сервера. Всё, что после — при его остановке.

Давайте посмотрим на эталонный пример подготовки ML-сервиса:

from contextlib import asynccontextmanager
from fastapi import FastAPI, Request
import torch

# Имитация тяжелой ML-модели
class MLModel:
    def load(self):
        print("Загрузка весов (5 ГБ) в VRAM...")
        self.device = "cuda" if torch.cuda.is_available() else "cpu"

    def predict(self, data: str) -> str:
        return f"Prediction for {data} on {self.device}"

    def release(self):
        print("Очистка VRAM...")
        torch.cuda.empty_cache()

@asynccontextmanager
async def lifespan(app: FastAPI):
    # 1. ФАЗА STARTUP: Инициализация ресурсов
    model = MLModel()
    model.load()

    # Сохраняем модель в глобальное состояние приложения
    yield {"ml_model": model}

    # 2. ФАЗА SHUTDOWN: Очистка ресурсов
    model.release()
    print("Сервер успешно остановлен, ресурсы освобождены.")

# Передаем lifespan при создании приложения
app = FastAPI(lifespan=lifespan)

Как передать модель в endpoint?

Обратите внимание на строку yield {"ml_model": model}. Словарь, который мы возвращаем (yield), автоматически помещается в request.state. Это изолированное хранилище состояния для конкретного экземпляра приложения.

Теперь в любом эндпоинте мы можем получить доступ к нашей тяжелой модели без использования глобальных переменных:

@app.post("/predict")
async def predict_endpoint(request: Request, payload: str):
    # Извлекаем модель из состояния приложения
    model = request.state.ml_model

    # Выполняем инференс
    result = model.predict(payload)
    return {"result": result}

Примечание: в реальных production-системах прямое обращение к request.state часто оборачивают в систему Dependency Injection (Depends), чтобы улучшить типизацию и автодополнение кода. Эту технику мы подробно разберем в одной из следующих глав.

Graceful Shutdown: почему это критично для MLOps

В мире микросервисов контейнеры постоянно перезапускаются (например, Kubernetes масштабирует поды или выкатывает новую версию). Когда оркестратор решает убить контейнер, он посылает сигнал SIGTERM.

Если приложение не умеет корректно обрабатывать этот сигнал (делать graceful shutdown), оно завершается жестко. Для обычного веб-сервера это означает пару оборванных HTTP-соединений. Для ML-сервиса последствия хуже:

  • Утечки VRAM: Видеокарты (особенно при работе с CUDA) очень чувствительны к незавершенным процессам. Жесткое убийство процесса может оставить куски памяти в GPU помеченными как «занятые». В итоге новый под просто не сможет стартовать из-за ошибки CUDA out of memory.
  • Потеря батчей: Если модель обрабатывала пакет данных (batch) в момент остановки, данные будут потеряны.

Блок кода после yield в lifespan гарантирует, что при получении SIGTERM FastAPI перестанет принимать новые запросы, дождется завершения текущих (в рамках таймаута), выполнит вашу логику очистки (model.release()) и только потом умрет.

Сквозная логика: от старта до инференса

Мы решили первую фундаментальную проблему — загрузили тяжелую модель один раз при старте и обеспечили ей безопасное удаление. Однако, если мы запустим наш predict_endpoint под нагрузкой прямо сейчас, мы столкнемся с новой проблемой.

Инференс нейросетей — это синхронная, блокирующая процессор (CPU-bound) задача. Пока наша модель высчитывает тензоры для одного пользователя, весь асинхронный event loop FastAPI будет заблокирован, и другие пользователи не смогут даже получить ответ «сервер жив». Как подружить синхронную математику с асинхронным веб-сервером — это следующая ступень проектирования, к которой мы перейдем далее.

Асинхронность в FastAPI: event loop, async/await и блокирующий CPU-код

Асинхронность в FastAPI: event loop, async/await и блокирующий CPU-код

В прошлой главе мы успешно загрузили тяжелую модель в память с помощью lifespan и сохранили её в request.state. Приложение стартует идеально. Но вот мы выводим сервис в продакшен. Первый же пользователь отправляет изображение на классификацию, модель начинает вычисления, и внезапно... все остальные пользователи получают ошибки таймаута. Сервер жив, но не отвечает ни на один запрос, даже на простейший /health. Почему фреймворк, знаменитый своей скоростью, завис от одного запроса?

Ответ кроется в конфликте парадигм: асинхронный веб-сервер оптимизирован для ожидания, а ML-модели созданы для тяжелой математической работы. Понимание того, как FastAPI маршрутизирует задачи, — один из самых частых вопросов на технических собеседованиях для MLOps.

Event loop и парадокс MLOps

В основе FastAPI (и библиотеки asyncio) лежит event loop (цикл событий). Это однопоточный механизм, который управляет выполнением задач. Его главная суперсила — эффективная работа с I/O-bound операциями (ввод-вывод). Когда код делает запрос к базе данных или скачивает файл по сети, event loop не ждет ответа. Он переключается на обслуживание следующего HTTP-запроса.

Но инференс нейросетей — это CPU-bound задача. Перемножение матриц под капотом модели имеет вычислительную сложность O(n3)O(n^3), где nn — размерность матрицы. Это непрерывный поток математических операций, который полностью захватывает процессорное ядро.

Если запустить CPU-bound задачу напрямую внутри event loop, цикл блокируется. Он физически не может переключиться на другие задачи, потому что процессор занят вычислениями. В результате новые HTTP-соединения ставятся в очередь, пока не истечет их время ожидания.

Ловушка async def в FastAPI

FastAPI заставляет разработчика принимать архитектурное решение прямо в сигнатуре эндпоинта: использовать def или async def. В контексте ML-инференса это решение критично.

Посмотрим, как фреймворк обрабатывает оба варианта.

Сценарий 1: Синхронный эндпоинт (def)

Если вы объявляете обработчик пути (endpoint) как обычную def функцию, FastAPI выполняет её в отдельном пуле потоков (thread pool), чтобы не блокировать event loop.

Официальная документация FastAPI

Это контринтуитивно, но для тяжелых ML-моделей обычный def часто бывает безопаснее. FastAPI использует библиотеку Starlette, которая отправляет синхронные функции в ThreadPoolExecutor.

@app.post("/predict_sync")
def predict_sync(request: Request, data: dict):
    # Извлекаем модель, загруженную в lifespan
    model = request.state.ml_model

    # Эта тяжелая операция НЕ заблокирует event loop,
    # так как FastAPI запустит весь этот эндпоинт в отдельном потоке
    prediction = model.predict(data["tensor"])

    return {"result": prediction}

Сценарий 2: Асинхронный эндпоинт (async def)

Если вы пишете async def, FastAPI предполагает, что вы знаете, что делаете, и запускает код напрямую в event loop.

Если внутри такого эндпоинта вызвать синхронную блокирующую функцию model.predict(), весь сервер остановится до завершения предсказания.

Тип эндпоинта Где выполняется код Риск для ML-инференса
def Во внешнем пуле потоков (Thread pool) Низкий. Event loop свободен, но есть накладные расходы на потоки.
async def Напрямую в главном потоке (Event loop) Критический. Синхронный вызов модели парализует весь веб-сервер.

Как правильно совмещать async и ML-модели

Часто нам необходимо использовать именно async def. Например, перед инференсом нужно асинхронно сходить в базу данных за профилем пользователя или отправить метрику в Redis.

Как в таком случае выполнить синхронный model.predict(), не убивая сервер? Нам нужно вручную отправить тяжелую задачу в пул потоков, используя asyncio.to_thread() (доступно с Python 3.9).

import asyncio
from fastapi import FastAPI, Request

@app.post("/predict_async")
async def predict_async(request: Request, data: dict):
    model = request.state.ml_model

    # 1. Асинхронный I/O (не блокирует)
    user_profile = await db.get_user(data["user_id"])

    # 2. Выносим CPU-bound задачу в отдельный поток,
    # а event loop через await ждет результата, продолжая обрабатывать другие запросы
    prediction = await asyncio.to_thread(model.predict, data["tensor"])

    return {"result": prediction, "user": user_profile}

В этом примере asyncio.to_thread берет синхронную функцию model.predict, передает ей аргумент data["tensor"] и выполняет её в фоновом потоке, возвращая управление event loop.

GIL и инференс: почему потоки работают?

На собеседованиях часто задают провокационный вопрос: "В Python есть GIL (Global Interpreter Lock), который запрещает параллельное выполнение Python-кода в нескольких потоках. Какой смысл выносить ML-модель в Thread pool, если GIL всё равно всё заблокирует?"

Секрет кроется в том, как написаны современные ML-фреймворки. Библиотеки вроде PyTorch, TensorFlow и NumPy написаны на C/C++. Когда вы вызываете model.predict(), интерпретатор Python передает управление C-коду.

В этот момент ML-фреймворк отпускает GIL. Пока процессор перемножает тензоры на уровне C++, GIL свободен, и event loop в основном потоке Python может беспрепятственно обрабатывать новые входящие HTTP-запросы. Потоки в Python отлично подходят для ML-инференса именно благодаря этому механизму.

Границы применимости пула потоков

Использование asyncio.to_thread или синхронных def эндпоинтов решает проблему блокировки сервера. Но у этого подхода есть архитектурные пределы:

  1. Ограничение пула. По умолчанию пул потоков в asyncio ограничен (обычно min(32, os.cpu_count() + 4)). Если придет 100 одновременных запросов, 32 из них займут потоки, а остальные 68 будут ждать в очереди.
  2. Таймауты HTTP. Если инференс занимает 15 секунд, браузер или клиентский сервис может разорвать соединение по таймауту (обычно 30-60 секунд) еще до того, как поток освободится.
  3. Сериализация данных. Прежде чем передать данные в модель, их нужно провалидировать и превратить в тензоры, что также требует ресурсов.

Мы защитили веб-сервер от зависания, но для построения по-настоящему масштабируемой системы этого недостаточно. Перед тем как переходить к архитектуре очередей для долгих задач, нам предстоит решить проблему передачи тяжелых многомерных массивов через API. Об оптимизации сериализации и валидации тензоров мы поговорим на следующем шаге.

Оптимизация сериализации: Pydantic v2 и валидация тяжелых тензоров

Оптимизация сериализации: Pydantic v2 и валидация тяжелых тензоров

Модель безопасно загружена в память через lifespan, а тяжелые матричные умножения изолированы в пуле потоков. Архитектура кажется готовой к высоким нагрузкам. Но при запуске нагрузочного тестирования с реалистичными данными — например, батчами изображений или объемными эмбеддингами — API внезапно начинает «захлебываться», а утилизация CPU упирается в 100%. Профайлер показывает, что event loop заблокирован еще до того, как запрос дошел до логики эндпоинта. Причина кроется в фазе десериализации и валидации входящего JSON.

Анатомия проблемы: почему JSON не подходит для тензоров

Тензоры — это плотные массивы чисел с плавающей точкой. Стандартный инстинкт разработчика при проектировании API — описать их передачу через встроенные типы, например List[float]. Для небольших конфигурационных данных это работает отлично, но в MLOps масштабы иные.

Оценим размер текстового представления массива в формате JSON:

S=N×cS = N \times c

  • SS — итоговый размер JSON-строки в байтах.
  • NN — количество чисел в массиве.
  • cc — среднее количество байт на одно число в текстовом виде (включая разделитель-запятую).

Пример: Клиент отправляет вектор на 10610^6 элементов (эквивалент одноканального мегапиксельного изображения). Одно число с плавающей точкой в тексте выглядит как 0.123456, и занимает около 10 байт. Итоговый размер JSON-строки составит 1010 МБ.

Чтобы обработать такой запрос, парсеру необходимо прочитать 10 мегабайт текста символ за символом, найти все запятые, конвертировать каждую подстроку в число и создать миллион независимых Python-объектов float. Эта операция выполняется в главном потоке, полностью блокируя event loop.

Pydantic v2 и ограничения pydantic-core

Вторая версия Pydantic получила ядро pydantic-core, написанное на языке Rust. Это дало колоссальный прирост производительности (в 5–50 раз) при валидации стандартных структур данных по сравнению с первой версией.

Однако Rust-движок не отменяет законов алгоритмической сложности. Парсинг огромного JSON-массива остается поэлементной операцией. Даже если Rust делает это быстро, итогом все равно становится создание тысяч тяжеловесных Python-объектов в памяти. Нам нужно передавать данные так, чтобы парсер обрабатывал их единым блоком.

Обход JSON-парсера: Base64 и бинарные данные

Вместо текста тензоры эффективнее представлять в виде непрерывного блока сырых байт.

V=N×bV = N \times b

  • VV — объем бинарных данных в оперативной памяти (в байтах).
  • NN — количество элементов тензора.
  • bb — размер одного элемента в байтах (например, 4 байта для типа float32).

Для того же вектора на 10610^6 элементов объем составит всего 44 МБ.

Поскольку стандартный JSON не поддерживает бинарные данные, используется кодировка Base64. Она переводит байты в ASCII-строку, увеличивая итоговый объем примерно на 33%, но главное преимущество сохраняется: мы избавляемся от поэлементного парсинга. Строка декодируется целиком за одну быструю C-операцию, результат которой можно напрямую загрузить в непрерывный участок памяти (например, в массив NumPy).

Характеристика List[float] в JSON Base64 + np.ndarray
Формат передачи Текстовый массив чисел Единая текстовая строка
Парсинг Поэлементный (очень медленно) Блочный C-уровень (очень быстро)
Потребление памяти Высокое (миллионы Python-объектов) Низкое (один непрерывный блок памяти)

Валидация ML-данных: принцип Fail-fast

Быстрая десериализация — только половина задачи. Перед тем как передать данные в ML-модель, необходимо убедиться, что они имеют правильную форму (shape) и тип.

Здесь вступает в силу принцип Fail-fast (падай быстро). Если передать в PyTorch тензор неправильной размерности, ошибка возникнет глубоко в C++ ядре во время выполнения графа вычислений. Это сложно перехватить, а ресурсы пула потоков уже будут потрачены впустую. Pydantic позволяет отсечь невалидный запрос еще на этапе маршрутизации, вернув клиенту понятный HTTP-статус 422 Unprocessable Entity.

В Pydantic v2 для создания сложных пайплайнов валидации используется тип Annotated в связке с BeforeValidator и AfterValidator. Это позволяет перехватить сырые данные до того, как Pydantic попытается применить к ним свои стандартные правила, а затем наложить специфичные для предметной области ограничения.

Практическая реализация

Рассмотрим создание кастомного типа данных, который принимает Base64-строку, мгновенно конвертирует ее в NumPy-массив и проверяет размерность.

import base64
import numpy as np
from pydantic import BaseModel, BeforeValidator, AfterValidator
from typing import Annotated

def decode_base64_to_numpy(v: str) -> np.ndarray:
    """
    Декодирует Base64 строку напрямую в numpy массив.
    Выполняется до стандартных проверок Pydantic.
    """
    if not isinstance(v, str):
        raise ValueError("Ожидается строка в формате Base64")

    # Быстрое декодирование всего блока данных
    raw_bytes = base64.b64decode(v)
    # Прямое отображение байт в массив float32 без поэлементного копирования
    return np.frombuffer(raw_bytes, dtype=np.float32)

def validate_shape(tensor: np.ndarray) -> np.ndarray:
    """
    Проверяет, что тензор имеет нужную размерность (Fail-fast).
    Выполняется после успешного создания массива.
    """
    if tensor.ndim != 1:
        raise ValueError(f"Ожидается 1D вектор, получено измерений: {tensor.ndim}")

    EXPECTED_DIM = 512
    if tensor.shape[0] != EXPECTED_DIM:
        raise ValueError(f"Ожидается размерность {EXPECTED_DIM}, получено: {tensor.shape[0]}")

    return tensor

# Конструируем кастомный тип с помощью Annotated
EmbeddingTensor = Annotated[
    np.ndarray,
    BeforeValidator(decode_base64_to_numpy),
    AfterValidator(validate_shape)
]

class PredictRequest(BaseModel):
    # Pydantic автоматически применит весь пайплайн при создании объекта
    embedding: EmbeddingTensor

В этом коде BeforeValidator берет на себя самую тяжелую работу: он перехватывает входящую строку и с помощью np.frombuffer создает массив без единого цикла в Python. Затем AfterValidator проверяет бизнес-логику (размерность тензора).

Устранив узкое место в сериализации и внедрив строгую проверку формы данных, мы обеспечили беспрепятственное и безопасное прохождение запросов от сетевого сокета до математического ядра приложения.

Паттерн Dependency Injection для переиспользования тяжелых ML-моделей

Паттерн Dependency Injection для переиспользования тяжелых ML-моделей

Объект request.state, куда механизм lifespan сохраняет загруженные в память веса нейросетей, имеет критический архитектурный недостаток. Для статических анализаторов (mypy) и среды разработки этот объект является «черной дырой»: любые извлеченные из него атрибуты получают тип Any. Когда в приложении десятки эндпоинтов для инференса, эмбеддингов и метрик, ручное извлечение модели вида model = request.state.ml_model не только плодит дублирование кода, но и лишает разработчика автодополнения, позволяя ошибкам (например, опечатке в имени метода model.predct()) дожить до продакшена.

Решением этой проблемы выступает встроенный в FastAPI механизм Dependency Injection (DI) — инъекция зависимостей.

Dependency Injection — паттерн проектирования, при котором функция или объект не создает и не ищет ресурсы для своей работы самостоятельно, а декларирует их в сигнатуре. Фреймворк берет на себя задачу найти, подготовить и передать (инъецировать) эти ресурсы в момент вызова.

От Any к строгой типизации

Вместо того чтобы заставлять каждый эндпоинт копаться во внутренностях HTTP-запроса, мы делегируем эту задачу отдельной функции-провайдеру. Эта функция извлекает модель и, что самое главное, возвращает ее с явной аннотацией типа.

from fastapi import Request, Depends
from my_ml_lib import ResNet50 # Тяжелый класс модели

# Функция-провайдер зависимости
def get_vision_model(request: Request) -> ResNet50:
    return request.state.vision_model

# Эндпоинт декларирует зависимость
@app.post("/predict")
async def predict(
    payload: PredictRequest,
    model: ResNet50 = Depends(get_vision_model)
):
    # IDE знает, что model — это ResNet50. Доступно автодополнение.
    result = await asyncio.to_thread(model.forward, payload.tensor)
    return {"predictions": result}

В момент HTTP-запроса FastAPI видит маркер Depends(get_vision_model). Фреймворк приостанавливает выполнение эндпоинта, вызывает get_vision_model, передает ей текущий объект Request, забирает результат и пробрасывает его в переменную model.

Каскадные зависимости и роутинг моделей (A/B-тестирование)

Зависимости в FastAPI могут зависеть от других зависимостей. Это позволяет строить сложные цепочки логики до того, как запрос попадет в обработчик.

Представим сценарий A/B-тестирования. В lifespan загружены две версии модели: стабильная и экспериментальная. Нам нужно направлять запросы на нужную версию в зависимости от HTTP-заголовка X-Model-Version. Если заголовок отсутствует, используется версия по умолчанию. Если запрошена несуществующая версия, клиент должен получить ошибку до начала вычислений.

from fastapi import Header, HTTPException

def get_model_version(x_model_version: str = Header(default="v1")) -> str:
    allowed_versions = {"v1", "v2"}
    if x_model_version not in allowed_versions:
        raise HTTPException(status_code=400, detail="Unknown model version")
    return x_model_version

# Каскадная зависимость: требует Request и результат get_model_version
def get_ab_model(
    request: Request,
    version: str = Depends(get_model_version)
) -> ResNet50:
    # Извлекаем нужную модель из словаря в state
    return request.state.models[version]

@app.post("/predict/ab")
async def predict_ab(
    payload: PredictRequest,
    model: ResNet50 = Depends(get_ab_model)
):
    result = await asyncio.to_thread(model.forward, payload.tensor)
    return {"predictions": result}

Такой подход полностью изолирует бизнес-логику выбора модели от логики инференса. Эндпоинт остается чистым: он просто получает готовую к работе модель нужной версии.

Сравним два подхода к архитектуре эндпоинтов:

Характеристика Прямое извлечение (request.state) Инъекция зависимостей (Depends)
Типизация Отсутствует (Any) Строгая (через сигнатуру провайдера)
Дублирование кода Высокое (проверки ключей в каждом роуте) Нулевое (логика инкапсулирована)
Обработка ошибок Внутри бизнес-логики эндпоинта На этапе маршрутизации (Fail-fast)

Изоляция для CI/CD: подмена зависимостей

Самая большая боль MLOps при написании юнит-тестов — необходимость загружать реальные веса. Если каждый запуск pytest будет инициализировать модели по 5 ГБ, CI/CD пайплайн станет недопустимо медленным и потребует GPU-раннеров.

Паттерн DI позволяет элегантно обойти это ограничение через механизм app.dependency_overrides. Мы можем указать FastAPI заменить реальную функцию-провайдер на легковесный Mock-объект во время тестов.

# Файл тестов: test_api.py
from fastapi.testclient import TestClient
from main import app, get_vision_model

# Создаем заглушку, которая не загружает веса
class MockModel:
    def forward(self, tensor):
        return [0.99, 0.01] # Фейковый ответ

def get_mock_model():
    return MockModel()

# Подменяем зависимость
app.dependency_overrides[get_vision_model] = get_mock_model

client = TestClient(app)

def test_predict_endpoint():
    # Запрос пройдет успешно, используя MockModel вместо ResNet50
    response = client.post("/predict", json={"tensor_data": "..."})
    assert response.status_code == 200

Словарь dependency_overrides перехватывает вызов Depends(get_vision_model) и подставляет get_mock_model. Это позволяет тестировать сериализацию Pydantic, логику HTTP-ответов и права доступа за миллисекунды на обычных CPU-серверах.

Сборка архитектуры воедино

Объединив концепции валидации, асинхронности и инъекции зависимостей, мы получаем эталонный эндпоинт для ML-инференса.

Предположим, математическое ограничение нашей модели требует, чтобы размер батча NN не превышал 3232, а порог уверенности предсказания α0.85\alpha \geq 0.85. Валидацию батча берет на себя Pydantic, порог α\alpha можно передать как Query-параметр через DI, а тяжелые вычисления делегируются пулу потоков.

@app.post("/predict/production")
async def production_predict(
    payload: PredictRequest,           # 1. Pydantic валидирует Base64 и shape (N <= 32)
    model: MLModel = Depends(get_model), # 2. DI выдает типизированную модель
    alpha: float = Query(0.85, ge=0.0) # 3. FastAPI валидирует Query-параметр
):
    # 4. Освобождаем Event Loop, отправляя CPU-bound задачу в поток
    raw_preds = await asyncio.to_thread(model.predict, payload.tensor)

    # 5. Фильтрация по порогу уверенности
    filtered = [p for p in raw_preds if p.confidence >= alpha]

    return {"results": filtered}

Этот код декларативен. Он описывает что нужно для работы, а не как это получить. Однако, даже идеально спроектированный синхронный инференс имеет физический предел. Если модель обрабатывает запрос дольше 30 секунд, клиентское соединение оборвется по HTTP-таймауту, независимо от того, насколько красиво написан код.

Архитектура асинхронного инференса: фоновые задачи BackgroundTasks и Celery

Архитектура асинхронного инференса: фоновые задачи BackgroundTasks и Celery

Представьте, что вы успешно обернули тяжелую модель машинного зрения в FastAPI. Вы настроили lifespan для загрузки весов, использовали Pydantic v2 для валидации тензоров и пробросили модель через Dependency Injection. Всё работает идеально на локальных тестах. Но когда вы выкатываете сервис в продакшен и пользователи начинают загружать минутные видео для обработки, логи заполняются ошибками 504 Gateway Timeout.

Проблема кроется в архитектуре HTTP. Большинство балансировщиков нагрузки и прокси-серверов (например, Nginx или AWS ALB) имеют жесткие лимиты на время ожидания ответа. Если время выполнения инференса T>30T > 30 секунд, соединение принудительно разрывается. Клиент получает ошибку, хотя ваш сервер всё ещё продолжает «молотить» видеокартой, впустую сжигая ресурсы.

Чтобы решить эту проблему, нам необходимо разорвать связь между приемом запроса и выдачей результата.

Иллюзия простого решения: BackgroundTasks

FastAPI предоставляет встроенный инструмент для отложенного выполнения кода — BackgroundTasks. Это класс, который позволяет добавить функцию в очередь на выполнение после того, как HTTP-ответ уже отправлен клиенту.

Синтаксис выглядит крайне привлекательно:

from fastapi import APIRouter, BackgroundTasks, Depends
from my_ml_module import get_model, MLModel
import asyncio

router = APIRouter()

def process_video_task(video_id: str, model: MLModel):
    # Тяжелая CPU-bound задача
    result = model.predict(video_id)
    # Сохранение результата в БД...

@router.post("/process")
async def process_video(
    video_id: str,
    background_tasks: BackgroundTasks,
    model: MLModel = Depends(get_model)
):
    # Передаем синхронную задачу в фон
    background_tasks.add_task(process_video_task, video_id, model)
    return {"message": "Видео принято в обработку"}

Клиент мгновенно получает ответ {"message": "Видео принято в обработку"}, HTTP-соединение закрывается, таймаута нет. Кажется, задача решена? Для MLOps — категорически нет.

Использование BackgroundTasks для тяжелого ML-инференса — это антипаттерн, ведущий к нестабильности системы.

Почему BackgroundTasks не подходит для тяжелых задач:

  1. Отсутствие отказоустойчивости (Persistence): Задачи хранятся в оперативной памяти процесса Uvicorn. Если под в Kubernetes перезагрузится из-за нехватки памяти (OOM Killed) или обновления, все фоновые задачи исчезнут навсегда.
  2. Монолитное масштабирование: Фоновые задачи выполняются на тех же серверах, что и API. Вы не можете масштабировать вычислительные узлы (с GPU) отдельно от узлов приема запросов (CPU).
  3. Блокировка пула потоков: Как мы разбирали ранее, синхронный код будет отправлен в Thread pool. Если придет 100 запросов, пул потоков API-сервера исчерпается, и он перестанет отвечать даже на простые запросы вроде /health.

BackgroundTasks идеален для легковесных I/O-операций: отправки email, записи метрики или очистки временного файла. Для инференса нужна иная парадигма.

Парадигма асинхронного инференса

Настоящая масштабируемость достигается физическим разделением системы на две независимые части: API-слой (быстрый, stateless, принимает запросы) и Worker-слой (тяжелый, stateful, выполняет ML-модели).

Связующим звеном между ними выступает Очередь задач (Task Queue).

В этой архитектуре API-сервер вообще не загружает ML-модель. Его единственная задача — принять данные, валидировать их, положить сообщение в очередь и вернуть клиенту уникальный идентификатор задачи (task_id).

Celery как стандарт распределенных очередей

Celery — это мощная асинхронная система очередей задач, де-факто стандарт в экосистеме Python.

Архитектура системы с Celery состоит из трех компонентов:

  1. Producer (FastAPI): Создает задачу и отправляет её в очередь.
  2. Broker (Брокер сообщений): Внешняя база данных (обычно RabbitMQ или Redis), которая надежно хранит очередь задач на диске или в памяти.
  3. Consumer / Worker (Celery процесс): Отдельный процесс (или сотни процессов на разных серверах), который берет задачу из брокера, выполняет её и записывает результат.

Паттерн Polling (Опрос состояния)

Поскольку HTTP-запрос завершается до того, как модель выдаст результат, клиенту нужен способ этот результат получить. Самый надежный подход в REST API — паттерн Polling.

Мы реализуем два эндпоинта. Первый ставит задачу в очередь:

from fastapi import APIRouter
from celery_app import ml_task # Импорт настроенной задачи Celery

router = APIRouter()

@router.post("/predict")
async def create_prediction_task(data: dict):
    # Отправляем задачу в брокер.
    # Celery возвращает объект AsyncResult, содержащий task_id
    task = ml_task.delay(data)

    # Мгновенно возвращаем ID клиенту
    return {"task_id": task.id, "status": "PENDING"}

Второй эндпоинт позволяет клиенту узнать статус по task_id. Поиск задачи в брокере по ID имеет вычислительную сложность O(1)O(1), поэтому этот эндпоинт работает молниеносно и не нагружает систему.

from fastapi import APIRouter
from celery.result import AsyncResult
from celery_app import celery_instance

@router.get("/status/{task_id}")
async def get_task_status(task_id: str):
    # Извлекаем состояние задачи из Celery
    task_result = AsyncResult(task_id, app=celery_instance)

    if task_result.ready():
        return {
            "task_id": task_id,
            "status": "SUCCESS",
            "result": task_result.result
        }

    return {
        "task_id": task_id,
        "status": task_result.status # PENDING или STARTED
    }

Клиент (например, фронтенд) делает POST-запрос, получает task_id, а затем каждые 5 секунд отправляет GET-запрос на /status/{task_id}, пока не получит SUCCESS.

Сравнение подходов

Чтобы закрепить понимание, сведем разницу между встроенными задачами и полноценной очередью в таблицу:

Характеристика BackgroundTasks (FastAPI) Celery + Broker
Хранение задач В оперативной памяти API-процесса В надежном внешнем брокере (Redis/RabbitMQ)
Отказоустойчивость При падении сервера задачи теряются При падении воркера задача возвращается в очередь
Масштабирование API и инференс масштабируются только вместе API (CPU) и Воркеры (GPU) масштабируются раздельно
Трекинг статуса Нет (задача просто выполняется в фоне) Да (PENDING, STARTED, SUCCESS, FAILURE)
Сложность инфраструктуры Нулевая (встроено в FastAPI) Высокая (нужно поднимать брокер и воркеры)

Переход от BackgroundTasks к Celery — это переход от скрипта к распределенной системе. Вы усложняете инфраструктуру, но взамен получаете гарантию того, что ни один запрос пользователя не потеряется, а дорогие GPU-серверы можно будет масштабировать независимо от легковесных API-шлюзов.

В следующей главе мы спустимся на уровень ниже и разберем, как именно FastAPI общается с брокерами сообщений, и почему выбор между RabbitMQ и Redis критически важен для гарантии доставки ML-задач.

Интеграция очередей сообщений: FastAPI, RabbitMQ и Redis для очередей задач

Интеграция очередей сообщений: FastAPI, RabbitMQ и Redis для очередей задач

Представьте, что ваш API принимает 500 запросов в секунду на генерацию изображений, но кластер GPU-воркеров способен обработать только 50. Куда деваются остальные 450 запросов? Если оставить их в оперативной памяти FastAPI-сервера, он быстро упадет с ошибкой Out Of Memory. Если сбрасывать — клиенты получат каскад ошибок 503 Service Unavailable. Разделение на API и Worker-слои требует надежной «комнаты ожидания» — инфраструктуры, которая гарантирует, что ни одна задача не потеряется при пиковых нагрузках или падении серверов.

Анатомия распределенной системы сообщений

Асинхронный инференс опирается на два независимых инфраструктурных компонента, которые часто путают: брокер сообщений и хранилище результатов.

Компонент Роль в MLOps Оптимальный инструмент Паттерн доступа
Message Broker (Брокер сообщений) Принимает задачи от FastAPI и распределяет их между воркерами. Гарантирует доставку. RabbitMQ FIFO-очередь (First In, First Out). Данные удаляются после успешной обработки.
Result Backend (Хранилище результатов) Хранит готовые предикты модели, пока клиент не заберет их по task_id. Redis Key-Value хранилище (Ключ-Значение). Данные живут ограниченное время.

Хотя Redis можно использовать и как брокер, и как бэкенд, в высоконагруженных ML-системах их роли строго разделяют. RabbitMQ спроектирован для сложной маршрутизации и надежного удержания невыполненных задач, а Redis идеален для молниеносного чтения результатов по ключу.

RabbitMQ: Гарантия доставки и защита от CUDA OOM

Главная проблема ML-воркеров — их нестабильность. Тяжелая модель может упасть из-за нехватки видеопамяти (CUDA Out Of Memory), если на вход пришел аномально большой тензор. Если в этот момент задача находилась в памяти воркера, она исчезнет навсегда.

RabbitMQ решает эту проблему через механизм подтверждений.

Acknowledgment (ACK) — сигнал, который воркер отправляет брокеру только после успешного завершения задачи. До получения этого сигнала брокер считает задачу «в работе», но не удаляет ее из очереди.

Если процесс с ML-моделью внезапно завершается (например, процесс убит OOM Killer операционной системы), TCP-соединение с RabbitMQ разрывается. Брокер понимает, что ACK не придет, и автоматически возвращает задачу в очередь, передавая её другому, живому воркеру.

Отравленные сообщения и Dead Letter Exchange

Механизм ACK создает новую угрозу: что если конкретное изображение повреждено так, что вызывает падение любого воркера? Задача вернется в очередь, попадет к следующему воркеру, убьет его, и так до бесконечности. Это называется «отравленным сообщением» (poison pill).

Для защиты применяется паттерн DLX:

  1. Воркер ловит исключение при обработке тензора.
  2. Вместо ACK он отправляет брокеру сигнал NACK (Negative Acknowledgment) без права возврата в основную очередь.
  3. RabbitMQ перенаправляет эту задачу в Dead Letter Exchange (DLX) — специальную очередь для «мертвых» задач.
  4. Инженеры могут позже проанализировать DLX, чтобы понять причину сбоя модели, не блокируя при этом работу продакшена.

Redis: Быстрая отдача результатов с TTL

Когда воркер успешно завершает инференс, результат (например, JSON с bounding boxes для объектов на фото) нужно где-то сохранить. FastAPI-сервер, обрабатывающий эндпоинт статуса, должен мгновенно получить к нему доступ по task_id.

Redis — это in-memory база данных, которая хранит информацию в оперативной памяти. Поиск по ключу celery-task-meta-{task_id} занимает доли миллисекунды.

Ключевая настройка Redis для MLOps — это TTL (Time-To-Live). Если клиент запросил генерацию, но закрыл браузер и никогда не вернулся за результатом, сохраненный тензор останется в памяти. Без автоматической очистки Redis быстро исчерпает RAM. Установка TTL (например, 3600 секунд) гарантирует, что невостребованные результаты самоуничтожатся через час.

Математика очередей: Закон Литтла

Чтобы правильно масштабировать RabbitMQ и воркеры, необходимо понимать пропускную способность системы. В теории массового обслуживания это описывается Законом Литтла:

L=λWL = \lambda W

Где:

  • LL — среднее количество задач в системе (в очереди RabbitMQ + в процессе обработки на GPU).
  • λ\lambda — интенсивность поступления задач от пользователей (запросов в секунду).
  • WW — среднее время нахождения задачи в системе (время ожидания + время инференса модели).

Практический пример: Ваш FastAPI принимает λ=20\lambda = 20 запросов в секунду на транскрибацию аудио. Каждая задача в среднем проводит в системе W=15W = 15 секунд (из них 14 секунд в очереди и 1 секунда на GPU). Согласно формуле, L=20×15=300L = 20 \times 15 = 300. Это значит, что в любой момент времени в вашей системе находится 300 активных задач. Если лимит длины очереди в RabbitMQ установлен в 100 сообщений, система неизбежно начнет отбрасывать новые запросы. Понимание этого закона позволяет точно настраивать лимиты очередей (Queue Length Limits) и правила автомасштабирования воркеров.

Интеграция в FastAPI

Связующим звеном между FastAPI, RabbitMQ и Redis чаще всего выступает библиотека Celery. FastAPI выступает в роли Producer (отправителя), а конфигурация брокера и бэкенда задается через URI.

from fastapi import FastAPI
from pydantic import BaseModel
from celery import Celery

# Инициализация Celery с указанием RabbitMQ как брокера и Redis как бэкенда
celery_app = Celery(
    "ml_tasks",
    broker="amqp://user:password@rabbitmq:5672//",
    backend="redis://redis:6379/0"
)

# Настройка TTL для результатов в Redis (1 час)
celery_app.conf.result_expires = 3600

app = FastAPI()

class InferenceRequest(BaseModel):
    text: str
    max_length: int = 50

@app.post("/generate")
async def generate_text(request: InferenceRequest):
    # Отправка задачи в RabbitMQ. FastAPI не ждет выполнения.
    task = celery_app.send_task(
        "worker.predict_text",
        args=[request.text, request.max_length]
    )
    return {"task_id": task.id, "status": "queued"}

В этом коде FastAPI мгновенно отвечает клиенту. Под капотом send_task сериализует аргументы и отправляет их по протоколу AMQP в RabbitMQ. Вся тяжелая работа с TCP-соединениями, повторными попытками отправки и пулами соединений скрыта внутри Celery.

Архитектура с брокером и хранилищем идеально решает задачи классического инференса, где результатом является единый конечный объект (JSON, изображение). Однако, если мы работаем с большими языковыми моделями (LLM), ожидание полного завершения генерации перед сохранением в Redis создает неприемлемую задержку для пользователя. В таких случаях парадигма меняется от очередей к прямым потоковым соединениям.

Потоковая передача данных: StreamingResponse для генеративных моделей (LLM)

Потоковая передача данных: StreamingResponse для генеративных моделей (LLM)

Представьте, что пользователь отправляет запрос к вашей новой большой языковой модели (LLM). Модель генерирует ответ со скоростью 20 токенов в секунду. Если итоговый текст состоит из 600 токенов, генерация займет 30 секунд. Если мы используем архитектуру очередей и поллинга, которую построили ранее, пользователь будет смотреть на пустой экран ровно полминуты, прежде чем получит готовый результат целиком. В современных реалиях интерфейсов это катастрофа: большинство пользователей закроют вкладку уже на пятой секунде ожидания.

Генеративные модели (текст, аудио) обладают фундаментальным отличием от задач классификации или обработки видео: их результат рождается последовательно. Нам не нужно ждать окончания вычислений, чтобы начать отдавать пользу клиенту. Мы можем отправлять данные по мере их появления.

Ограничения очередей и смена парадигмы

Архитектура с RabbitMQ и Redis идеально решает проблему надежности для монолитных вычислительных задач. Брокер гарантирует, что задача не потеряется, а Result Backend хранит итоговый артефакт.

Однако для LLM нам требуется мгновенная обратная связь. Здесь парадигма асинхронного инференса через очереди уступает место потоковой передаче данных (Streaming). Вместо того чтобы разрывать HTTP-соединение и заставлять клиента опрашивать сервер, мы удерживаем одно HTTP-соединение открытым и непрерывно «докидываем» в него новые порции данных.

Характеристика Очереди задач (Celery, RabbitMQ) Потоковая передача (StreamingResponse)
Характер результата Монолитный (весь ответ целиком) Последовательный (токены, чанки)
Жизненный цикл HTTP Короткий (принял задачу — ответил 202 Accepted) Длинный (соединение открыто до конца генерации)
UX клиента Ожидание 100% готовности Чтение текста по мере его «печатания»
Риск при сбое Минимальный (задача вернется в очередь) Потеря текущей генерации (клиенту придется повторить запрос)

Механика потоковой передачи: Chunked Transfer Encoding

В основе потоковой отдачи в FastAPI лежит механизм стандарта HTTP/1.1 — Chunked Transfer Encoding.

Обычно сервер отправляет заголовок Content-Length, сообщая клиенту точный размер ответа в байтах. Но при генерации текста LLM мы заранее не знаем, сколько токенов выдаст модель до того, как сгенерирует токен остановки (EOS).

При использовании механизма чанков сервер опускает заголовок Content-Length и вместо этого отправляет заголовок Transfer-Encoding: chunked. Данные передаются кусками (чанками), где перед каждым куском указывается его размер. Клиентский браузер или библиотека читает этот поток непрерывно, пока сервер не пришлет чанк нулевой длины, что означает конец передачи.

В FastAPI за этот процесс отвечает класс StreamingResponse. Он принимает на вход асинхронный генератор (функцию, использующую yield вместо return) и автоматически оборачивает отдаваемые данные в нужные HTTP-заголовки.

Метрики генеративного инференса

При проектировании потокового инференса фокус оптимизации смещается с общей пропускной способности на две специфические временные метрики:

Ttotal=TTFT+N×TPOTT_{total} = TTFT + N \times TPOT

Разберем элементы этой формулы:

  • TtotalT_{total} — общее время ответа системы для пользователя.
  • TTFTTTFT (Time To First Token) — время до первого токена. Это задержка от момента отправки запроса до получения первого слова. Включает в себя сетевую задержку и время, которое модель тратит на чтение и осмысление промпта (prefill phase).
  • NN — количество сгенерированных токенов ответа.
  • TPOTTPOT (Time Per Output Token) — время генерации каждого последующего токена.

Практический пример: Допустим, пользователь отправил длинный промпт. Модель обрабатывала его 0.5 секунды (TTFT=0.5TTFT = 0.5), а затем начала выдавать ответ из 100 токенов (N=100N = 100) со скоростью 0.05 секунды на токен (TPOT=0.05TPOT = 0.05). Без потоковой передачи пользователь ждал бы 0.5+100×0.05=5.50.5 + 100 \times 0.05 = 5.5 секунд. С потоковой передачей он увидит первое слово уже через 0.5 секунды, а остальной текст будет плавно появляться на экране.

Интеграция CPU-bound генерации и асинхронного потока

Здесь мы сталкиваемся с главным архитектурным вызовом. Как мы помним, инференс нейросетей — это тяжелая математика, блокирующая процессор (CPU-bound задача). Если мы просто напишем цикл генерации токенов с использованием yield внутри асинхронного эндпоинта, мы намертво заблокируем Event loop. Ни один другой пользователь не сможет получить даже статус от сервера.

Нам необходимо построить мост: запустить синхронную тяжелую генерацию в пуле потоков, но при этом асинхронно забирать оттуда токены и отправлять их клиенту.

Для этого идеально подходит паттерн «Производитель-Потребитель» с использованием потокобезопасной очереди asyncio.Queue.

Рассмотрим реализацию этого паттерна в коде:

import asyncio
from fastapi import APIRouter
from fastapi.responses import StreamingResponse

router = APIRouter()

# 1. Синхронная функция генерации (Производитель)
# Выполняется в отдельном системном потоке, чтобы не блокировать Event loop
def generate_tokens_sync(prompt: str, queue: asyncio.Queue, loop: asyncio.AbstractEventLoop):
    # Имитация работы LLM
    tokens = ["FastAPI", " ", "отлично", " ", "подходит", " ", "для", " ", "LLM", "."]

    for token in tokens:
        import time
        time.sleep(0.1) # Имитация TPOT (CPU-bound вычисления)

        # Безопасная передача токена из системного потока в асинхронную очередь
        asyncio.run_coroutine_threadsafe(queue.put(token), loop)

    # Сигнал об окончании генерации (EOS)
    asyncio.run_coroutine_threadsafe(queue.put(None), loop)

# 2. Асинхронный генератор (Потребитель)
# Работает в Event loop, не блокируя его при ожидании новых токенов
async def token_generator(prompt: str):
    queue = asyncio.Queue()
    loop = asyncio.get_running_loop()

    # Запускаем тяжелую генерацию в пуле потоков
    asyncio.create_task(
        asyncio.to_thread(generate_tokens_sync, prompt, queue, loop)
    )

    # Асинхронно читаем очередь и отдаем токены клиенту
    while True:
        token = await queue.get()
        if token is None: # Проверка на сигнал остановки
            break
        yield token

# 3. Эндпоинт FastAPI
@router.get("/generate")
async def generate_text(prompt: str):
    # StreamingResponse читает генератор и отправляет чанки по сети
    return StreamingResponse(
        token_generator(prompt),
        media_type="text/event-stream"
    )

Разбор архитектуры моста

В этом коде мы элегантно связываем две парадигмы:

  1. Функция generate_tokens_sync делает тяжелую математику. Она ничего не знает об HTTP. Ее единственная задача — вычислять токены и складывать их в трубу (queue.put). Поскольку она работает в другом системном потоке, мы обязаны использовать asyncio.run_coroutine_threadsafe, чтобы безопасно передать данные обратно в главный Event loop.
  2. Функция token_generator является асинхронной. Она мгновенно засыпает на строке await queue.get(), освобождая Event loop для других запросов. Как только математический поток кладет токен в очередь, генератор просыпается, делает yield и отдает токен в StreamingResponse.

Использование media_type="text/event-stream" в эндпоинте указывает клиенту, что мы используем формат Server-Sent Events (SSE) — надстройку над потоковой передачей, которая стандартизирует чтение таких потоков в веб-браузерах с помощью встроенного JavaScript-интерфейса EventSource.

Эта архитектура позволяет масштабировать эндпоинты генерации. Вы можете заменить заглушку generate_tokens_sync на реальный вызов HuggingFace generate() или интегрировать специализированные движки вроде vLLM, сохранив при этом отзывчивость самого веб-сервера FastAPI.

Мониторинг и наблюдаемость: интеграция Prometheus метрик и OpenTelemetry

Мониторинг и наблюдаемость: интеграция Prometheus метрик и OpenTelemetry

Представьте ситуацию: вы успешно внедрили потоковую генерацию токенов и настроили асинхронные очереди для тяжелых задач. Система работает, но в пятницу вечером метрика времени ответа резко возрастает, а клиенты начинают получать ошибку 504. Что именно сломалось? Завис Event loop в FastAPI? Переполнилась очередь RabbitMQ? Или GPU на воркере начал троттлить из-за перегрева? Без внедренной системы наблюдаемости поиск ответа превращается в гадание на кофейной гуще.

В этой главе мы свяжем разрозненные компоненты нашей MLOps-архитектуры в единую прозрачную систему, используя метрики Prometheus и распределенную трассировку OpenTelemetry.

Наблюдаемость против классического мониторинга

Мониторинг отвечает на вопрос «Что сломалось?», а наблюдаемость — «Почему это сломалось?».

Наблюдаемость (Observability) — это мера того, насколько хорошо внутренние состояния системы могут быть поняты на основе её внешних выходных данных.

Основы теории управления, Рудольф Калман

Для достижения наблюдаемости в MLOps мы опираемся на два базовых столпа:

  1. Метрики (Metrics) — агрегированные числовые показатели системы за интервал времени.
  2. Трейсы (Traces) — детализированный путь одного конкретного запроса через все микросервисы.
Инструмент Сущность Фокус Пример вопроса
Prometheus Метрики Агрегация и тренды Какой процент запросов обрабатывается дольше 2 секунд?
OpenTelemetry Трейсы Контекст и путь На каком этапе завис запрос пользователя с ID 8472?

Prometheus: Pull-модель и типы метрик

Prometheus работает по Pull-модели: он не ждет, пока приложение пришлет ему данные, а сам периодически опрашивает специальный HTTP-эндпоинт (обычно /metrics), который должно предоставлять ваше приложение.

Чтобы метрики были полезными, важно выбрать правильный тип для каждой задачи.

1. Counter (Счетчик)

Это монотонно возрастающее значение. Оно никогда не уменьшается (кроме случаев перезапуска сервиса). Идеально подходит для подсчета общего количества запросов или ошибок.

Сами по себе абсолютные значения счетчика редко бывают полезны. В Prometheus мы вычисляем скорость изменения счетчика — рейт. Математически частота запросов вычисляется как Rate=ΔCΔtRate = \frac{\Delta C}{\Delta t}, где ΔC\Delta C — изменение значения счетчика за интервал времени Δt\Delta t.

2. Gauge (Датчик)

Значение, которое может как увеличиваться, так и уменьшаться. В контексте инференса это критически важный тип.

Примеры использования Gauge:

  • Текущее количество задач в очереди RabbitMQ.
  • Объем занятой видеопамяти (VRAM) на GPU.
  • Количество активных потоковых соединений (SSE).

3. Histogram (Гистограмма)

Гистограмма распределяет наблюдаемые значения по заранее заданным корзинам (buckets) и подсчитывает их количество. Это единственный правильный способ измерять задержки (latency), такие как TTFT или TPOT.

Почему нельзя использовать среднее арифметическое? Если 99 запросов выполнились за 10 мс, а один завис на 10 секунд, среднее время составит около 110 мс. Эта цифра скроет катастрофический сбой для одного пользователя и покажет ложное ухудшение для остальных. Гистограммы позволяют вычислять перцентили (например, P99P_{99} — значение, хуже которого работает только 1% запросов), давая реальную картину пользовательского опыта.

Интеграция метрик в FastAPI

Для создания эндпоинта /metrics мы воспользуемся официальной библиотекой prometheus_client. Реализуем сбор метрики TPOT (Time Per Output Token), с которой мы работали при настройке LLM-стриминга.

import time
from fastapi import FastAPI, Response
from prometheus_client import Histogram, generate_latest, CONTENT_TYPE_LATEST

app = FastAPI(lifespan=lifespan)

# Определяем гистограмму с кастомными корзинами (в секундах)
# Корзины подобраны под типичные значения TPOT: от 10 мс до 1 секунды
TPOT_HISTOGRAM = Histogram(
    "llm_tpot_seconds",
    "Time Per Output Token in seconds",
    buckets=[0.01, 0.05, 0.1, 0.2, 0.5, 1.0]
)

@app.get("/metrics")
async def metrics():
    """Эндпоинт, который будет опрашивать Prometheus."""
    return Response(content=generate_latest(), media_type=CONTENT_TYPE_LATEST)

@app.post("/generate")
async def generate_text():
    # Имитация генерации токена
    start_time = time.perf_counter()

    # ... здесь происходит инференс модели ...
    time.sleep(0.045)

    generation_time = time.perf_counter() - start_time

    # Записываем наблюдение в гистограмму
    TPOT_HISTOGRAM.observe(generation_time)

    return {"status": "ok"}

Теперь, если Prometheus обратится к /metrics, он получит текстовый снимок текущего состояния гистограммы, готовый к агрегации.

OpenTelemetry: Распределенная трассировка

Метрики покажут нам, что P99P_{99} для инференса вырос до 5 секунд. Но они не скажут, где именно произошла задержка: при сериализации Pydantic, в сети между FastAPI и Redis, или внутри самой модели. Здесь в игру вступает OpenTelemetry (OTel).

Архитектура трассировки строится на двух терминах:

  • Trace (Трейс) — полное дерево выполнения одного бизнес-запроса от начала до конца. Представляет собой направленный ациклический граф.
  • Span (Спан) — один логический шаг внутри трейса (например, «запрос к БД» или «вызов функции predict»). Каждый спан имеет время начала, время конца и метаданные.

Проблема потери контекста (Context Propagation)

Когда клиент делает HTTP-запрос к FastAPI, OTel автоматически создает корневой спан и генерирует уникальный trace_id. Но как только FastAPI отправляет задачу в RabbitMQ для асинхронного инференса, связь обрывается. Воркер Celery ничего не знает о trace_id из FastAPI и начнет свой собственный независимый трейс.

Чтобы этого избежать, применяется Context Propagation (Проброс контекста). FastAPI должен внедрить (inject) текущий trace_id в метаданные сообщения RabbitMQ, а Celery-воркер должен извлечь (extract) его перед началом работы, чтобы продолжить существующий трейс.

Автоматическая инструментация FastAPI

OpenTelemetry предоставляет готовые инструменты, которые оборачивают стандартные библиотеки Python и автоматически создают спаны для HTTP-запросов, запросов к базам данных и брокерам сообщений.

from fastapi import FastAPI
from opentelemetry import trace
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor, ConsoleSpanExporter

# 1. Настраиваем провайдер трассировки
provider = TracerProvider()
processor = BatchSpanProcessor(ConsoleSpanExporter())
provider.add_span_processor(processor)
trace.set_tracer_provider(provider)

app = FastAPI()

# 2. Инструментируем FastAPI
# Это автоматически создаст спаны для каждого входящего HTTP-запроса
FastAPIInstrumentor.instrument_app(app)

tracer = trace.get_tracer(__name__)

@app.post("/predict")
async def predict_heavy_model():
    # 3. Создание кастомного внутреннего спана для тяжелой операции
    with tracer.start_as_current_span("tensor_preprocessing") as span:
        span.set_attribute("tensor.shape", "[1, 3, 224, 224]")
        # ... логика предобработки ...

    return {"result": "success"}

В этом примере FastAPIInstrumentor берет на себя всю рутину: он перехватывает входящие запросы, читает HTTP-заголовки (например, traceparent, если запрос пришел от другого микросервиса), создает корневой спан для эндпоинта /predict и фиксирует время ответа. А блок with tracer.start_as_current_span позволяет нам детализировать внутреннюю логику, измеряя конкретно этап подготовки тензора.

Таким образом, объединяя агрегированные метрики Prometheus для выявления аномалий и детализированные трейсы OpenTelemetry для поиска первопричин, мы получаем полностью наблюдаемую систему. Мы больше не гадаем, почему инференс стал медленным — мы видим это на графиках и в деревьях вызовов.

Контейнеризация и оптимизация Docker-манифестов для FastAPI с ML-зависимостями

Контейнеризация и оптимизация Docker-манифестов для FastAPI с ML-зависимостями

Стандартная команда docker build для проекта с FastAPI, PyTorch и HuggingFace Transformers часто рождает монстра размером 6–8 гигабайт. Загрузка такого образа в реестр и его скачивание на продакшен-серверы занимает минуты. В условиях динамической нагрузки, когда система пытается масштабироваться под наплыв пользователей, эти минуты превращаются в отказы в обслуживании.

Время готовности новой реплики сервиса = (Размер образа / Пропускная способность сети) + Время старта приложения. Например, при масштабировании воркера в облаке с каналом 125 МБ/с скачивание образа весом 6000 МБ займет 48 секунд. Добавим 15 секунд на загрузку весов модели в фазе lifespan, и получим больше минуты простоя. Уменьшение образа до 600 МБ сократит сетевую задержку до 4.8 секунд.

Анатомия слоев и инвалидация кэша

Docker собирает образы послойно. Каждая инструкция в манифесте (RUN, COPY, ADD) создает новый неизменяемый слой, который сохраняется в локальном кэше.

Кэш Docker работает по принципу домино: если слой изменяется, все последующие слои автоматически инвалидируются и пересобираются с нуля.

Официальная документация Docker

В контексте MLOps установка зависимостей через pip install — самая долгая операция, которая может занимать до 10 минут из-за компиляции C-расширений и скачивания гигабайтных библиотек. Порядок команд критически влияет на скорость итераций разработки.

Подход Структура манифеста Результат
Наивный COPY . .<br>RUN pip install -r requirements.txt Любая правка в main.py ломает кэш на шаге COPY. Установка тяжелых ML-библиотек запускается заново при каждом билде.
Оптимальный COPY requirements.txt .<br>RUN pip install -r requirements.txt<br>COPY . . Исходный код копируется строго после установки зависимостей. Правки в коде не вызывают переустановку библиотек.

Диета для ML-зависимостей: CPU против GPU

Архитектура асинхронного инференса физически разделяет легковесный API-слой и тяжелый Worker-слой. Эти компоненты требуют принципиально разных Docker-образов.

По умолчанию команда pip install torch скачивает пакет, включающий в себя скомпилированные бинарные файлы CUDA-рантайма. Это добавляет около 2.5 ГБ к размеру образа.

Если FastAPI-контейнер выполняет только роль шлюза — принимает HTTP-запросы, валидирует тензоры через Pydantic-модели и отправляет задачи в RabbitMQ — ему не нужен доступ к GPU. Для API-узла необходимо принудительно устанавливать CPU-версии библиотек, используя специализированные индексы пакетов:

# Установка легковесной CPU-версии PyTorch для API-шлюза
RUN pip install torch --extra-index-url https://download.pytorch.org/whl/cpu

Для Worker-узла, который фактически выполняет инференс и требует доступа к видеокарте, базовый образ должен содержать драйверы NVIDIA, а зависимости устанавливаются в стандартном режиме.

Паттерн Multi-stage сборки

Даже при использовании CPU-версий библиотек, в процессе установки часто требуются системные компиляторы (gcc, g++) и заголовочные файлы (python3-dev). Оставлять их в финальном образе небезопасно и неэффективно.

Multi-stage сборка решает эту проблему путем создания нескольких изолированных сред в рамках одного Dockerfile.

  1. Builder (Сборщик): содержит все необходимые утилиты для компиляции. Здесь устанавливаются зависимости и создается виртуальное окружение (Virtual Environment).
  2. Runner (Исполнитель): использует минималистичный базовый образ (например, python:3.11-slim). В него копируется только готовое виртуальное окружение из сборщика и исходный код приложения.

Оптимизированный манифест FastAPI-приложения

Объединим управление кэшем, изоляцию зависимостей и безопасность (запуск от имени непривилегированного пользователя) в единый Dockerfile для нашего API-слоя.

# ==========================================
# Stage 1: Builder
# ==========================================
FROM python:3.11-slim as builder

# Установка системных зависимостей для компиляции
RUN apt-get update && apt-get install -y --no-install-recommends \
    build-essential \
    && rm -rf /var/lib/apt/lists/*

# Создание виртуального окружения
RUN python -m venv /opt/venv
ENV PATH="/opt/venv/bin:$PATH"

# Кэширование установки зависимостей
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt \
    --extra-index-url https://download.pytorch.org/whl/cpu

# ==========================================
# Stage 2: Runner
# ==========================================
FROM python:3.11-slim as runner

# Создание непривилегированного пользователя для безопасности
RUN useradd -m -r appuser

# Копирование только готового окружения из builder
COPY --from=builder /opt/venv /opt/venv
ENV PATH="/opt/venv/bin:$PATH"

# Настройка рабочей директории
WORKDIR /app
COPY --chown=appuser:appuser . /app

# Переключение на безопасного пользователя
USER appuser

# Запуск приложения
CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "8000"]

Этот манифест гарантирует минимальный размер образа, отсутствие уязвимостей, связанных с правами root, и мгновенную пересборку при изменении бизнес-логики в Python-коде.

Оркестрация инфраструктуры

Разрозненные компоненты системы инференса объединяются в единую сеть через Docker Compose. Это позволяет настроить проброс контекста трассировки и сбор метрик Prometheus без сложной настройки сетевых доступов — сервисы обращаются друг к другу по внутренним DNS-именам, совпадающим с названиями контейнеров.

version: '3.8'

services:
  api:
    build:
      context: .
      dockerfile: Dockerfile.api
    ports:
      - "8000:8000"
    environment:
      - BROKER_URL=amqp://guest:guest@rabbitmq:5672//
      - RESULT_BACKEND=redis://redis:6379/0
    depends_on:
      - rabbitmq
      - redis

  worker:
    build:
      context: .
      dockerfile: Dockerfile.worker
    deploy:
      resources:
        reservations:
          devices:
            - driver: nvidia
              count: 1
              capabilities: [gpu]
    environment:
      - BROKER_URL=amqp://guest:guest@rabbitmq:5672//
      - RESULT_BACKEND=redis://redis:6379/0

  rabbitmq:
    image: rabbitmq:3-management
    ports:
      - "15672:15672"

  redis:
    image: redis:7-alpine

  prometheus:
    image: prom/prometheus
    volumes:
      - ./prometheus.yml:/etc/prometheus/prometheus.yml
    ports:
      - "9090:9090"

В этой конфигурации Worker-узел получает эксклюзивный доступ к GPU через директиву deploy.resources, в то время как API-узел остается легковесным и может масштабироваться горизонтально независимо от вычислительных мощностей инференса.

Стратегии масштабирования: ASGI-серверы, Gunicorn, Uvicorn и балансировка нагрузки

Стратегии масштабирования: ASGI-серверы, Gunicorn, Uvicorn и балансировка нагрузки

Мы упаковали наш API и ML-воркеры в легковесные Docker-контейнеры, разделив сборку и среду выполнения. Теперь представьте: вы разворачиваете контейнер с FastAPI на мощном сервере с 32 ядрами CPU. Запускаете систему, трафик резко возрастает до тысяч запросов в секунду. Внезапно API начинает отдавать таймауты. Вы заходите на сервер, открываете мониторинг ресурсов и видите парадоксальную картину: одно ядро загружено на 100%, а остальные 31 ядро абсолютно простаивают.

Почему асинхронный код не спас ситуацию и как заставить систему утилизировать всё доступное железо?

Анатомия ASGI и предел одного процесса

FastAPI сам по себе не является веб-сервером. Это фреймворк, который определяет маршруты и логику. Чтобы HTTP-запрос из сети превратился в словарь Python, нужен посредник — ASGI-сервер.

ASGI (Asynchronous Server Gateway Interface) — стандарт взаимодействия между асинхронными веб-серверами и Python-приложениями.

В экосистеме FastAPI стандартом де-факто стал Uvicorn. Его задача — слушать сетевой сокет, парсить входящие байты по протоколу HTTP, формировать объекты запросов и передавать их в Event loop нашего приложения.

Проблема заключается в том, что Uvicorn — это один процесс Python. Из-за механизма GIL этот процесс жестко привязан к одному ядру процессора. Асинхронность позволяет этому ядру эффективно переключаться между задачами при ожидании I/O (например, пока мы ждем ответа от Redis или RabbitMQ), но вычислительная мощность для парсинга JSON, валидации Pydantic-моделей и маршрутизации ограничена ровно одним ядром.

Если мы запустим контейнер командой CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0"], мы искусственно ограничим пропускную способность API-узла.

Gunicorn как дирижер процессов

Чтобы задействовать все ядра, нам нужно запустить несколько экземпляров Uvicorn. Можно было бы запустить несколько Docker-контейнеров на одной машине, но это создает избыточный оверхед. Правильный паттерн — использование менеджера процессов.

В мире Python эту роль исторически выполняет Gunicorn. Хотя изначально он создавался для синхронных WSGI-приложений, его архитектура позволяет использовать сторонние классы воркеров.

Мы назначаем Gunicorn на роль Process Manager (управляющего процессами). Он запускает один главный процесс (Master), который:

  1. Открывает сетевой сокет на порту (например, 8000).
  2. Порождает несколько дочерних рабочих процессов (Workers).
  3. Следит за их здоровьем и автоматически перезапускает упавшие воркеры (например, если процесс убит OOM Killer-ом).

В качестве самих воркеров мы указываем класс uvicorn.workers.UvicornWorker. Таким образом, Gunicorn управляет жизненным циклом, а Uvicorn внутри каждого воркера крутит асинхронный Event loop.

Сколько воркеров нужно?

Для расчета оптимального количества рабочих процессов используется устоявшаяся эмпирическая формула:

W=2C+1W = 2C + 1

Где:

  • WW — итоговое количество воркеров.
  • CC — количество доступных ядер CPU.

Пример: Если наш API-узел развернут на виртуальной машине с 4 ядрами, нам потребуется W=2×4+1=9W = 2 \times 4 + 1 = 9 воркеров.

Почему именно так? Поскольку наши API-узлы в MLOps архитектуре занимаются I/O-bound задачами (валидируют тензоры и отправляют их в RabbitMQ, не выполняя тяжелую математику напрямую), воркеры часто простаивают в ожидании сети. Множитель «2» гарантирует, что пока часть воркеров заблокирована сетевыми вызовами, другие готовы принимать новые HTTP-запросы. Единица добавляется для того, чтобы всегда оставался один резервный процесс для обработки пиковых всплесков, пока остальные заняты.

В Dockerfile из предыдущей главы финальная команда запуска примет следующий вид:

CMD ["gunicorn", "app.main:app", "--workers", "9", "--worker-class", "uvicorn.workers.UvicornWorker", "--bind", "0.0.0.0:8000"]

Балансировка нагрузки: от узла к кластеру

Gunicorn решает проблему масштабирования в рамках одной машины. Но что, если 9 воркеров перестанут справляться? Мы начинаем добавлять новые серверы (горизонтальное масштабирование).

Теперь у нас есть три сервера по 9 воркеров. Клиент не может знать IP-адреса всех серверов. Перед ними встает Load Balancer (балансировщик нагрузки) — например, Nginx, HAProxy или Ingress-контроллер в Kubernetes.

В MLOps критически важно правильно выбрать алгоритм балансировки.

Алгоритм Как работает Применимость в MLOps
Round Robin Распределяет запросы строго по кругу: 1-му серверу, 2-му, 3-му, затем снова 1-му. Плохо подходит для генеративных моделей. Если 1-й сервер получил запрос на генерацию 2000 токенов, а 2-й — на 10 токенов, 2-й освободится моментально. Round Robin всё равно отправит следующий запрос 1-му серверу, перегружая его.
Least Connections Отправляет запрос тому серверу, у которого в данный момент меньше всего активных открытых соединений. Идеально. Особенно при использовании StreamingResponse (SSE), где соединение удерживается открытым всё время генерации. Балансировщик видит, какой узел сейчас занят долгими генерациями, и направляет трафик на свободные.

Финальная архитектура масштабируемого инференса

Пройдя путь от первого скрипта с lifespan до кластерной балансировки, мы построили полноценную промышленную архитектуру. Давайте проследим путь одного сложного запроса через всю систему, чтобы увидеть, как взаимодействуют все изученные концепции.

  1. Вход: Клиент отправляет тяжелый тензор в формате Base64. Запрос попадает на Load Balancer кластера.
  2. Балансировка: Балансировщик по алгоритму Least Connections находит наименее загруженный контейнер API.
  3. Маршрутизация внутри узла: Запрос попадает в Gunicorn Master, который передает его одному из свободных Uvicorn-воркеров.
  4. Валидация (Fail-fast): Pydantic v2 мгновенно проверяет размерность тензора. Если она неверна, клиент сразу получает 422 ошибку, не нагружая кластер.
  5. Асинхронность: FastAPI-эндпоинт формирует задачу и асинхронно отправляет её в RabbitMQ. Event loop не блокируется, воркер готов принять следующий запрос.
  6. Очередь: RabbitMQ надежно хранит задачу. Если один из GPU-серверов падает с CUDA OOM, механизм ACK возвращает задачу в очередь.
  7. Инференс: Celery-воркер на GPU-узле берет задачу, извлекает модель из памяти (загруженную ранее через lifespan) и выполняет вычисления.
  8. Результат: Итог сохраняется в Redis с заданным TTL, а API-узел отдает его клиенту (либо через Polling, либо потоком через Server-Sent Events).
  9. Наблюдаемость: Весь этот путь покрыт единым Trace ID в OpenTelemetry, а метрики (TTFT, TPOT) собираются Prometheus-ом.

Именно так stateless веб-серверы примиряются со stateful ML-моделями. Разделение слоев, правильная работа с Event loop, надежные очереди сообщений и многоуровневое управление процессами — это фундамент, который позволяет системе выдерживать любые нагрузки, оставаясь предсказуемой и отказоустойчивой.