Перейти к содержанию

Выделенный дренер (Dedicated Drainer)

Бенчмарки опционального режима dedicated_drainer: отдельный фоновый поток берёт на себя чтение сокета, чтобы паблишеры (потоки, которые публикуют сообщения) не выстраивались в очередь за чужим блокирующим вызовом drain_events().

Локальный прогон, а не гарантированный порог

Все цифры на этой странице получены за один локальный прогон (macOS, RabbitMQ в Docker за Toxiproxy) — это ориентир, а не SLA. Там, где разброс между запусками существенный, это отдельно оговорено.

Что это такое

По умолчанию соединение читает кадры (frame — единица данных протокола AMQP) из сокета лениво: чтением занимается тот прикладной поток, который в данный момент вызвал drain_events(). Этот поток держит блокировку транспорта (transport lock) на всё время блокирующего чтения из сокета — до истечения собственного таймаута вызова, — и любой другой поток, использующий то же соединение (в первую очередь паблишеры), встаёт в очередь за ним.

Опция transport_options={"dedicated_drainer": True} (по умолчанию False) даёт соединению собственный фоновый поток, который забирает на себя чтение сокета. Он опрашивает сокет через select, перекладывает все доступные кадры в буферы соединения и берёт блокировку транспорта только на уже полностью готовые кадры — никогда на блокирующее чтение. Остальные потоки просто ждут в wait_activity(), а не соревнуются за сокет самостоятельно. Полный контракт (жизненный цикл, heartbeat, семантика закрытия) описан в разделе README — Dedicated drainer thread.

Бенчмарк 1: задержка публикации при зависшем консьюмере

Фоновый поток крутится в блокирующем цикле drain_events() на том же соединении, которым пользуется паблишер, — это моделирует простаивающего консьюмера (получателя сообщений), делящего соединение с паблишером. В легаси-режиме каждое такое блокирующее чтение держит блокировку транспорта на весь свой таймаут, и публикация может застрять у него за спиной.

Режим P50 P95 P99 Max
Легаси 0,04 мс 1001 мс 1956 мс 1956 мс
dedicated_drainer 0,02 мс 0,08 мс 2,1 мс 2,1 мс

В этом прогоне P99 (99-й перцентиль — значение, хуже которого попадает только 1% измерений) улучшился примерно в 900 раз. От сессии к сессии это соотношение колебалось от 380 до 1566 раз: и числитель, и знаменатель здесь измеряются в единицах миллисекунд, поэтому само соотношение шумное, хотя направление изменения стабильно. Бенчмарк проверяет один жёсткий порог для обоих режимов: P99 публикации не должен превышать 5 секунд; в остальном он не задуман как строгий детектор регрессий.

Воспроизвести: pytest tests/benchmarks/test_drainer_benchmarks.py -v -m benchmark

Бенчмарк 2: паритет по пропускной способности

Тот же сценарий публикации сообщений, что и в разделе Throughput, прогнанный подряд в легаси-режиме и с dedicated_drainer.

Потоков Режим Пропускная способность P99
10 Легаси 28 202 msg/s 3,3 мс
10 dedicated_drainer 36 626 msg/s 1,5 мс
100 Легаси 24 681 msg/s 7,0 мс
100 dedicated_drainer 23 756 msg/s 8,9 мс

При умеренном числе потоков dedicated_drainer быстрее и по пропускной способности, и по хвостовой задержке — здесь нет простаивающего консьюмера, но дренер всё равно держит блокировку транспорта меньше времени на кадр, чем блокирующее чтение. На 100 потоках оба режима показывают результат в пределах шума друг относительно друга: при таком числе потоков узким местом становится конкуренция за блокировку и встречное давление брокера (backpressure), а на это ни один из режимов не влияет. Опция не жертвует пропускной способностью ради решения проблемы зависшего консьюмера из бенчмарка выше.

Воспроизвести: pytest tests/benchmarks/bench_throughput.py -v -m benchmark

Бенчмарк 3: восстановление после сброса соединения

Симулирует падение RabbitMQ (TCP RST через toxic reset_peer в Toxiproxy) при нескольких рабочих потоках, разделяющих одно соединение, и измеряет, сколько времени требуется, чтобы ошибка дошла до всех потоков.

