сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Асинхронность в Python

Координация корутин: 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, остальные разбираются в тренажёре.

  1. #async_coordination1 / 5
    Нужно скачать 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.

  2. #async_coordination2 / 5
    Что делает asyncio.wait_for(coro, timeout=2) при превышении таймаута?
    A)Дождётся всё равно — таймаут это совет
    B)Блокирует цикл на 2 секунды
    C)Возвращает None
    D)Отменяет coro и поднимает TimeoutError
    показать ответ и разбор
    +D)Отменяет coro и поднимает TimeoutError

    // разбор: wait_for по истечении таймаута отменяет внутреннюю задачу (шлёт CancelledError в точку await) и поднимает asyncio.TimeoutError. Следствие: coro должна переживать отмену — закрывать ресурсы в finally. Таймаут жёсткий, а не рекомендательный.

  3. #async_coordination3 / 5
    Нужно обрабатывать результаты по мере готовности, не дожидаясь всех. Что взять?
    A)asyncio.wait_for
    B)asyncio.gather
    C)Последовательный await каждой задачи по очереди
    D)asyncio.as_completed над задачами
    показать ответ и разбор
    +D)`asyncio.as_completed` над задачами

    // разбор: asyncio.as_completed отдаёт итератор футур в порядке ЗАВЕРШЕНИЯ: for fut in as_completed(tasks): x = await fut — обрабатываешь первый освободившийся, не дожидаясь медленных. gather вернёт всё разом в порядке аргументов. Для стриминга результатов as_completed (или очередь) — то, что нужно.

  4. #async_coordination4 / 5
    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+).

  5. #async_coordination5 / 5
    Что делает await asyncio.gather(a(), b(), c())?
    A)Исполняет корутины строго по очереди: сперва a, потом b, потом c
    B)Запускает их в трёх отдельных потоках и собирает результаты
    C)Запускает корутины конкурентно, результаты — в порядке аргументов
    D)Возвращает результат первой завершившейся корутины, отменив остальные
    показать ответ и разбор
    +C)Запускает корутины конкурентно, результаты — в порядке аргументов

    // разбор: gather планирует все переданные корутины в цикле и ждёт завершения всех. Результаты возвращаются списком в порядке аргументов, а не в порядке готовности. Это базовый способ распараллелить независимые ожидания (несколько сетевых запросов) внутри одного цикла.

дальше

Теорию прочитали. Навык ставится повторением

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