Skip to content

DistributedLock

Helper for distributed locking using the configured backend.

分布式锁

基于缓存后端的分布式锁实现,支持超时和自动释放。

使用示例: >>> lock = DistributedLock(manager, "resource:123", timeout=10) >>> if lock.acquire(): ... try: ... # 处理资源 ... pass ... finally: ... lock.release()

Source code in src/symphra_cache/locks.py
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
class DistributedLock:
    """
    分布式锁

    基于缓存后端的分布式锁实现,支持超时和自动释放。

    使用示例:
        >>> lock = DistributedLock(manager, "resource:123", timeout=10)
        >>> if lock.acquire():
        ...     try:
        ...         # 处理资源
        ...         pass
        ...     finally:
        ...         lock.release()
    """

    def __init__(
        self,
        manager: CacheManager,
        name: str,
        timeout: int = 10,
        blocking: bool = True,
        blocking_timeout: float | None = None,
    ) -> None:
        """
        初始化分布式锁

        Args:
            manager: 缓存管理器
            name: 锁名称
            timeout: 锁超时时间(秒)
            blocking: 是否阻塞等待
            blocking_timeout: 阻塞超时(秒)
        """
        self.manager = manager
        self.name = f"lock:{name}"
        self.timeout = timeout
        self.blocking = blocking
        self.blocking_timeout = blocking_timeout
        self.identifier = str(uuid.uuid4())  # 唯一标识符
        self._locked = False

    def acquire(self) -> bool:
        """
        获取锁

        Returns:
            是否成功获取锁
        """
        start_time = time.time()

        while True:
            # 尝试设置锁(使用 TTL 防止死锁)
            existing = self.manager.get(self.name)

            if existing is None:
                # 锁不存在,尝试获取
                self.manager.set(self.name, self.identifier, ttl=self.timeout)
                self._locked = True
                return True

            if not self.blocking:
                return False

            # 检查阻塞超时
            if (
                self.blocking_timeout is not None
                and time.time() - start_time >= self.blocking_timeout
            ):
                return False

            # 短暂休眠后重试
            time.sleep(0.01)

    def release(self) -> None:
        """释放锁"""
        if not self._locked:
            return

        # 验证是否是自己的锁
        current_value = self.manager.get(self.name)
        if current_value == self.identifier:
            self.manager.delete(self.name)
            self._locked = False

    def __enter__(self) -> DistributedLock:
        """上下文管理器:进入"""
        self.acquire()
        return self

    def __exit__(self, exc_type: object, exc_val: object, exc_tb: object) -> None:
        """上下文管理器:退出"""
        self.release()

__enter__()

上下文管理器:进入

Source code in src/symphra_cache/locks.py
113
114
115
116
def __enter__(self) -> DistributedLock:
    """上下文管理器:进入"""
    self.acquire()
    return self

__exit__(exc_type, exc_val, exc_tb)

上下文管理器:退出

Source code in src/symphra_cache/locks.py
118
119
120
def __exit__(self, exc_type: object, exc_val: object, exc_tb: object) -> None:
    """上下文管理器:退出"""
    self.release()

__init__(manager, name, timeout=10, blocking=True, blocking_timeout=None)

初始化分布式锁

Parameters:

Name Type Description Default
manager CacheManager

缓存管理器

required
name str

锁名称

required
timeout int

锁超时时间(秒)

10
blocking bool

是否阻塞等待

True
blocking_timeout float | None

阻塞超时(秒)

None
Source code in src/symphra_cache/locks.py
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
def __init__(
    self,
    manager: CacheManager,
    name: str,
    timeout: int = 10,
    blocking: bool = True,
    blocking_timeout: float | None = None,
) -> None:
    """
    初始化分布式锁

    Args:
        manager: 缓存管理器
        name: 锁名称
        timeout: 锁超时时间(秒)
        blocking: 是否阻塞等待
        blocking_timeout: 阻塞超时(秒)
    """
    self.manager = manager
    self.name = f"lock:{name}"
    self.timeout = timeout
    self.blocking = blocking
    self.blocking_timeout = blocking_timeout
    self.identifier = str(uuid.uuid4())  # 唯一标识符
    self._locked = False

acquire()

获取锁

Returns:

Type Description
bool

是否成功获取锁

Source code in src/symphra_cache/locks.py
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
def acquire(self) -> bool:
    """
    获取锁

    Returns:
        是否成功获取锁
    """
    start_time = time.time()

    while True:
        # 尝试设置锁(使用 TTL 防止死锁)
        existing = self.manager.get(self.name)

        if existing is None:
            # 锁不存在,尝试获取
            self.manager.set(self.name, self.identifier, ttl=self.timeout)
            self._locked = True
            return True

        if not self.blocking:
            return False

        # 检查阻塞超时
        if (
            self.blocking_timeout is not None
            and time.time() - start_time >= self.blocking_timeout
        ):
            return False

        # 短暂休眠后重试
        time.sleep(0.01)

release()

释放锁

Source code in src/symphra_cache/locks.py
102
103
104
105
106
107
108
109
110
111
def release(self) -> None:
    """释放锁"""
    if not self._locked:
        return

    # 验证是否是自己的锁
    current_value = self.manager.get(self.name)
    if current_value == self.identifier:
        self.manager.delete(self.name)
        self._locked = False