Раз 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: первый выполняет, остальные ждут, запись удаляется по завершении.
Что почитать дальше
- Таймауты и обратное давление — как семафоры и очереди складываются в защиту сервиса под нагрузкой.
- Блокирующий код и потоки — где заканчивается однопоточная безопасность и начинаются потоки.
- Защита от лавины запросов к кешу — тот же приём на уровне Redis для нескольких процессов.