Программирование PythonКонкурентность и asyncioPython-разработчик асинхронных сервисов

Практическая ситуация: определите, когда завершится третий вызов put и почему, если потребитель начнёт рабо...

Практическая ситуация: определите, когда завершится третий вызов put и почему, если потребитель начнёт работу через секунду.

import asyncio

async def producer(queue):
    for item in range(3):
        print("before", item)
        await queue.put(item)
        print("after", item)

async def main():
    queue = asyncio.Queue(maxsize=2)
    task = asyncio.create_task(producer(queue))
    await asyncio.sleep(1)
    print("got", await queue.get())
    await task

asyncio.run(main())
Проходите собеседования с ИИ помощником Hintsage

Краткий ответ

Третий вызов queue.put(2) сначала заблокируется: очередь уже содержит два элемента и достигла maxsize=2. Он продолжится после того, как queue.get() освободит место; это механизм обратного давления, или backpressure, от потребителя к производителю.

Исторический контекст

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

В asyncio эта модель позволяет координировать корутины без блокировки потока операционной системы. Ожидающая put() корутина уступает управление event loop, поэтому другие готовые задачи могут продолжать выполняться.

Постановка проблемы

В примере первые два элемента помещаются в очередь немедленно. На третьем элементе producer достигает await queue.put(2) и приостанавливается, потому что свободных слотов нет.

Это не ошибка и не зависание event loop: корутина ждёт конкретное событие — извлечение элемента потребителем. Если потребитель вообще не вызовет get(), производитель останется ожидающим, а если использовать неограниченную очередь, память может расти без контроля.

Подробное решение

Параметр maxsize=2 задаёт максимальное число элементов, находящихся в очереди. Когда очередь заполнена, put() не возвращается сразу, а регистрирует ожидающую корутину; после get() event loop возобновляет её, и третий элемент добавляется.

Важное различие: await queue.put(item) не блокирует весь поток, в отличие от синхронной блокировки. При ожидании освобождения места event loop может выполнять другие корутины.

import asyncio async def producer(q): for item in range(3): await q.put(item) print("добавлен", item) async def main(): q = asyncio.Queue(maxsize=2) producer_task = asyncio.create_task(producer(q)) await asyncio.sleep(0) print("извлечён", await q.get()) await producer_task asyncio.run(main())

maxsize=0 означает неограниченную очередь, поэтому put() не создаёт обратного давления из-за размера. Очередь asyncio.Queue предназначена для взаимодействия корутин в одном event loop и не является универсальной потокобезопасной очередью для разных потоков.

get() только извлекает элемент. Если нужно дождаться завершения обработки всех помещённых элементов, применяют queue.task_done() после обработки каждого элемента и await queue.join() для ожидания опустошения счётчика незавершённой работы.

Следует учитывать отмену ожидающей операции, ошибки потребителя и корректное завершение производителей. Ограниченный размер снижает риск переполнения памяти, но может уменьшить пропускную способность и увеличить задержку производителя.

Ситуация из практики

В сервисе приёма событий HTTP-обработчики быстро складывают события в очередь, а несколько корутин отправляют их во внешнюю систему. При временной недоступности внешней системы потребители замедляются.

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

Практичное решение — ограниченная asyncio.Queue, явная политика отказа или тайм-аут для помещения элемента и контролируемое число потребителей. Производитель получает возможность дождаться свободного места либо сообщить вызывающему коду, что сервис временно перегружен; память остаётся ограниченной, а нагрузка на внешний сервис — управляемой.

Что кандидаты часто упускают

  1. Что изменится при maxsize=0?

    Значение 0 означает отсутствие ограничения размера, а не очередь нулевой ёмкости. put() не будет ждать освобождения места из-за переполнения, поэтому обратное давление через размер очереди исчезает. Ограничение памяти в таком случае нужно реализовывать отдельно — например, через конечный размер очереди, лимит входных запросов или другую политику управления нагрузкой.

  2. Освобождает ли get() место только после завершения обработки элемента?

    Нет. get() удаляет элемент из очереди сразу, поэтому место для нового put() освобождается до фактической обработки. Для отслеживания завершения работы потребитель должен вызвать task_done() после обработки, а join() ждёт именно эти отметки, а не просто отсутствие элементов в очереди.

  3. Что произойдёт, если ожидающий put() отменить?

    Если отмена доставлена во время ожидания свободного места, корутина получает asyncio.CancelledError, а конкретная операция put() не считается успешно выполненной. Вызывающий код должен решить, нужно ли пробросить отмену, повторить операцию или зафиксировать отказ; без такой обработки можно потерять элемент или некорректно завершить конвейер.