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

192 уроков, 14 библиотек и челлендж «Что выведет код?» — бесплатно, код прямо в браузере

Начать обучение
Урок 14 из 20 Средний 35 мин 140 XP

Semaphore: когда конкурентов должно быть пять, а не пятьдесят

Тысячу задач gather запустит одновременно — и завалит API. Semaphore строит окно: максимум N конкурентов, остальные стоят в очереди. Считаем пики и подбираем размер.

Редакция Питоники

К этому моменту вы умеете запускать задачи пачками: gather честно стартует все переданные корутины разом. На трёх задачах это подарок, на трёх тысячах — диверсия: тысяча одновременных соединений выжжет лимиты API, надёжный сервер ответит отказом, а часть запросов утонет в таймаутах. Конкурентность — ресурс, который тоже надо нормировать. Инструмент — Semaphore, счётчик одновременных входов.

Аналогия из жизни — проходная с турникетами: сколько бы людей ни стояло в очереди, внутрь одновременно проходят столько, сколько турникетов. Никто не разворачивается, никто не теряется — просто темп задаёт пропускной узел, а не толпа. Семафор в asyncio — цифровой турникет: N разрешений, каждое выдают одному входу и забирают на выходе. И как с проходной, размер подбирается не по желанию толпы, а по вместимости помещения — то есть по выносливости той стороны, куда вы идёте.

Сначала посмотрим на лавину

Десять «запросов» с сетевой задержкой 2 секунды без всяких ограничений. Счётчик active считает вошедших, peak запоминает максимум — так мы увидим честный пик одновременности:

gather без ограничений: пик 10 из 10
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())
Нужен aiohttp и сеть. Семафор живёт рядом с сессией и передаётся каждой корутине: gather запланировал все 20, но в сеть одновременно уходят максимум 5. Тот же приём внутри хендлеров бота спасает от лавины апдейтов.

Что дальше

Вы собрали полный арсенал регулировки: 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())
Проверь себя
0 / 5

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?

Карточки терминов
Запомнено: 0 / 6
Практика

Укроти пачку: шесть задач, каждая печатает «задача N начала» и «задача N завершилась» вокруг двухсекундного «запроса». Одновременно внутри может работать максимум две — оборачивай тело задачи в async with семафора. В конце выведи сводку: пик одновременности и общее время в секундах.

practice.py
Вопросы и ответы по уроку

Как ограничить количество одновременных запросов в 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-канал или свой блог — так о проекте узнают новые читатели.

TelegramVK

Похожие уроки по темам

Подобраны автоматически по пересечению тем и ключевых слов.