← все задачи

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

Очередь: производитель и потребители

Продвинутый 25 минут asyncio.Queueворкерыgraceful shutdown

Условие

Данные приходят потоком (например, читаются из файла или из API постранично). Обработать каждый элемент нужно медленной асинхронной операцией. Сделайте пул из N воркеров, которые разбирают очередь, и корректно завершите работу, когда данные кончились.

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

  • Производитель не должен забивать память: очередь ограничена по размеру
  • Воркеров N, они работают параллельно
  • После окончания данных программа завершается, не подвисая

Пример

await run(source=read_pages(), handle=save_to_db, workers=5, queue_size=100)
# 5 воркеров разбирают очередь, в очереди не больше 100 элементов

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

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

  • Что делать с ошибкой в воркере: пропустить элемент, повторить, остановить всё?
  • Важен ли порядок обработки или элементы независимы?
  • Нужна ли гарантия «обработано хотя бы раз» — тогда нужен ack и хранилище, а не память процесса
Показать решение Скрыть решение

Решение

import asyncio
import logging

logger = logging.getLogger(__name__)


async def run(source, handle, workers=5, queue_size=100):
    queue = asyncio.Queue(maxsize=queue_size)

    async def worker(number):
        while True:
            item = await queue.get()
            try:
                await handle(item)
            except Exception:
                logger.exception("воркер %s не смог обработать %r", number, item)
            finally:
                queue.task_done()      # даже при ошибке: иначе join() не дождётся

    tasks = [asyncio.create_task(worker(i)) for i in range(workers)]

    async for item in source:
        await queue.put(item)          # ждёт, если очередь заполнена

    await queue.join()                 # дожидаемся обработки всего, что положили

    for task in tasks:
        task.cancel()
    await asyncio.gather(*tasks, return_exceptions=True)

Почему так

Зачем ограничивать размер очереди

  • Безразмерная очередь превращает быстрого производителя в утечку памяти: миллион строк из файла окажется в RAM
  • maxsize делает put ожидающим — это и есть обратное давление, производитель сам притормаживает
  • Размер очереди — регулятор: слишком мал — воркеры простаивают, слишком велик — растёт память и задержка

Почему task_done в finally

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

Почему воркеры отменяются в конце

  • Воркеры написаны как бесконечный цикл — после join они висят на await queue.get()
  • Без явного cancel программа не завершится, а при выходе вы получите предупреждение о незавершённых задачах
  • Альтернатива — послать в очередь по «стоп-элементу» на каждого воркера; вариант с cancel короче и не смешивает данные с управлением

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

  • Спросят про ошибки: сейчас элемент теряется — правильное продолжение это повтор или очередь неудачных
  • Спросят, чем это отличается от Celery или Kafka: здесь всё в памяти одного процесса, при падении данные теряются
  • Спросят про сигнал остановки (SIGTERM): нужно перестать читать источник, дать воркерам доработать очередь и только потом выходить

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

Повторы для асинхронного вызова — Тот же retry, но в асинхронном мире — и с вопросом, что будет с отменой во время паузы.