Skip to content

Commit 15ed732

Browse files
committed
优化锁的问题
1 parent a0a1e84 commit 15ed732

5 files changed

Lines changed: 42 additions & 16 deletions

File tree

feapder/core/scheduler.py

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -263,10 +263,7 @@ def _start(self):
263263
if self.wait_lock:
264264
# 将添加任务处加锁,防止多进程之间添加重复的任务
265265
with RedisLock(
266-
key=self._spider_name,
267-
timeout=3600,
268-
wait_timeout=60,
269-
redis_cli=RedisDB().get_redis_obj(),
266+
key=self._spider_name, redis_cli=RedisDB().get_redis_obj()
270267
) as lock:
271268
if lock.locked:
272269
self.__add_task()

feapder/core/spiders/batch_spider.py

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -946,10 +946,7 @@ def task_is_done(self):
946946
if is_done: # 检查任务表中是否有没做的任务 若有则is_done 为 False
947947
# 比较耗时 加锁防止多进程同时查询
948948
with RedisLock(
949-
key=self._spider_name,
950-
timeout=3600,
951-
wait_timeout=0,
952-
redis_cli=RedisDB().get_redis_obj(),
949+
key=self._spider_name, redis_cli=RedisDB().get_redis_obj()
953950
) as lock:
954951
if lock.locked:
955952
log.info("批次表标记已完成,正在检查任务表是否有未完成的任务")

feapder/dedup/bloomfilter.py

Lines changed: 16 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -205,6 +205,8 @@ class ScalableBloomFilter(object):
205205
BASE_MEMORY = BloomFilter.BASE_MEMORY
206206
BASE_REDIS = BloomFilter.BASE_REDIS
207207

208+
__redis_cli = None
209+
208210
def __init__(
209211
self,
210212
initial_capacity: int = 100000000,
@@ -237,6 +239,13 @@ def _setup(self, initial_capacity, error_rate, name, bitarray_type, redis_url):
237239
def __repr__(self):
238240
return "<ScalableBloomFilter: {}>".format(self.filters[-1].bitarray)
239241

242+
@property
243+
def _redis_cli(self):
244+
if self.__class__.__redis_cli is None:
245+
self.__class__.__redis_cli = RedisDB(url=self.redis_url).get_redis_obj()
246+
247+
return self.__class__.__redis_cli
248+
240249
def create_filter(self):
241250
filter = BloomFilter(
242251
capacity=self.initial_capacity,
@@ -267,12 +276,13 @@ def check_filter_capacity(self):
267276

268277
self._check_capacity_time = time.time()
269278
else:
270-
with RedisLock(
271-
key="ScalableBloomFilter",
272-
timeout=300,
273-
wait_timeout=300,
274-
redis_cli=RedisDB(url=self.redis_url).get_redis_obj(),
275-
) as lock: # 全局锁 同一时间只有一个进程在真正的创建新的filter,等这个进程创建完,其他进程只是把刚创建的filter append进来
279+
# 全局锁 同一时间只有一个进程在真正的创建新的filter,等这个进程创建完,其他进程只是把刚创建的filter append进来
280+
key = (
281+
f"ScalableBloomFilter:{self.name}"
282+
if self.name
283+
else "ScalableBloomFilter"
284+
)
285+
with RedisLock(key=key, redis_cli=self._redis_cli) as lock:
276286
if lock.locked:
277287
while True:
278288
if self.filters[-1].is_at_capacity:

feapder/utils/redis_lock.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -14,13 +14,13 @@
1414

1515

1616
class RedisLock(object):
17-
def __init__(self, key, redis_cli, wait_timeout=0, lock_timeout=0):
17+
def __init__(self, key, redis_cli, wait_timeout=0, lock_timeout=86400):
1818
"""
1919
redis超时锁
2020
:param key: 存储锁的key redis_lock:[key]
2121
:param redis_cli: redis客户端对象
2222
:param wait_timeout: 等待加锁超时时间,为0时则不等待加锁,加锁失败
23-
:param lock_timeout: 锁超时时间 为0时则不会超时,直到锁释放或意外退出
23+
:param lock_timeout: 锁超时时间 为0时则不会超时,直到锁释放或意外退出,默认超时为1天
2424
2525
用法示例:
2626
with RedisLock(key="test", redis_cli=redis_obj) as _lock:

tests/test_lock.py

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
# -*- coding: utf-8 -*-
2+
"""
3+
Created on 2021/7/15 5:00 下午
4+
---------
5+
@summary:
6+
---------
7+
@author: Boris
8+
@email: boris_liu@foxmail.com
9+
"""
10+
11+
from feapder.utils.redis_lock import RedisLock
12+
from feapder.db.redisdb import RedisDB
13+
import time
14+
15+
def test_lock():
16+
with RedisLock(key="test", redis_cli=RedisDB().get_redis_obj(), wait_timeout=10) as _lock:
17+
if _lock.locked:
18+
print(1)
19+
time.sleep(100)
20+
21+
if __name__ == '__main__':
22+
test_lock()

0 commit comments

Comments
 (0)