Потоков Режим Время распространения ошибки
10 Легаси 468 мс
10 dedicated_drainer 2,2 мс
50 Легаси 523 мс
50 dedicated_drainer 4,1 мс

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

Воспроизвести: pytest tests/benchmarks/bench_recovery_latency.py -v -k reset_peer -m benchmark

Бенчмарк 4: ограниченная по времени остановка под нагрузкой

Вызов close() на соединении с одновременно активными фоновым консьюмером и паблишером, 10 свежих соединений на каждый режим.

Режим Завершилось Дедлоков P50 P95
Легаси 10/10 0/10 1952 мс 1990 мс
dedicated_drainer 10/10 0/10 508 мс 513 мс

Раньше легаси-режим здесь зависал — исправлено в v0.7.1

До версии v0.7.0 включительно close() в легаси-режиме держал блокировку транспорта на всё время ожидания ответа CloseOk от брокера. Если в этот самый момент у фонового консьюмера было в разгаре собственное блокирующее чтение, он ждал ту же самую блокировку, чтобы это чтение выполнить, — взаимное ожидание по кругу (классический дедлок), подтверждённое дампами стека потоков прямо в этом бенчмарке: до исправления легаси зависал в 10 прогонах из 10. Исправление повторяет то, что dedicated_drainer делал всегда: close() больше не держит блокировку транспорта во время этого ожидания, и поток, читающий сокет, свободно доставляет CloseOk. Легаси по-прежнему останавливается медленнее (вызвавший поток ждёт своей очереди за блокирующим чтением консьюмера), но время остановки теперь ограничено.

Бенчмарк проверяет нулевое число дедлоков для обоих режимов, а для dedicated_drainer — ещё и P95 меньше 4 секунд.

Воспроизвести: pytest tests/benchmarks/test_drainer_shutdown.py -v -m benchmark

Фрагментированные кадры

Отдельная проверка, не сравнение легаси и дренера: что происходит, если тело одного AMQP-кадра (близкое к frame_max, 131 072 байта) приходит не одним чтением, а множеством мелких порций с задержкой — это смоделировано парой toxic'ов Toxiproxy: bandwidth (16 КБ/с) и slicer (порции по ~512 байт) на одном прокси.

Сообщение доставляется целиком в каждом запуске. Пока оно приходит, загрузка CPU потоком дренера остаётся низкой: доля занятости (busy ratio) одного ядра — 0,003–0,005 за всё окно доставки (измерено через psutil). Раньше, до исправления в обработке таймаута чтения посреди кадра, тот же сценарий забирал почти целое ядро (busy ratio ≈ 0,97) — дренер крутился в плотном цикле в ожидании остатка кадра вместо того, чтобы уступать управление между порциями.

Воспроизвести: pytest tests/benchmarks/test_drainer_fragmentation.py -v -m benchmark

Heartbeat и намеренное закрытие

При включённой опции соединение само отправляет heartbeat-кадры (служебные "пульсы", которыми клиент и брокер подтверждают друг другу, что соединение живо) по расписанию — этим занимается тот же цикл дренера, отдельный вызов heartbeat_check не нужен.

Учтите: это один дополнительный фоновый поток на соединение, живущий всё время его жизни. Он стартует в момент реального подключения (для соединения, которое только публикует, — при первой публикации: kombu подключается лениво) и останавливается вместе с close()/collect() — вручную управлять им не нужно. Для «только публикующих» соединений это и есть главное практическое отличие: в легаси-режиме между публикациями сокет никто не читает, heartbeat в паузах не отправляется и не проверяется, а закрытие соединения со стороны брокера остаётся незамеченным до следующей неудачной публикации; с дренером соединение поддерживает себя само.

Намеренный вызов close()/collect() будит любой поток, застрявший в drain_events(), исключением ConnectionClosedIntentionally — благодаря этому механизм повторных попыток kombu (ensure()/Consumer) не принимает осознанное закрытие соединения за сбой, после которого нужно переподключаться. Подробности по обоим пунктам — в разделе README, на который есть ссылка выше.