Semaphore: когда конкурентов должно быть пять, а не пятьдесят
Тысячу задач gather запустит одновременно — и завалит API. Semaphore строит окно: максимум N конкурентов, остальные стоят в очереди. Считаем пики и подбираем размер.
Редакция Питоники
К этому моменту вы умеете запускать задачи пачками: gather честно стартует все переданные корутины разом. На трёх задачах это подарок, на трёх тысячах — диверсия: тысяча одновременных соединений выжжет лимиты API, надёжный сервер ответит отказом, а часть запросов утонет в таймаутах. Конкурентность — ресурс, который тоже надо нормировать. Инструмент — Semaphore, счётчик одновременных входов.
Аналогия из жизни — проходная с турникетами: сколько бы людей ни стояло в очереди, внутрь одновременно проходят столько, сколько турникетов. Никто не разворачивается, никто не теряется — просто темп задаёт пропускной узел, а не толпа. Семафор в asyncio — цифровой турникет: N разрешений, каждое выдают одному входу и забирают на выходе. И как с проходной, размер подбирается не по желанию толпы, а по вместимости помещения — то есть по выносливости той стороны, куда вы идёте.
Сначала посмотрим на лавину
Десять «запросов» с сетевой задержкой 2 секунды без всяких ограничений. Счётчик active считает вошедших, peak запоминает максимум — так мы увидим честный пик одновременности:
import asyncio
active = 0
peak = 0
async def request(n):
global active, peak
active += 1
peak = max(peak, active)
print("запрос", n, "стартовал, одновременно:", active)
await asyncio.sleep(2) # имитация сетевой задержки
active -= 1
async def main():
await asyncio.gather(*(request(i) for i in range(1, 11)))
print("Пик одновременных запросов:", peak)
asyncio.run(main())
запрос 1 стартовал, одновременно: 1 запрос 2 стартовал, одновременно: 2 запрос 3 стартовал, одновременно: 3 запрос 4 стартовал, одновременно: 4 запрос 5 стартовал, одновременно: 5 запрос 6 стартовал, одновременно: 6 запрос 7 стартовал, одновременно: 7 запрос 8 стартовал, одновременно: 8 запрос 9 стартовал, одновременно: 9 запрос 10 стартовал, одновременно: 10 Пик одновременных запросов: 10
Все десять стартовали в один момент: gather не знает про лимиты, он честно делает то, что вы просили. Против слабенького API это выглядело бы как DDoS: десять соединений, десять нагрузок на чужую квоту, десять шансов поймать 429 Too Many Requests. Программа отработала быстро — и именно в этом проблема: скорость здесь куплена ценой агрессии.
Что при этом чувствует сервер? У любой стороны есть свои очереди и лимиты: рабочие процессы, пул соединений к базе, квоты тарифа. Десять запросов разом — десять записей в этих очередях; если лимит исчерпан, сервер не обслуживает «помедленнее», а отвечает отказом или ставит в хвост, откуда запросы выходят с просроченными таймаутами. Получается злой парадокс: чем агрессивнее клиент, тем хуже обслуживается каждый его запрос — все вместе давят в дверь и застревают в проёме. Ретраи поверх такой лавины только усугубляют картину, поэтому первым делом нормируют одновременность, и лишь потом рассуждают о повторах.
Semaphore: окно на три
Теперь то же самое, но каждый запрос входит в свою критическую секцию через async with sem, где sem = asyncio.Semaphore(3). Внутри окна — максимум три запроса; остальные ждут у входа. Плюс замер общего времени:
import asyncio
import time
active = 0
peak = 0
async def request(n, sem):
global active, peak
async with sem: # больше трёх внутрь не пройдут
active += 1
peak = max(peak, active)
print("запрос", n, "стартовал, одновременно:", active)
await asyncio.sleep(2) # имитация сетевой задержки
active -= 1
print("запрос", n, "завершился, одновременно:", active)
async def main():
sem = asyncio.Semaphore(3)
t0 = time.monotonic()
await asyncio.gather(*(request(i, sem) for i in range(1, 11)))
print("Пик одновременных запросов:", peak)
print("Секунд всего:", round(time.monotonic() - t0))
asyncio.run(main())
запрос 1 стартовал, одновременно: 1 запрос 2 стартовал, одновременно: 2 запрос 3 стартовал, одновременно: 3 запрос 1 завершился, одновременно: 2 запрос 2 завершился, одновременно: 1 запрос 3 завершился, одновременно: 0 запрос 4 стартовал, одновременно: 1 запрос 5 стартовал, одновременно: 2 запрос 6 стартовал, одновременно: 3 запрос 4 завершился, одновременно: 2 запрос 5 завершился, одновременно: 1 запрос 6 завершился, одновременно: 0 запрос 7 стартовал, одновременно: 1 запрос 8 стартовал, одновременно: 2 запрос 9 стартовал, одновременно: 3 запрос 7 завершился, одновременно: 2 запрос 8 завершился, одновременно: 1 запрос 9 завершился, одновременно: 0 запрос 10 стартовал, одновременно: 1 запрос 10 завершился, одновременно: 0 Пик одновременных запросов: 3 Секунд всего: 8
Десять запросов вышли аккуратными волнами по три: стартовала тройка, отработала, впустила следующую. Пик одновременности — ровно 3, как и заказано, ни один запрос не проскочил четвёртым. Цена вежливости видна в последней строке: 8 секунд вместо 2 — три полных волны плюс остаток. Семафор не ускоряет работу, он её нормирует; иногда (лимитированное API) это единственный способ уложиться без банов.
Механика внутри проста: у семафора есть счётчик с тремя начальными разрешениями. async with sem при входе уменьшает счётчик — если там ноль, корутина засыпает в очереди ожидания. Выход (конец блока with, даже по исключению или отмене) увеличивает счётчик и будит следующего. По сути это очередь с пропусками — родня asyncio.Queue, только очередь не заданий, а прав на вход.
Ручной вариант полезно один раз увидеть, чтобы понимать, от чего вас защищает async with. Один-семафор превращается в замок: один вход одновременно, остальные ждут:
import asyncio
async def worker(n, sem):
await sem.acquire() # то же, что async with, но руками
print("задача", n, "вошла")
await asyncio.sleep(1)
sem.release()
print("задача", n, "вышла")
async def main():
sem = asyncio.Semaphore(1)
await asyncio.gather(worker(1, sem), worker(2, sem))
asyncio.run(main())
задача 1 вошла задача 1 вышла задача 2 вошла задача 2 вышла
Обратите внимание на порядок: «вышла» печатается раньше, чем следующая «вошла» — разрешение освобождается release-ом, и второй воркер просыпается уже после выхода первого. Для реального кода предпочтительнее async with: он делает тот же release в finally, при любых развязках.
Окно на единицу — это ещё и замок для общего ресурса. Сначала посмотрим, что бывает без него: операция «прочитал — подождал — записал» в asyncio легко теряет обновления, ведь между await выполняться может кто угодно:
import asyncio
total = 0
async def deposit(amount):
global total
current = total # прочитали общий счётчик
await asyncio.sleep(0) # точка переключения посреди операции
total = current + amount # записали устаревшее значение
async def main():
await asyncio.gather(*(deposit(1) for _ in range(5)))
print("Итого:", total)
asyncio.run(main())
Итого: 1
Пять пополнений по единице — а в итоге единица: все пять корутин прочитали ноль до первой точки переключения и записали единицу поверх друг друга. Теперь тот же код под замком из семафора на одно место:
import asyncio
total = 0
async def deposit(amount, lock):
global total
async with lock: # читай-меняй-пиши без гонок
current = total
await asyncio.sleep(0) # точка переключения посреди операции
total = current + amount
async def main():
lock = asyncio.Semaphore(1)
await asyncio.gather(*(deposit(1, lock) for _ in range(5)))
print("Итого:", total)
asyncio.run(main())
Итого: 5
Критическая секция «прочитал-изменил-записал» больше не разрывается чужим await-ом: замок держится через внутреннюю паузу, и каждая корутина видит свежее значение. Так семафор работает не только против сетевых лавин, но и против логических гонок на общих данных.
Подбор размера окна
Размер окна — параметр с последствиями, и его полезно измерить. Одна и та же пачка из десяти односекундных запросов при разных окнах:
import asyncio
import time
async def fetch(n, sem):
async with sem:
await asyncio.sleep(1)
return n
async def measure(size):
sem = asyncio.Semaphore(size)
t0 = time.monotonic()
await asyncio.gather(*(fetch(i, sem) for i in range(1, 11)))
return round(time.monotonic() - t0)
async def main():
print("Окно 10:", await measure(10), "с")
print("Окно 5:", await measure(5), "с")
print("Окно 2:", await measure(2), "с")
asyncio.run(main())
Окно 10: 1 с Окно 5: 2 с Окно 2: 5 с
Арифметика волн: время примерно равно числу задач, делённому на размер окна, умноженному на длительность одной задачи. Окно 10 — всё разом, окно 2 — пять волн по секунде. Как выбирать: смотрите на требования источника — у API обычно написано «не более N запросов в секунду» или «не более N соединений»; окно семафора ставится в эти границы с запасом. Для человеческого сайта — от единиц до десятка; для своего внутреннего сервиса — сколько он реально тянет под нагрузкой.
Финальный прогон в продакшен-стиле: двадцать вызовов, окно 4, без печатей на каждый запрос — только сводка в конце. Такой однострочный отчёт удобно вешать в логи:
import asyncio
import time
active = 0
peak = 0
async def call(n, sem):
global active, peak
async with sem:
active += 1
peak = max(peak, active)
await asyncio.sleep(1)
active -= 1
async def main():
sem = asyncio.Semaphore(4)
t0 = time.monotonic()
await asyncio.gather(*(call(i, sem) for i in range(1, 21)))
print("Вызовов: 20 | Пик:", peak, "| Секунд:", round(time.monotonic() - t0))
asyncio.run(main())
Вызовов: 20 | Пик: 4 | Секунд: 5
Пятнадцать секунд экономии против честной последовательности — и ноль превышений лимита против лавины. Это и есть рабочий диапазон семафора: между «по одному» и «все разом» лежит окно, которое подбирается измерением, а не верой.
Пара слов о сочетании с предыдущими уроками. Семафор не меняет ни порядок результатов gather (он остаётся порядком аргументов — урок 8), ни поведение при падениях (изоляция ошибок — урок 11): задачи просто дольше ждут входа. А если какой-то источник просит «притормози» ответом 429 с заголовком Retry-After, окно можно уменьшать на лету — семафор живёт в переменной, и никто не мешает пересоздавать конвейер с новым размером по сигналу от сервера. Регулировка — это цикл: замерил пик и время, подвинул окно, замерил снова.
Пиковый счётчик из этих примеров — не только учебный приём. Две переменные (active, peak) и одна строка в конце пачки — а на выходе самая честная характеристика нагрузки: сколько соединений вы реально открывали одновременно и как долго шла пачка. Заведите привычку добавлять такую сводку в любой конвейер: когда через месяц встанет вопрос «а не мы ли завалили чужой API», ответ будет лежать в логе, а не в догадках. Метрика из лога легко превращается в график — и окно семафора подгоняется уже по данным, а не по ощущениям.
Как это выглядит в боевом коде парсера — семафор делится между воркерами, каждый запрос проходит через окно (сеть в песочнице недоступна, скелет для локального запуска):
import asyncio
import aiohttp
async def fetch(session, sem, url):
async with sem:
async with session.get(url) as resp:
return url, resp.status
async def main():
sem = asyncio.Semaphore(5) # максимум 5 одновременных запросов
urls = ["https://example.com/"] * 20
async with aiohttp.ClientSession() as session:
results = await asyncio.gather(*(fetch(session, sem, u) for u in urls))
for url, status in results[:3]:
print(status, url)
print("Всего запросов:", len(results))
asyncio.run(main())
Что дальше
Вы собрали полный арсенал регулировки: gather для пачек, очередь для конвейера, семафор для окна. Дальше язык asyncio углубляется: асинхронные итераторы async for — потоки данных, которые отдают элементы по мере поступления; после них — асинхронные контекст-менеджеры async with, чью механику вы только что видели изнутри семафора. А в финальном проекте всё это соберётся в диспетчер задач с таймаутами и лимитами.
Скорость без лимита — это не скорость, а лотерея: семафор превращает лавину запросов в расписание, которое переживает прод.
Сначала предскажи ответ в голове — это главный навык программиста.
import asyncio
async def job(n, sem):
async with sem:
print(n, "вошла")
await asyncio.sleep(1)
print(n, "вышла")
async def main():
sem = asyncio.Semaphore(2)
await asyncio.gather(job(1, sem), job(2, sem), job(3, sem))
asyncio.run(main())
import asyncio
async def main():
sem = asyncio.Semaphore(1)
async with sem:
print("внутри")
async with sem:
print("снова внутри")
asyncio.run(main())
import asyncio
import time
async def fetch(sem):
async with sem:
await asyncio.sleep(1)
async def main():
sem = asyncio.Semaphore(5)
t0 = time.monotonic()
await asyncio.gather(*(fetch(sem) for _ in range(10)))
print(round(time.monotonic() - t0), "с")
asyncio.run(main())
1. Что делает async with sem при sem = asyncio.Semaphore(3), когда уже работают три корутины внутри блока?
2. Двадцать задач с семафором на 4, каждая занимает 1 секунду. Сколько примерно займёт пачка?
3. Почему ручной sem.acquire() без try/finally опаснее async with sem?
4. API в правилах разрешает не более 10 одновременных соединений. Как ограничить пачку из 500 запросов?
5. Что покажет счётчик пика, если все запросы обёрнуты в async with sem на 3?
Укроти пачку: шесть задач, каждая печатает «задача N начала» и «задача N завершилась» вокруг двухсекундного «запроса». Одновременно внутри может работать максимум две — оборачивай тело задачи в async with семафора. В конце выведи сводку: пик одновременности и общее время в секундах.
Как ограничить количество одновременных запросов в asyncio?
Семафором: создайте один asyncio.Semaphore(N) и оберните тело каждой корутины в async with sem — внутрь одновременно попадут максимум N задач, остальные подождут освобождения. Пачку планируйте обычным gather: он запланирует все задачи, а окно отрегулирует одновременность. Так один семафор превращает лавину в волны.
Как выбрать размер окна семафора?
От лимитов источника: у API в правилах обычно указаны одновременные соединения или запросы в секунду — окно ставится в эти границы с небольшим запасом. Для тяжёлого своего сервиса размер подбирается нагрузочным тестом. Ориентир по времени: пачка занимает примерно (число задач / размер окна) умножить на длительность одной задачи.
Чем Semaphore отличается от Queue?
Queue хранит сами задания и раздаёт их свободным воркерам — это конвейер. Semaphore не хранит ничего: он лишь ограничивает, сколько корутин одновременно находится в заданном блоке. Их часто используют вместе: очередь отвечает за темп и буфер, семафор — за одновременность внешних вызовов.
Что происходит с семафором при отмене задачи внутри async with?
Разрешение корректно возвращается: async with реализует release в finally, который выполняется и при CancelledError. А вот ручной вариант await sem.acquire() ... sem.release() без try/finally на отмене разрешение теряет — после нескольких таких потерь семафор блокируется навсегда, поэтому руками его пишут только с finally.
Понравился урок? Сошлитесь на него
«Скорость без лимита — это не скорость, а лотерея: семафор превращает лавину запросов в расписание, которое переживает прод.»
Скопируйте готовую ссылку в формате HTML, Markdown или чистый адрес и вставьте в статью на Habr, VC, Telegram-канал или свой блог — так о проекте узнают новые читатели.
Что читать дальше
asyncio · Урок 13
Очередь asyncio.Queue: producer-consumer
Продюсер кладёт задания, консьюмеры разбирают — классический конвейер на asyncio.Queue: put, get, task_done и join, плюс maxsize, чтобы производитель не завалил обработку.
asyncio · Урок 16
Асинхронные контекст-менеджеры: async with
Соединение надо не только открыть, но и закрыть — даже если посреди работы задачу отменили. async with берёт эту уборку на себя.
asyncio · Урок 20
Проект: диспетчер задач с таймаутами и лимитами
Финал курса: очередь заданий, семафор-лимит, wait_for-дедлайны и честный отчёт — успехи, таймауты, ошибки. Мини-версия движка настоящего бота.
Похожие уроки по темам
Подобраны автоматически по пересечению тем и ключевых слов.
aiogram · Урок 7
Асинхронность в aiogram: asyncio для бота без страха
Разбираем на живом коде, зачем aiogram асинхронный: async и await, asyncio.sleep против time.sleep и gather для параллельных чатов.
asyncio для начинающихaiogram async
aiogram · Урок 8
Middleware в aiogram: антиспам, логирование и общие данные
Пишем middleware в aiogram: антиспам-троттлинг, логирование всех апдейтов и передачу общих данных в хендлеры — механику запускаем прямо на странице.
ограничение частоты сообщений телеграмзащита телеграм бота от флуда
asyncio · Урок 17
Мини-проект: загрузка страниц без сети
Сеть необязательна: шесть «страниц», фиксированные задержки и asyncio — и ты своими глазами увидишь, как конкурентная загрузка сжимает десять секунд до трёх.
asyncio имитация запросов примерasyncio имитация запросов