Стандартный asyncio.Lock не поддерживает дедлайны и приоритеты — в production это приводит к инверсии приоритетов и зависанию воркеров. Без наследования приоритетов низкоприоритетная задача может блокировать критичный поток, а без дедлайна блокировка навсегда замораживает воркер.
Проблема и решение
В high-throughput системе с асинхронными воркерами (например, Celery + asyncio) конкуренция за общий ресурс без управления приоритетами вызывает голодание. Решение — кастомный диспетчер поверх asyncio.Lock с очередью на heapq и механизмом приоритетного наследования: когда высокоприоритетный воркер ждёт блокировку, его приоритет временно передаётся текущему владельцу.
Архитектура диспетчера
Ключевые элементы:
* Очередь_Request с priority и deadline, реализованная через heapq. Приоритет инвертируется (отрицание) для корректной сортировки.
* Рекурсивный счётчик для RLock-совместимости — блокировка может быть захвачена тем же воркером повторно.
* Future для пробуждения: каждый запрос ассоциирован с asyncio.Future, которая вызывается при освобождении.
* При наследовании приоритет владельца повышается до максимального из ожидающих, при освобождении сбрасывается.
Упрощённый код ядра
import asyncio, heapq
from dataclasses import dataclass, field
@dataclass(order=True)
class _LockRequest:
deadline: float
priority: int
future: asyncio.Future = field(compare=False)
class PriorityRLock:
def __init__(self):
self._owner = None
self._current_priority = 0
self._queue = []
self._recursion_level = 0
self._cond = asyncio.Condition()
async def acquire(self, priority=0, timeout=None):
loop = asyncio.get_event_loop()
deadline = loop.time() + timeout if timeout else float('inf')
future = loop.create_future()
heapq.heappush(self._queue, _LockRequest(deadline, -priority, future))
async with self._cond:
if self._owner and priority > self._current_priority:
self._current_priority = priority # наследование
while True:
if self._owner is None and self._queue[0].future == future:
self._owner = id(asyncio.current_task())
self._current_priority = priority
self._recursion_level += 1
return True
try:
await asyncio.wait_for(future, timeout)
except asyncio.TimeoutError:
self._queue.remove(future)
heapq.heapify(self._queue)
return False
async def release(self):
async with self._cond:
self._recursion_level -= 1
if self._recursion_level == 0:
self._owner = None
self._current_priority = 0
if self._queue:
top = self._queue[0]
loop.call_soon(top.future.set_result, None)
Типичная ошибка и практический совет
Ошибка: забыть про сортировку heap при удалении — use _LockRequest с compare=False для future, иначе сравнение объектов ломает кучу. Практический совет: в production замените asyncio.wait_for на кастомный таймер с асинхронным сном, чтобы избежать оверхеда исключений при частых таймаутах. Для простых кейсов (один пул, без жёстких дедлайнов) используйте Lock с timeout — не усложняйте.
Когда это критично?
* Очереди с приоритетами в Celery/arq для воркеров к общему ресурсу (БД, кэш).
* Системы с жёсткими дедлайнами, где превышение времени блокировки порождает цепную реакцию.
* High-throughput сценарии, где инверсия приоритетов проявляется только под нагрузкой. Минусы: сложность отладки и оверхед на сортировку Heap (O(log n)). Но в продакшене без этого ловятся странные гонки и зависания, не воспроизводимые локально.
Вывод: Проектируйте диспетчер блокировок с приоритетным наследованием и дедлайнами только под реальную нагрузку, иначе ст