← все задачи
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, но в асинхронном мире — и с вопросом, что будет с отменой во время паузы.