← назад к разделу

Раз asyncio выполняет код в одном потоке, зачем в нём блокировки? Потому что гонки бывают не только между потоками. Две задачи, которые по очереди проверяют условие, уходят на await и возвращаются, могут обе увидеть «записи нет» и обе её создать. Разберём, где в однопоточном коде появляется гонка, какие примитивы её закрывают и чем они отличаются от своих тёзок из threading.

Обязательно

Где гонка в одном потоке

Между двумя await код атомарен: никто другой не выполняется. Гонка возникает ровно тогда, когда между проверкой и действием есть await:

_cache: dict[str, Rate] = {}

async def get_rate(currency: str) -> Rate:
    if currency not in _cache:                       # проверка
        _cache[currency] = await fetch_rate(currency) # await: сюда войдут и вторая, и третья задача
    return _cache[currency]

Десять одновременных запросов за курсом доллара при пустом кеше сделают десять вызовов fetch_rate, потому что каждая задача проверила кеш до того, как первая его заполнила. Это та же гонка, что и в потоках, только переключение происходит в известных местах, на await. Правило чтения кода: любое «проверил, подождал, записал» подозрительно.

Lock: критическая секция с ожиданием внутри

asyncio.Lock делает блок кода, содержащий await, выполняющимся по одной задаче за раз:

_lock = asyncio.Lock()

async def get_rate(currency: str) -> Rate:
    if currency in _cache:
        return _cache[currency]
    async with _lock:
        if currency not in _cache:                     # повторная проверка под замком
            _cache[currency] = await fetch_rate(currency)
    return _cache[currency]

Первая проверка без замка нужна, чтобы горячий путь не ждал; вторая под замком закрывает гонку. Один замок на весь кеш означает, что запрос за евро ждёт запроса за доллар; для независимых ключей держат словарь замков по ключу или Future на ключ, который ждут все, кроме первого.

Отличия от threading.Lock: asyncio.Lock не потокобезопасен и работает только внутри одного цикла событий; он не реентерабельный, повторный async with _lock в той же задаче заблокирует её навсегда; захват через await lock.acquire() или async with, обычный with не подходит. И главное: замок нужен только вокруг блока с await. Обычное присваивание в словарь защищать не нужно, оно и так атомарно относительно других задач.

Semaphore: не взаимное исключение, а лимит

asyncio.Semaphore(n) пускает внутрь не больше n задач одновременно. Это основной инструмент ограничения параллелизма на внешние системы: не больше двадцати одновременных запросов к соседнему сервису, не больше пяти одновременных выгрузок в S3.

_s3_slots = asyncio.Semaphore(5)

async def upload(key: str, data: bytes) -> None:
    async with _s3_slots:
        await s3.put_object(Bucket=BUCKET, Key=key, Body=data)

Задачи сверх лимита ждут в очереди семафора. Это обратное давление в чистом виде, и о его месте в архитектуре обработчика рассказывает статья про таймауты и обратное давление. BoundedSemaphore дополнительно падает, если освободить его больше раз, чем захватили, что ловит ошибку парности при ручном acquire и release.

Event и Condition: сигналы между задачами

asyncio.Event это флаг с ожиданием: одни задачи ждут await event.wait(), кто-то один делает event.set(), и все ждущие просыпаются. Типичное применение: «приложение готово», «пора останавливаться». Фоновый цикл опроса завершается по событию остановки, а не по cancel, если нужно доделать текущую итерацию:

stop = asyncio.Event()

async def poll_loop():
    while not stop.is_set():
        await poll_once()
        try:
            await asyncio.wait_for(stop.wait(), timeout=5)   # пауза, прерываемая остановкой
        except TimeoutError:
            pass

asyncio.Condition это замок плюс ожидание условия с уведомлением: await cond.wait_for(lambda: len(buffer) > 0) и cond.notify() при изменении. Нужен, когда условие сложнее флага, но в прикладном коде он редкий гость: почти всегда вместо Condition правильнее Queue.

Queue: конвейер с обратным давлением

