
Продолжение прошлого поста
После прошлого поста я вдохновился на продолжение, помимо того, что я изначально хотел его доработать, я также увидел, что количество людей увидевших мой пост перевалило за 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 секунд.
Это явный баг, который я пока не смог объяснить полностью, поэтому, если у вас есть мысли, можете написать в комментарии.

