Выделенный дренер (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,
на который есть ссылка выше.