asyncio.Queue связывает производителей и потребителей. С maxsize она даёт обратное давление: await queue.put(item) ждёт, пока потребители не освободят место, и производитель замедляется до скорости потребителей вместо того, чтобы копить миллион элементов в памяти. Проверено: у Queue(maxsize=2) третий put_nowait поднимает QueueFull, а await put ждёт.

async def process_all(items: AsyncIterator[Item], workers: int = 10) -> None:
    queue: asyncio.Queue[Item | None] = asyncio.Queue(maxsize=100)

    async def worker():
        while (item := await queue.get()) is not None:
            try:
                await handle(item)
            finally:
                queue.task_done()

    async with asyncio.TaskGroup() as tg:
        for _ in range(workers):
            tg.create_task(worker())
        async for item in items:
            await queue.put(item)
        for _ in range(workers):
            await queue.put(None)                       # сигнал завершения каждому потребителю

Десять потребителей, очередь на сто элементов, производитель не убегает вперёд. task_done и join нужны, когда надо дождаться обработки всех элементов, а не только их постановки. Для приоритетов есть PriorityQueue, для стека LifoQueue.

Чего примитивы asyncio не умеют

Все они работают внутри одного цикла событий и не защищают от потоков. Если данные меняет и задача asyncio, и поток из to_thread, asyncio.Lock не поможет; нужен threading.Lock или передача результата в цикл через call_soon_threadsafe, об этом статья про блокирующий код и потоки. Они не защищают от других процессов и подов: лимит «не больше пяти выгрузок» действует внутри одного процесса, и при четырёх воркерах uvicorn их двадцать; общий лимит на сервис живёт в Redis или в семафоре на стороне принимающей системы. И они не переживают перезапуск: очередь asyncio это память процесса, всё, что в ней было при падении, потеряно; надёжная очередь это Kafka, RabbitMQ или таблица в PostgreSQL.

Дополнительно: при первом чтении можно пропустить

Глубже: примитив на уровне цикла или на уровне ключарасширенное

Один Lock на модуль прост, но превращает независимые операции в последовательные. Там, где операции независимы по ключу (кеш по валюте, дедупликация по идентификатору заказа), нужен замок на ключ, и наивный словарь dict[str, asyncio.Lock] течёт: замки создаются и не удаляются. Приём с Future решает и гонку, и утечку: первая задача кладёт в словарь незавершённый Future и выполняет работу, остальные находят Future и ждут его, по завершении результат раздаётся всем, а запись удаляется:

_inflight: dict[str, asyncio.Future[Rate]] = {}

async def get_rate(currency: str) -> Rate:
    if currency in _cache:
        return _cache[currency]
    if currency in _inflight:
        return await _inflight[currency]
    fut = asyncio.get_running_loop().create_future()
    _inflight[currency] = fut
    try:
        rate = await fetch_rate(currency)
        _cache[currency] = rate
        fut.set_result(rate)
        return rate
    except BaseException as e:
        fut.set_exception(e)
        raise
    finally:
        _inflight.pop(currency, None)

Это защита от лавины запросов к источнику при протухании кеша (cache stampede) без единого замка, и она же лежит в основе дедупликации одновременных одинаковых запросов. Нюанс с set_exception при CancelledError: если первую задачу отменили, ждущие получат CancelledError как исключение; для них честнее повторить попытку самим, поэтому в боевом коде ожидающие оборачивают await fut в проверку и при отмене лидера становятся лидером сами.

Коротко

  • Гонка в asyncio возникает между проверкой и действием, если между ними есть await; код между двумя await атомарен.
  • asyncio.Lock для критической секции с await внутри, с повторной проверкой под замком; он не потокобезопасен и не реентерабелен.
  • Semaphore(n) это лимит одновременных операций на внешнюю систему, основной инструмент обратного давления внутри процесса.
  • Event для сигналов «готово» и «останавливаемся», Condition почти всегда заменяется Queue.
  • Queue(maxsize=...) даёт конвейер с обратным давлением: put ждёт потребителей; завершение потребителей через сигнальный элемент или отмену группы.
  • Примитивы asyncio не защищают от потоков, других процессов и перезапуска; общие лимиты и надёжные очереди живут снаружи.
  • Замок на ключ без утечек делают через словарь Future: первый выполняет, остальные ждут, запись удаляется по завершении.

Что почитать дальше