Я думал, что 16 воркеров ускорят обработку задач в 16 раз. Но что‑то пошло не так

в 14:29, , рубрики: backend, celery, docker, optimisation, python, worker, асинхронность, многопоточность, распределенные системы, тестирование

Недавно я написал свою систему распределённой обработки задач. Сначала это был учебный проект. Мне хотелось лучше разобраться в очередях задач, координации воркеров, блокировках строк в PostgreSQL, retry‑механизмах и планировании фоновых задач. Когда основная функциональность была готова, я понял, что не знаю, как система ведёт себя под нагрузкой.
Поэтому вместо того, чтобы добавлять очередную фичу, я решил провести простой эксперимент. Взять одинаковый набор задач и постепенно увеличивать количество воркеров. Первые результаты выглядели именно так, как я ожидал, а потом масштабирование внезапно остановилось.

Увеличение количества воркеров с 16 до 32 практически не дало прироста производительности, а при 64 воркерах ситуация стала ещё интереснее: задачи, которые должны были выполняться около 0.5 секунды, начали занимать больше секунды. Первой моей мыслью было, что я упёрся в PostgreSQL. Оказалось, что нет.

Что именно я тестировал

Система состоит из API, PostgreSQL и набора асинхронных воркеров. Воркеры получают задачи из PostgreSQL, используя SELECT ... FOR UPDATE SKIP LOCKED, выполняют их и обновляют состояние задачи.

Для эксперимента я написал небольшой load generator.
Во всех тестах использовались одинаковые условия:

  • 350 задач;

  • задачи отправлялись одновременно;

  • каждая задача выполняла await asyncio.sleep(0.5);

  • менялось только количество воркеров;

  • остальные параметры оставались неизменными.

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

Тестирование проводилось локально на ноутбуке:

  • Intel Core i7-1065G7;

  • Intel Iris Plus;

  • 16 GB RAM.

Поэтому абсолютные значения throughput не стоит воспринимать как показатель производительности production‑сервера. В этом эксперименте меня интересовало прежде всего относительное масштабирование системы.


От одного воркера к шестнадцати

Я начал с одного воркера и постепенно увеличивал их количество:

1 -> 2 -> 4 -> 8 -> 16

Результаты оказались довольно приятными.
При увеличении количества воркеров throughput рос почти пропорционально:

Воркеры

Throughput, tasks/s

1

1.85

2

3.64

4

6.91

8

13.18

16

23.78

В диапазоне от 1 до 16 воркеров каждое удвоение количества воркеров давало примерно пропорциональный прирост — от 1.97x до 1.80x.

Одновременно резко сокращалась задержка очереди:

  • 1 worker — около 92.5 секунд;

  • 16 workers — около 6.7 секунды.

На этом этапе всё выглядело довольно предсказуемо.

Больше воркеров, значит больше одновременно выполняющихся задач, соответственно выше throughput и меньше времени ожидания в очереди.

Поэтому следующим логичным шагом было проверить, что произойдёт после 16 воркеров.


А потом масштабирование остановилось

Я повторил тот же эксперимент для 32 и 64 воркеров.
Количество задач осталось тем же — 350. Искусственная задержка тоже осталась 0.5 секунды.
Изменилось только количество воркеров. И здесь график внезапно изменился.
Переход: 16 -> 32 workers дал всего около 3% прироста throughput.
Увеличение: 32 → 64 workersтоже практически ничего не изменило.

Получилась довольно странная картина: количество воркеров увеличивается в четыре раза, а производительность почти не меняется. При этом latency очереди продолжала уменьшаться.

То есть дополнительные воркеры действительно забирали задачи из очереди быстрее, но это почти не превращалось в увеличение общей пропускной способности системы.
Но ещё интереснее было посмотреть не на очередь, а непосредственно на время выполнения задачи.


Почему полсекунды превратились в 1.22?

Каждая тестовая задача делала практически следующее:

await asyncio.sleep(0.5)

До 16 воркеров медианное время обработки оставалось примерно на уровне ожидаемых 0.52 секунды. Но дальше произошло следующее:

Воркеры

Медианное время обработки

16

~0.52 s

32

~0.79 s

64

~1.22 s

При 64 воркерах задача, которая должна была ждать около 0.5 секунды, в среднем стала занимать больше 1.2 секунды. Это было уже сложно объяснить простой конкуренцией за строки в базе. И в итоге я начал искать bottleneck.


Первое предположение: PostgreSQL

Первой моей гипотезой был PostgreSQL.

Первое что я подумал — PostgreSQL. Может он был проблемой, потому что сама суть работы воркеров была основана на конструкции:

SELECT ...
FOR UPDATE SKIP LOCKED

Поэтому вполне логично было предположить, что при большом количестве воркеров база становится узким местом. Именно это я и ожидал увидеть.
Поэтому первым делом посмотрел на загрузку ресурсов.
Во время запуска с 64 воркерами docker stats показал примерно:

  • application container ~47% CPU;

  • PostgreSQL container ~13% CPU.

Это уже было неожиданно, ведь PostgreSQL явно не выглядел перегруженным. Да и само приложение не использовало CPU полностью. Значит, проблема могла находиться где‑то между этими двумя очевидными вариантами.


Что происходит с event loop?

В этот момент я решил перестать гадать и попробовать измерить то, что непосредственно отвечает за выполнение всех этих корутин — я написал небольшой монитор задержки event loop.

Идея очень простая: запускается корутина, которая каждые 50 миллисекунд делает asyncio.sleep(), а затем сравнивает ожидаемое и фактическое время ожидания.

async def check_lag_monitor(interval=0.05):
    while True:
        start = time.monotonic()

        await asyncio.sleep(interval)

        actual = time.monotonic() - start

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

В идеальном случае:

