forked from Boris-code/feapder
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathredis_lock.py
More file actions
127 lines (112 loc) · 3.87 KB
/
Copy pathredis_lock.py
File metadata and controls
127 lines (112 loc) · 3.87 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
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
121
122
123
124
125
126
127
# -*- coding: utf-8 -*-
"""
Created on 2019/11/5 5:25 PM
---------
@summary:
---------
@author: Boris
@email: boris@bzkj.tech
"""
import time
import feapder.utils.log as log
class RedisLock(object):
def __init__(
self,
key,
timeout=300,
wait_timeout=8 * 3600,
break_wait=None,
redis_cli=None,
logger=None,
):
"""
redis超时锁
:param key: 关键字 不同项目区分
:param timeout: 锁超时时间
:param wait_timeout: 等待加锁超时时间 默认8小时 防止多线程竞争时可能出现的 某个线程无限等待
<=0 则不等待 直接加锁失败
:param break_wait: 可自定义函数 灵活控制 wait_timeout 时间 当此函数返回True时 不再wait
:param redis_cli: redis客户端
用法示例:
with RedisLock(key="test", timeout=10, wait_timeout=100, redis_uri="") as _lock:
if _lock.locked:
# 用来判断是否加上了锁
# do somethings
"""
self.redis_index = -1
if not key:
raise Exception("lock key is empty")
if not redis_cli:
raise Exception("redis_cli is empty")
self.redis_conn = redis_cli
self.logger = logger or log.get_logger(__file__)
self.lock_key = "redis_lock:{}".format(key)
# 锁超时时间
self.timeout = timeout
# 等待加锁时间
self.wait_timeout = wait_timeout
# wait中断函数
self.break_wait = break_wait
if self.break_wait is None:
self.break_wait = lambda: False
if not callable(self.break_wait):
raise TypeError(
"break_wait must be function or None, but: {}".format(
type(self.break_wait)
)
)
self.locked = False
def __enter__(self):
if not self.locked:
self.acquire()
return self
def __exit__(self, exc_type, exc_val, exc_tb):
self.release()
def __repr__(self):
return "<RedisLock: {} index: {}>".format(self.lock_key, self.redis_index)
def acquire(self):
start = time.time()
self.logger.debug("准备获取锁{} ...".format(self))
while 1:
# 尝试加锁
if self.redis_conn.setnx(self.lock_key, time.time()):
self.redis_conn.expire(self.lock_key, self.timeout)
self.locked = True
self.logger.debug("加锁成功: {}".format(self))
break
else:
# 修复bug: 当加锁时被干掉 导致没有设置expire成功 锁无限存在
if self.redis_conn.ttl(self.lock_key) < 0:
self.redis_conn.delete(self.lock_key)
if self.wait_timeout > 0:
if time.time() - start > self.wait_timeout:
break
else:
# 不等待
break
if self.break_wait():
self.logger.debug("break_wait 生效 不再等待加锁")
break
self.logger.debug("等待加锁: {} wait:{}".format(self, time.time() - start))
if self.wait_timeout > 10:
time.sleep(5)
else:
time.sleep(1)
return
def release(self):
if self.locked:
self.redis_conn.delete(self.lock_key)
self.locked = False
return
def prolong_life(self, life_time: int) -> int:
"""
延长这个锁的超时时间
:param life_time: 延长时间
:return:
"""
expire = self.redis_conn.ttl(self.lock_key)
if expire < 0:
return expire
expire += life_time
self.redis_conn.expire(self.lock_key, expire)
return self.redis_conn.ttl(self.lock_key)