Рост throughput относительно количества workers
Рост throughput относительно количества workers

Продолжение прошлого поста

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

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

Я уже думал писать статью, но решил мимолетом запустить 175 воркеров. Это решение создало неожиданную проблему — задача оказалась обработана за 11 секунд до ее создания.

Переработка логирования

Если раньше под логирование работала отдельная функция:

async def check_lag_monitor(interval=0.1):
    while True:
        start = time.monotonic()
        await asyncio.sleep(interval)

        actual = time.monotonic() - start

        if actual > interval * 1.5:
            worker_logger.warn(f"Event loop lag detected: expected {interval}s, got {actual:.3f}s")

То сейчас, я решил полностью переделать логику EventLoopMonitor, сделал отдельный класс:

import asyncio
import time
from core.utils import logger as worker_logger

class EventLoopMonitor:
    def __init__(self, interval=0.05, threshold_multiplier=1.5):
        self.interval = interval
        self.threshold_multiplier = threshold_multiplier
        self.lags = []

    async def run(self):
        while True:
            start = time.monotonic()
            await asyncio.sleep(self.interval)
            actual = time.monotonic() - start
            if actual > self.interval * self.threshold_multiplier:
                worker_logger.warn(f"Event loop lag detected: expected {self.interval}s, got {actual:.3f}s")
                self.lags.append(actual - self.interval)

    def get_stats(self):
        return {
            "count": len(self.lags),
            "max": max(self.lags, default=0),
            "sum": sum(self.lags),
            "avg": sum(self.lags) / len(self.lags) if self.lags else 0,
        }

Причиной стало то, что мне нужны были конкретные данные по воркерам, а не просто warn. С классом я смог получать более общие и удобные для анализа данные.

И затем уже использую его там где это нужно через app.state.event_loop = event_loop
Если кратко, в классе происходит создание монитора, который работает как отдельная корутина в том же процесса, из которой можно спокойно получать json‑формат статистики после теста.

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

@router.get("/debug/event_loop_stats", status_code=200)
async def get_event_loop_stats(request: Request):
    event_loop = request.app.state.event_loop
    return {
        "interval": event_loop.interval,
        "threshold_multiplier": event_loop.threshold_multiplier,
        "lags": event_loop.get_stats()
    }

Он в свою очередь является обычным эндпоинтом, поэтому через curl или через веб‑интерфейс я могу считывать статистику после тестирования воркеров.

Что конкретно изменилось в вызовах

Раннее логирование воркеров производилось напрямую и блокировало весь event loop из‑за того что было синхронным и весь поток ждал, пока вывод завершится:

worker_logger.info(f"Got task: {task.name}\nWith id: {task.id}\nWorker id: {self.worker_id}")

Сейчас, все происходит интереснее. Для начала получаем текущий event loop:

loop = asyncio.get_running_loop()

А затем, где это необходимо, запускаем логирование через loop.run_in_executor(), что позволяет выполнить обычную синхронную функцию вне event loop, чтобы она не блокировала его.

И соответственно получаем вызов функций формата:

await loop.run_in_executor(None, worker_logger.info, f"Got task: {task.name}\nWith id: {task.id}\nWorker id: {self.worker_id}")

Накопление задач в БД

После этих изменений, я ожидал, что throughput вырастет до небес, но этого не произошло. При тестировании throughput только падал. И я думал, что же не так, пока не решил воспользоваться давно известным EXPLAIN (ANALYZE, BUFFERS):

Seq Scan on tasks
  Filter: (status = 'PENDING'::status)
  Rows Removed by Filter: 2450

Оказалось, что я просто не очищал tasks после прогонов. Решение оказалось достаточно простым конкретно для этого случая — добавить clear_tasks в начало:

async def clear_tasks(client: httpx.AsyncClient) -> None:
    response = await client.delete("/tasks/")
    response.raise_for_status()

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

Интересная деталь с docker

Кроме того, я попал в ловушку своей невнимательности. Ввел как обычно: docker-compose run loadtest python generate.py И... ошибка SyntaxError. Оказалось, что во‑первых, я попал в интерактивную консоль Python своего предыдущего теста, во‑вторых, все это время я писал без флага --rm, что генерировало и не удаляло отработавшие контейнеры. Заглянув в Docker Desktop я обнаружил 11 «забытых» контейнеров, которые висели в системе. С этого момента, я навсегда запомню про такой флаг.

Итоговая таблица

Стена до изменений
Стена до изменений
Стена после изменений
Стена после изменений

P. S. 8,16 и 32 воркера тестировались по 1 разу, соответственно медианна может отличаться от этих значений.

Итоги изменений

Исходя из этих итогов, я могу сказать, что стена действительно сдвинулась благодаря изменению самого логирования, однако, при тестировании изменений я заметил баг. При указывании воркеров в размере явно более 64 и не попадающему линейному развитию, queue latency становится отрицательным. Например, среднее время выполнение ожидания в очереди на 175 воркерах составило -1.1 секунды, а пиковое значение вообще около -11 секунд.
Это явный баг, который я пока не смог объяснить полностью, поэтому, если у вас есть мысли, можете написать в комментарии.