expected: 0.050s
actual:   ~0.050s

Если event loop занят и не может вовремя продолжить выполнение coroutine:

expected: 0.050s
actual:   0.100s

Такой монитор не показывает непосредственно причину задержки, но позволяет увидеть сам факт того, что event loop не успевает выполнять задачи в ожидаемое время.

И именно это я и хотел проверить.


Что показал монитор

Во время самого загруженного участка benchmark'а monitor начал практически постоянно фиксировать задержки.

Причём предупреждения хорошо совпадали с периодами интенсивного логирования:

task received
task completed
task received
task completed
...

И тут я заметил одну неприятную деталь. Логирование в этом месте выполнялось синхронно. То есть coroutine могла сделать что‑то вроде:

logger.info("task completed")

и запись в файл выполнялась непосредственно в том же event loop, где находились остальные coroutine. Каждая отдельная операция была дешёвой.

Но при большом количестве событий их накопленный эффект становился заметным. Это объясняло часть проблемы.


Но логирование объясняло не всё

Я мог бы на этом остановиться и сказать:

«Проблема была в синхронном логировании».

Но это было бы слишком сильным выводом. После того как основная очередь задач опустошалась и интенсивность логирования падала, event loop monitor всё ещё иногда фиксировал задержки. Они становились реже, но полностью не исчезали.

А в процессе выполнения у меня одновременно работали:

  • десятки worker coroutine;

  • scheduler;

  • reaper;

  • heartbeat;

  • другие фоновые coroutine.

Все они находились в одном event loop. Поэтому появилась ещё одна гипотеза: при таком количестве конкурентных coroutine начинает становиться заметным сам overhead координации и планирования.

Важно отметить: я не измерил этот overhead напрямую.

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

  1. синхронное I/O при логировании;

  2. overhead выполнения и координации большого количества coroutine в одном event loop.


Почему дополнительные coroutine не бесплатны

Здесь есть довольно простой, но важный момент. Легко представить asyncio примерно так:

«Если coroutine большую часть времени ждёт, значит она почти ничего не стоит. Можно просто создать ещё несколько десятков worker'ов».

До определённого момента это действительно работает. Если задача большую часть времени проводит в ожидании I/O, добавление конкурентных worker'ов позволяет эффективнее использовать время ожидания. Но coroutine не исчезает полностью после await.

Event loop должен:

  • отслеживать готовность coroutine;

  • возобновлять её выполнение;

  • обрабатывать callbacks;

  • переключаться между большим количеством задач;

  • выполнять таймеры;

  • обрабатывать результаты I/O;

  • запускать другие фоновые задачи.

При небольшом количестве coroutine этот overhead практически незаметен. Но при увеличении их количества стоимость координации тоже растёт. В какой‑то момент дополнительные worker'ы перестают давать тот же эффект, потому что всё больше времени уходит не на полезную работу, а на обслуживание самой конкурентности. В моём случае это особенно хорошо видно по времени выполнения. Сама операция ничего тяжёлого не делает. Если она начинает занимать 0.79 или 1.22 секунды, значит проблема находится не внутри этой операции, а где‑то вокруг неё — в том числе в том, насколько быстро event loop может возобновить coroutine после истечения таймера.


И ещё одна неожиданная проблема — логирование

Вторая вещь, которую я забрал из этого эксперимента, оказалась даже более практичной. Синхронное I/O может выглядеть совершенно безобидно пока система работает под небольшой нагрузкой, запись нескольких сообщений в лог не вызывает никаких заметных проблем. Но когда количество событий возрастает в десятки раз, эти маленькие блокирующие операции начинают конкурировать с основной работой event loop.

Именно поэтому проблема не проявлялась в обычном запуске системы и стала очевидной только после того, как я начал измерять event loop latency.


Что я бы попробовал дальше

У этого эксперимента осталось несколько логичных продолжений.

1. Разнести компоненты по процессам

Сейчас API, worker'ы и фоновые сервисы работают в рамках одной Python‑процесса/event loop. Следующий эксперимент — разделить их на несколько процессов. Это позволит операционной системе распределять CPU‑нагрузку между ядрами и одновременно убрать часть конкуренции за один event loop. При этом интересно проверить, насколько далеко после этого сдвинется scaling wall.


2. Убрать синхронное логирование

Второй очевидный эксперимент — убрать блокирующую запись логов из основного event loop. Например, можно отправлять сообщения в очередь, а запись выполнять отдельным writer task или отдельным потоком/процессом. Тогда worker не должен ждать завершения операции записи. После этого можно повторить тот же benchmark и сравнить event loop lag.


Что я вынес из этого эксперимента

Самый интересный результат оказался не в том, что система перестала масштабироваться после 16 воркеров. Гораздо интереснее оказалось то, где я сначала ожидал увидеть bottleneck и где он оказался на самом деле.
Я предполагал, что первым ограничением станет PostgreSQL из‑за конкуренции воркеров за задачи. Но измерения показали, что PostgreSQL использовал всего около 13% CPU во время проблемного запуска.

Если бы я начал оптимизировать SQL, индексы или SKIP LOCKED, я, скорее всего, потратил бы время на компонент, который вообще не был главным ограничением. Вместо этого простой монитор event loop показал, что проблема находится значительно ближе к самому приложению.
В моём случае несколько строк дополнительной instrumentation оказались полезнее, чем очередная попытка оптимизировать запрос к базе данных. И, возможно, это одна из главных вещей, которую я понял после этого эксперимента:

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


Исходный проект

Исходный код системы доступен на GitHub:

https://github.com/Shjryoku/Distributed‑Task‑and‑Job‑Processing

Автор: Balaklav

Источник

* - обязательные к заполнению поля


https://ajax.googleapis.com/ajax/libs/jquery/3.4.1/jquery.min.js