41 lines
1.2 KiB
Python
41 lines
1.2 KiB
Python
"""进程内聚合锁原语(B7 · 挂账①:从 alert_service 公共化,profile_l3 同用)。
|
|
|
|
单进程内按 key 串行减少并发首单冲突(多进程部署换 Redis SET NX,接口不变);
|
|
拿锁超时降级独立执行,冲突安全由调用方兜底(预警聚合:同日单条锚点查询在
|
|
锁内重查;L3:乐观锁重试),宁多勿漏。
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import threading
|
|
from typing import Any, Callable
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
LOCK_TIMEOUT_SECONDS = 2.0
|
|
|
|
_locks: dict[str, threading.Lock] = {}
|
|
_locks_guard = threading.Lock()
|
|
|
|
|
|
def lock_for(key: str) -> threading.Lock:
|
|
with _locks_guard:
|
|
lock = _locks.get(key)
|
|
if lock is None:
|
|
lock = threading.Lock()
|
|
_locks[key] = lock
|
|
return lock
|
|
|
|
|
|
def run_locked(key: str, fn: Callable[[bool], Any]) -> Any:
|
|
"""锁内执行 fn(locked=True);获取超时降级 fn(locked=False)。"""
|
|
lock = lock_for(key)
|
|
if not lock.acquire(timeout=LOCK_TIMEOUT_SECONDS):
|
|
logger.warning("agg lock timeout, run without lock (conflicts bounded by caller): %s", key)
|
|
return fn(locked=False)
|
|
try:
|
|
return fn(locked=True)
|
|
finally:
|
|
lock.release()
|