Skip to content

Commit a00a951

Browse files
author
Boris
committed
添加去重库文档
1 parent 6acda7f commit a00a951

10 files changed

Lines changed: 329 additions & 91 deletions

File tree

docs/source_code/dedup.md

Lines changed: 118 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,120 @@
11
# Dedup
22

3-
未完待续
3+
Dedup是feapder大数据去重模块,内置3种去重机制,使用方式一致,可容纳的去重数据量与内存有关。不同于BloomFilter,去重受槽位数量影响,Dedup使用了弹性的去重机制,可容纳海量的数据去重。
4+
5+
6+
## 去重方式
7+
8+
### 临时去重
9+
10+
> 基于redis,支持批量,去重有时效性。去重一万条数据约0.26秒,一亿条数据占用内存约1.43G
11+
12+
```python
13+
from feapder.dedup import Dedup
14+
15+
data = {"xxx": 123, "xxxx": "xxxx"}
16+
datas = ["xxx", "bbb"]
17+
18+
def test_ExpireFilter():
19+
dedup = Dedup(
20+
Dedup.ExpireFilter, expire_time=10, redis_url="redis://@localhost:6379/0"
21+
)
22+
23+
# 逐条去重
24+
assert dedup.add(data) == 1
25+
assert dedup.get(data) == 1
26+
27+
# 批量去重
28+
assert dedup.add(datas) == [1, 1]
29+
assert dedup.get(datas) == [1, 1]
30+
```
31+
32+
33+
### 内存去重
34+
35+
> 基于内存,支持批量。去重一万条数据约0.5秒,一亿条数据占用内存约285MB
36+
37+
```python
38+
from feapder.dedup import Dedup
39+
40+
data = {"xxx": 123, "xxxx": "xxxx"}
41+
datas = ["xxx", "bbb"]
42+
43+
def test_MemoryFilter():
44+
dedup = Dedup(Dedup.MemoryFilter) # 表名为test 历史数据3秒有效期
45+
46+
# 逐条去重
47+
assert dedup.add(data) == 1
48+
assert dedup.get(data) == 1
49+
50+
# 批量去重
51+
assert dedup.add(datas) == [1, 1]
52+
assert dedup.get(datas) == [1, 1]
53+
```
54+
55+
### 永久去重
56+
57+
> 基于redis,支持批量,永久去重。 去重一万条数据约3.5秒,一亿条数据占用内存约285MB
58+
59+
```python
60+
from feapder.dedup import Dedup
61+
62+
def test_BloomFilter():
63+
dedup = Dedup(Dedup.BloomFilter, redis_url="redis://@localhost:6379/0")
64+
65+
# 逐条去重
66+
assert dedup.add(data) == 1
67+
assert dedup.get(data) == 1
68+
69+
# 批量去重
70+
assert dedup.add(datas) == [1, 1]
71+
assert dedup.get(datas) == [1, 1]
72+
```
73+
74+
## 过滤数据
75+
76+
Dedup可以通过如下方法,过滤掉已存在的数据
77+
78+
79+
```python
80+
from feapder.dedup import Dedup
81+
82+
def test_filter():
83+
dedup = Dedup(Dedup.BloomFilter, redis_url="redis://@localhost:6379/0")
84+
85+
# 制造已存在数据
86+
datas = ["xxx", "bbb"]
87+
dedup.add(datas)
88+
89+
# 过滤掉已存在数据 "xxx", "bbb"
90+
datas = ["xxx", "bbb", "ccc"]
91+
dedup.filter_exist_data(datas)
92+
assert datas == ["ccc"]
93+
```
94+
95+
## Dedup参数
96+
97+
- **filter_type**:去重类型,支持BloomFilter、MemoryFilter、ExpireFilter三种
98+
- **redis_url**不是必须传递的,若项目中存在setting.py文件,且已配置redis连接方式,则可以不传递redis_url
99+
100+
![-w294](http://markdown-media.oss-cn-beijing.aliyuncs.com/2021/03/07/16151133801599.jpg?x-oss-process=style/markdown-media)
101+
102+
```
103+
import feapder
104+
from feapder.dedup import Dedup
105+
106+
class TestSpider(feapder.Spider):
107+
def __init__(self, *args, **kwargs):
108+
self.dedup = Dedup() # 默认是永久去重
109+
```
110+
111+
- **name**: 过滤器名称 该名称会默认以dedup作为前缀 `dedup:expire_set:[name]`或`dedup:bloomfilter:[name]`。 默认ExpireFilter name=过期时间,BloomFilter name=`dedup:bloomfilter:bloomfilter`
112+
113+
![-w499](http://markdown-media.oss-cn-beijing.aliyuncs.com/2021/03/07/16151136442498.jpg?x-oss-process=style/markdown-media)
114+
115+
若对不同数据源去重,可通过name参数来指定不同去重库
116+
117+
- **absolute_name**:过滤器绝对名称 不会加dedup前缀
118+
- **expire_time**:ExpireFilter的过期时间 单位为秒,其他两种过滤器不用指定
119+
- **error_rate**:BloomFilter/MemoryFilter的误判率 默认为0.00001
120+
- **to_md5**:去重前是否将数据转为MD5,默认是

feapder/VERSION

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1 @@
1-
1.1.7
1+
1.1.8

feapder/db/redisdb.py

Lines changed: 30 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -15,18 +15,18 @@
1515
from feapder.utils.log import log
1616

1717

18-
class Singleton(object):
19-
def __init__(self, cls):
20-
self._cls = cls
21-
self._instance = {}
22-
23-
def __call__(self, *args, **kwargs):
24-
if self._cls not in self._instance:
25-
self._instance[self._cls] = self._cls(*args, **kwargs)
26-
return self._instance[self._cls]
27-
28-
29-
@Singleton
18+
# class Singleton(object):
19+
# def __init__(self, cls):
20+
# self._cls = cls
21+
# self._instance = {}
22+
#
23+
# def __call__(self, *args, **kwargs):
24+
# if self._cls not in self._instance:
25+
# self._instance[self._cls] = self._cls(*args, **kwargs)
26+
# return self._instance[self._cls]
27+
#
28+
#
29+
# @Singleton
3030
class RedisDB:
3131
def __init__(
3232
self,
@@ -64,6 +64,9 @@ def __init__(
6464

6565
try:
6666
if not url:
67+
if not ip_ports:
68+
raise Exception("未设置redis连接信息")
69+
6770
ip_ports = (
6871
ip_ports if isinstance(ip_ports, list) else ip_ports.split(",")
6972
)
@@ -74,7 +77,7 @@ def __init__(
7477
startup_nodes.append({"host": ip, "port": port})
7578

7679
if service_name:
77-
log.debug("使用redis哨兵模式")
80+
# log.debug("使用redis哨兵模式")
7881
hosts = [(node["host"], node["port"]) for node in startup_nodes]
7982
sentinel = Sentinel(hosts, socket_timeout=3, **kwargs)
8083
self._redis = sentinel.master_for(
@@ -88,7 +91,7 @@ def __init__(
8891
)
8992

9093
else:
91-
log.debug("使用redis集群模式")
94+
# log.debug("使用redis集群模式")
9295
self._redis = RedisCluster(
9396
startup_nodes=startup_nodes,
9497
decode_responses=decode_responses,
@@ -117,10 +120,11 @@ def __init__(
117120
except Exception as e:
118121
raise
119122
else:
120-
if not url:
121-
log.debug("连接到redis数据库 %s db%s" % (ip_ports, db))
122-
else:
123-
log.debug("连接到redis数据库 %s" % (url))
123+
# if not url:
124+
# log.debug("连接到redis数据库 %s db%s" % (ip_ports, db))
125+
# else:
126+
# log.debug("连接到redis数据库 %s" % (url))
127+
pass
124128

125129
self._ip_ports = ip_ports
126130
self._db = db
@@ -137,6 +141,14 @@ def __repr__(self):
137141

138142
@classmethod
139143
def from_url(cls, url):
144+
"""
145+
146+
Args:
147+
url: redis://[[username]:[password]]@localhost:6379/0
148+
149+
Returns:
150+
151+
"""
140152
return cls(url=url)
141153

142154
def sadd(self, table, values):

feapder/dedup/README.md

Lines changed: 72 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -1,65 +1,93 @@
1+
# Dedup
12

2-
# 大数据去重
3+
Dedup是feapder大数据去重模块,内置3种去重机制,使用方式一致,可容纳的去重数据量与内存有关。不同于BloomFilter,去重受槽位数量影响,Dedup使用了弹性的去重机制,可容纳海量的数据去重。
34

4-
## 功能
55

6-
1. 基于redis临时去重:
7-
指定数据时效性,时效性之外的历史数据不参与去重
8-
2. 基于内存去重:
9-
使用可扩展的bloomfilter方式,适用于程序运行到结束生命周期内去重
10-
3. 基于redis永久去重
11-
使用可扩展的bloomfilter方式,永久去重海量数据
12-
4. 支持批量去重,输入列表数据,返回列表结果(如[0,1] 0不存在 1 已存在)
6+
## 去重方式
137

14-
## 使用方法
8+
### 临时去重
159

16-
### 临时去重
10+
> 基于redis,支持批量,去重有时效性。去重一万条数据约0.26秒,一亿条数据占用内存约1.43G
11+
12+
```
13+
from feapder.dedup import Dedup
14+
15+
data = {"xxx": 123, "xxxx": "xxxx"}
16+
datas = ["xxx", "bbb"]
17+
18+
def test_ExpireFilter():
19+
dedup = Dedup(
20+
Dedup.ExpireFilter, expire_time=10, redis_url="redis://@localhost:6379/0"
21+
)
22+
23+
# 逐条去重
24+
assert dedup.add(data) == 1
25+
assert dedup.get(data) == 1
26+
27+
# 批量去重
28+
assert dedup.add(datas) == [1, 1]
29+
assert dedup.get(datas) == [1, 1]
30+
```
1731

18-
> 支持批量。速度快,一万条数据约0.26秒。 去重1亿条数据占用内存约1.43G,不适合永久去重
1932

20-
from spider.dedup import Dedup
21-
22-
datas = {
23-
"xxx": xxx,
24-
"xxxx": "xxxx",
25-
}
26-
27-
dedup = Dedup('test', 3) # 表名为test 历史数据3秒有效期
28-
29-
print(dedup) # <ExpireSet: dedup:expire_set:test>
30-
print(dedup.add(datas)) # 0 不存在
31-
print(dedup.get(datas)) # 1 存在
32-
3333
### 内存去重
3434

35-
> 支持批量。一万条数据约0.5秒。 去重一亿条数据占用内存约285MB
36-
37-
from spider.dedup import Dedup
35+
> 基于内存,支持批量。去重一万条数据约0.5秒,一亿条数据占用内存约285MB
36+
37+
```
38+
from feapder.dedup import Dedup
39+
40+
data = {"xxx": 123, "xxxx": "xxxx"}
41+
datas = ["xxx", "bbb"]
42+
43+
def test_MemoryFilter():
44+
dedup = Dedup(Dedup.MemoryFilter) # 表名为test 历史数据3秒有效期
45+
46+
# 逐条去重
47+
assert dedup.add(data) == 1
48+
assert dedup.get(data) == 1
49+
50+
# 批量去重
51+
assert dedup.add(datas) == [1, 1]
52+
assert dedup.get(datas) == [1, 1]
53+
```
3854

39-
datas = {
40-
"xxx": xxx,
41-
"xxxx": "xxxx",
42-
}
43-
44-
dedup = Dedup(use_memory=True)
45-
46-
print(dedup) # <ScalableBloomFilter: MemoryBitArray: 2396264597> (2396264597 为位数组的大小)
47-
print(dedup.add(datas)) # 0 不存在
48-
print(dedup.get(datas)) # 1 存在
49-
5055
### 永久去重
5156

52-
> 支持批量。 一万条数据约3.5秒。去重一亿条数据占用内存约285MB
57+
> 基于redis,支持批量,永久去重。 去重一万条数据约3.5秒,一亿条数据占用内存约285MB
5358
54-
from spider.dedup import Dedup
59+
from feapder.dedup import Dedup
5560

5661
datas = {
5762
"xxx": xxx,
5863
"xxxx": "xxxx",
5964
}
60-
65+
6166
dedup = Dedup()
62-
67+
6368
print(dedup) # <ScalableBloomFilter: RedisBitArray: dedup:bloomfilter:bloomfilter>
6469
print(dedup.add(datas)) # 0 不存在
65-
print(dedup.get(datas)) # 1 存在
70+
print(dedup.get(datas)) # 1 存在
71+
72+
## 过滤数据
73+
74+
Dedup可以通过如下方法,过滤掉已存在的数据
75+
76+
77+
```python
78+
from feapder.dedup import Dedup
79+
80+
def test_filter():
81+
dedup = Dedup(Dedup.BloomFilter, redis_url="redis://@localhost:6379/0")
82+
83+
# 制造已存在数据
84+
datas = ["xxx", "bbb"]
85+
dedup.add(datas)
86+
87+
# 过滤掉已存在数据 "xxx", "bbb"
88+
datas = ["xxx", "bbb", "ccc"]
89+
dedup.filter_exist_data(datas)
90+
assert datas == ["ccc"]
91+
```
92+
93+

0 commit comments

Comments
 (0)