← все задачи

Async Python · задача 2 из 10

Ограничить число одновременных запросов

Средний 15 минут Semaphoreобратное давлениеограничения чужого API

Условие

Тысяча URL, но чужой API разрешает не больше десяти одновременных соединений. Заберите все адреса, соблюдая ограничение.

Что требуется

  • Одновременно в полёте не более десяти запросов
  • Все тысяча адресов в итоге обработаны
  • Ограничение соблюдается и при исключениях

Пример

results = await fetch_all(urls, limit=10)
# в каждый момент времени активно не больше 10 запросов

Сначала уточните

Вопросы до кода — половина оценки. Молча начать печатать хуже, чем задать два вопроса.

  • Ограничение на соединения или на частоту (столько-то запросов в секунду)? Это разные вещи
  • Нужны ли повторы при 429 и уважение к заголовку Retry-After?
  • Сохранять ли порядок результатов или можно обрабатывать по мере готовности?
Показать решение Скрыть решение

Решение

import asyncio
import aiohttp


async def fetch_all(urls, limit=10):
    semaphore = asyncio.Semaphore(limit)

    async def fetch_one(session, url):
        async with semaphore:              # ждём свободный слот
            async with session.get(url) as response:
                response.raise_for_status()
                return await response.json()

    async with aiohttp.ClientSession() as session:
        tasks = [fetch_one(session, url) for url in urls]
        return await asyncio.gather(*tasks, return_exceptions=True)


# Альтернатива для Python 3.11+: очередь задач с фиксированным числом воркеров
async def fetch_all_workers(urls, limit=10):
    queue = asyncio.Queue()
    for index, url in enumerate(urls):
        queue.put_nowait((index, url))

    results = [None] * len(urls)

    async def worker(session):
        while not queue.empty():
            index, url = await queue.get()
            try:
                async with session.get(url) as response:
                    results[index] = await response.json()
            finally:
                queue.task_done()

    async with aiohttp.ClientSession() as session:
        await asyncio.gather(*[worker(session) for _ in range(limit)])

    return results

Почему так

Почему семафор, а не нарезка на пачки

  • Пачками по десять весь батч ждёт самый медленный запрос — девять слотов простаивают
  • Семафор освобождает слот сразу, как только запрос закончился, и в полёте всегда ровно десять
  • На хвосте распределения (один запрос из десяти медленный) разница по времени получается кратная

Почему async with semaphore, а не acquire/release

  • Ручной release легко потерять при исключении — слот утечёт, и через тысячу запросов всё встанет
  • async with освобождает слот в любом случае, включая отмену задачи
  • Ровно та же логика, что и с обычным Lock: контекстный менеджер снимает целый класс ошибок

Зачем return_exceptions=True

  • Без него первое же исключение всплывёт из gather, а остальные задачи останутся работать «в фоне» без присмотра
  • С ним ошибки приходят в списке рядом с результатами, и можно решить, что делать: повторить или пропустить
  • На тысяче запросов «упасть целиком из-за одного 500» — почти всегда неверное поведение

Что спросят дальше

  • Спросят разницу: семафор ограничивает одновременность, rate limiter — частоту; для «100 запросов в секунду» нужен второй
  • Спросят про 429: нужны повторы с паузой, иначе ограничение вы формально соблюли, а API вас всё равно забанит
  • Спросят про TaskGroup из 3.11: он отменяет остальные задачи при ошибке — это другое поведение, и его надо выбирать осознанно

Следующая задача

Таймаут и отмена — Проверяют понимание отмены: что происходит с корутиной, которую прервали на середине.