Координация корутин: gather и семафоры
Запустил 200 обращений к чужому сервису одновременно. Пик одновременных вызовов - 200, время 53 миллисекунды. Красиво ровно до того момента, как этот сервис ляжет от такой заботы. Тот же набор с ограничителем на 10 одновременных: пик 10, время 1011 миллисекунд. Медленнее в двадцать раз - и это осознанная плата за то, что сосед остался жив.
Стержень: асинхронность снимает предел «сколько потоков», поэтому предел приходится ставить руками - иначе вы просто быстрее убьёте того, к кому ходите.
// Формулировки: «как ограничить одновременные запросы?», «зачем очередь с пределом?», «бывают ли гонки в async?»
Ограничитель одновременности
В синхронном коде число одновременных обращений ограничено само собой - числом потоков в пуле. В асинхронном такого предела нет: можно запустить тысячу заданий, и все они одновременно постучатся наружу. Мой замер это буквально показывает: пик 200 одновременных вызовов из 200 запущенных.
Ограничитель - это счётчик пропусков: задание берёт пропуск перед работой и возвращает после. Нет свободного - ждёт. С ограничителем на 10 пик стал ровно 10, а общее время выросло с 53 до 1011 миллисекунд. Это и есть управляемая нагрузка на соседа.
// Правильное число подбирают не из головы, а из того, сколько сосед выдерживает, и часто согласуют явно. И держат ограничитель ОДИН на весь клиент, а не создают новый на каждый запрос - иначе он ничего не ограничивает. Рядом обычно живут предел ожидания на каждый вызов и повторы с разбросом, о которых стоит думать сразу вместе.
- ограничитель одновременности
- счётчик пропусков: больше N заданий одновременно работать не будут
- управляемая нагрузка
- осознанно выбранный предел давления на чужой сервис
Очередь как регулятор скорости
Очередь с ограниченным размером связывает быстрого производителя и медленных потребителей. Проверил на стенде: очередь на пять мест, один производитель и три потребителя, 50 сообщений - обработаны все 50 за 175 миллисекунд.
Главная ценность тут не в скорости, а в том, что происходит при переполнении: производитель ЖДЁТ, пока освободится место. Без предела он положил бы в память все 50 тысяч сообщений и упёрся бы в память процесса. Очередь с пределом - это встроенная защита от того, чтобы захлебнуться.
// Типовой рисунок такой: несколько заданий-потребителей крутятся в цикле и разбирают очередь, а производитель складывает в неё работу. Завершение делают через специальную метку в очереди («больше не будет») по числу потребителей - именно так я и останавливал стенд. И не забывают про отмену: при остановке сервиса потребителей отменяют явно, иначе они будут вечно ждать следующего сообщения.
- очередь с пределом
- производитель ждёт, когда мест нет; защита от переполнения памяти
- метка завершения
- специальное сообщение «работы больше не будет» по числу потребителей
Гонки бывают и в асинхронном коде
Распространённое заблуждение: раз всё в одном потоке, гонок нет. Проверил и опроверг. Пять заданий, каждое тысячу раз читает счётчик, потом делает await, потом записывает прочитанное плюс один. Результат: 1000 вместо 5000. Четыре тысячи инкрементов потерялись.
Причина в том, что await - это точка переключения. Между чтением и записью цикл успевает отдать поток другому заданию, и оно читает то же самое старое значение. Ровно та же механика, что и в потоках, просто точки переключения тут видны глазами - это и есть слова await.
// Тот же код под асинхронным замком дал ровно 5000. Правило простое: если между чтением и записью общего состояния есть await, нужен замок. А лучше - не держать изменяемое общее состояние вовсе: пусть каждое задание работает со своими данными, а обмен идёт через очередь.
- точка переключения
- каждый await; между ним и следующей строкой успевают поработать другие
- асинхронный замок
- пропускает в участок кода одно задание за раз
Как отвечать: «Как ограничить число одновременных запросов к внешнему сервису?»
Ограничителем одновременности - счётчиком пропусков, который держится один на весь клиент. Задание берёт пропуск перед вызовом и возвращает после, а если свободных нет - ждёт. Я мерил разницу: двести запросов без ограничения дают пик в двести одновременных вызовов и укладываются в пятьдесят миллисекунд, а с ограничителем на десять пик равен десяти, и всё занимает секунду. Медленнее в двадцать раз - и это осознанная плата за то, что чужой сервис остаётся жив. Число подбирают не из головы, а из того, сколько сосед выдерживает, и обычно согласуют явно. Рядом сразу ставлю предел ожидания на каждый вызов и повторы с разбросом, потому что без разброса все клиенты вернутся одновременно. И слежу, чтобы ограничитель создавался один раз, а не на каждый запрос: новый ограничитель в каждом вызове не ограничивает ничего.
Ответ даёт механизм, честно называет цену в цифрах и добавляет две соседние настройки. Замечание про один ограничитель на клиент - типовая ошибка, которую видно только на практике.
На чём валятся
- −Запускают неограниченное число одновременных вызовов и роняют чужой сервис.
- −Создают ограничитель внутри функции, и он ничего не ограничивает.
- −Делают очередь без предела размера и упираются в память процесса.
- −Считают, что в одном потоке гонок не бывает: между чтением и записью есть await.
- −Не отменяют задания-потребители при остановке, и сервис не завершается.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 13, остальные разбираются в тренажёре.
- Нужно скачать 1000 URL, но не больше 10 одновременно. Что взять?A)
asyncio.Semaphore(10)вокруг каждой загрузкиB)asyncio.gatherсо всеми 1000 сразуC)multiprocessing.Pool(10)D)Ставитьtime.sleepмежду запросами для троттлингапоказать ответ и разбор
+A)`asyncio.Semaphore(10)` вокруг каждой загрузки// разбор: asyncio.Semaphore(10) пропускает не более 10 корутин в критическую секцию:
async with sem: await fetch(u), остальные ждут на acquire. gather со всеми 1000 сразу откроет тысячу соединений. Пул процессов тут лишний: задача I/O-bound, не CPU. - Что делает
asyncio.wait_for(coro, timeout=2)при превышении таймаута?A)Дождётся всё равно — таймаут это советB)Блокирует цикл на 2 секундыC)Возвращает NoneD)Отменяет coro и поднимает TimeoutErrorпоказать ответ и разбор
+D)Отменяет coro и поднимает TimeoutError// разбор: wait_for по истечении таймаута отменяет внутреннюю задачу (шлёт CancelledError в точку await) и поднимает asyncio.TimeoutError. Следствие: coro должна переживать отмену — закрывать ресурсы в finally. Таймаут жёсткий, а не рекомендательный.
- Нужно обрабатывать результаты по мере готовности, не дожидаясь всех. Что взять?A)
asyncio.wait_forB)asyncio.gatherC)Последовательныйawaitкаждой задачи по очередиD)asyncio.as_completedнад задачамипоказать ответ и разбор
+D)`asyncio.as_completed` над задачами// разбор: asyncio.as_completed отдаёт итератор футур в порядке ЗАВЕРШЕНИЯ:
for fut in as_completed(tasks): x = await fut— обрабатываешь первый освободившийся, не дожидаясь медленных. gather вернёт всё разом в порядке аргументов. Для стриминга результатов as_completed (или очередь) — то, что нужно. await asyncio.gather(boom(), slow()), boom сразу кидает исключение. Что со slow?A)Исключение всплывает, но slow продолжает работатьB)gather вернёт частичный список результатовC)slow отменяется автоматическиD)slow тоже получит это исключениепоказать ответ и разбор
+A)Исключение всплывает, но slow продолжает работать// разбор: По умолчанию (return_exceptions=False) gather пробрасывает первое исключение сразу, НО не отменяет остальные — slow доработает «в фоне», и её ошибка может уйти в «never retrieved». Чтобы собрать всё не роняя — return_exceptions=True; чтобы падать и отменять сиблингов — TaskGroup (3.11+).
- Что делает
await asyncio.gather(a(), b(), c())?A)Исполняет корутины строго по очереди: сперва a, потом b, потом cB)Запускает их в трёх отдельных потоках и собирает результатыC)Запускает корутины конкурентно, результаты — в порядке аргументовD)Возвращает результат первой завершившейся корутины, отменив остальныепоказать ответ и разбор
+C)Запускает корутины конкурентно, результаты — в порядке аргументов// разбор: gather планирует все переданные корутины в цикле и ждёт завершения всех. Результаты возвращаются списком в порядке аргументов, а не в порядке готовности. Это базовый способ распараллелить независимые ожидания (несколько сетевых запросов) внутри одного цикла.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.