Пишем на Python, создаём нейросети и ИИ-агентов.
Алгоритмы, задачи и вайбкодинг.
Личный блог автора - @just_genych
По вопросам рекламы или разработки: @g_abashkin
Post #2794
345
Асинхронные краны сообщений через
Когда Celery — это перебор, а Redis-очередь лень поднимать, многие кидаются на
Датакласс с сортировкой
Используем
Воркер с тайм-аутом
Обёртка
Production-пример
В реальном проекте (обработка пайплайна алертов) приоритеты определяют важность: priority=1 для критических, priority=5 для отчётов. Тайм-аут на 2 секунды не даёт медленному внешнему API тянуть всю очередь. Код минимален и работает на Python 3.7+.
Типичная ошибка
Начинающие забывают про
Трейд-оффы
Вся очередь живёт в памяти — упадёт процесс, задачи потеряны. Нет распределённости: воркеры работают в одном процессе. На высоких нагрузках (10k+ задач/сек) вставка O(log n) через heapq может просадить производительность. Для повышения надёжности добавляйте
Вывод: Для простых однопроцессных сценариев с приоритетами и тайм-аутами
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 на кластерных нагрузках.



