TGViewer
About Python [ru] About Python [ru] @python_tesst · 6.45K subscribers
Post #2794 345
⁣Асинхронные краны сообщений через asyncio.Queue с приоритетами и тайм-аутами: диспетчер задач без Celery

Когда Celery — это перебор, а Redis-очередь лень поднимать, многие кидаются на asyncio.Queue. Но голая FIFO не отдаст приоритет срочным задачам, и зависший воркер подвесит всю систему. Решение — собственный диспетчер с приоритетами и тайм-аутами.

Датакласс с сортировкой
Используем @dataclass(order=True), чтобы задачи сортировались по приоритету автоматически. Поле timeout задаёт лимит на выполнение, а payload — полезная нагрузка. Это даёт чистый интерфейс без внешних зависимостей.

Воркер с тайм-аутом
Обёртка asyncio.wait_for в цикле воркера режет задачу по тайм-ауту. Если не уложилась — кидаем TimeoutError, но очередь не ломается, и воркер спокойно переходит к следующей. Это ключевой приём для production, где один зависший IO-запрос не должен блокировать остальные.

Production-пример
В реальном проекте (обработка пайплайна алертов) приоритеты определяют важность: priority=1 для критических, priority=5 для отчётов. Тайм-аут на 2 секунды не даёт медленному внешнему API тянуть всю очередь. Код минимален и работает на Python 3.7+.

import asyncio
from dataclasses import dataclass, field

@dataclass(order=True)
class Task:
priority: int
timeout: int = field(default=10, compare=False)
payload: str = field(default=None, compare=False)

async def worker(queue: asyncio.PriorityQueue, name: str):
while True:
task: Task = await queue.get()
try:
await asyncio.wait_for(process(task.payload), timeout=task.timeout)
except asyncio.TimeoutError:
print(f"[{name}] Timed out priority {task.priority}")
finally:
queue.task_done()


Типичная ошибка
Начинающие забывают про queue.task_done() — это ведёт к зависанию queue.join(). Или не оборачивают задачу в asyncio.wait_for, тогда одна долгая операция блокирует всех воркеров.

Трейд-оффы
Вся очередь живёт в памяти — упадёт процесс, задачи потеряны. Нет распределённости: воркеры работают в одном процессе. На высоких нагрузках (10k+ задач/сек) вставка O(log n) через heapq может просадить производительность. Для повышения надёжности добавляйте asyncio.Semaphore для лимита параллельных задач и retry с повышением приоритета.

Вывод: Для простых однопроцессных сценариев с приоритетами и тайм-аутами asyncio.PriorityQueue — работающий лёгкий инструмент, но без персистентности и распределённости он не конкурент Celery на кластерных нагрузках.
More from @python_tesst
  1. Sep 28, 2026Post #3191
  2. Sep 27, 2026OpenAI метит в подписку за 500 баксов Что там в описании тарифа? Пока что от ChatGPT Pro о…
  3. Sep 27, 2026Post #3189
  4. Sep 27, 2026Свежая обложка The Economist подъехала Журналисты: да мы вообще не сгущаем краски Те же жу…
  5. Sep 27, 2026Ночная годнота: Docker выкатила первые официальные скиллы для ИИ-агентов, которые ковыряют…
  6. Sep 26, 2026Claude Code больше не будет рубить всё на полуслове, если 5-часовой лимит прилетит прямо п…
Threads Profile ViewerView any public Threads profile without an account.Open ThreadLook →Writing with AI? Make it sound human.Metric37 rewrites AI drafts so they read naturally. Free AI detector, 1,500 words free.Try Metric37 →