Skip to content

Commit 601f3fc

Browse files
author
Boris
committed
删除无用的配置
1 parent 8a01432 commit 601f3fc

24 files changed

Lines changed: 132 additions & 113 deletions

File tree

docs/_sidebar.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@
2222
* [Spider进阶](source_code/Spider进阶.md)
2323
* [BatchSpider进阶](source_code/BatchSpider进阶.md)
2424
* [配置文件](source_code/配置文件.md)
25-
* [数据管道-pipline](source_code/pipline.md)
25+
* [数据管道-pipeline](source_code/pipeline.md)
2626
* [Item](source_code/Item.md)
2727
* [UpdateItem](source_code/UpdateItem.md)
2828
* [MysqlDB](source_code/MysqlDB.md)

docs/source_code/pipline.md

Lines changed: 13 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,21 +1,21 @@
1-
# Pipline
1+
# Pipeline
22

3-
Pipline是数据入库时流经的管道,默认为使用mysql入库,用户可自定义。
3+
Pipeline是数据入库时流经的管道,默认为使用mysql入库,用户可自定义。
44

55
注:AirSpider不支持
66

77
## 使用方式
88

9-
### 1. 编写pipline
9+
### 1. 编写pipeline
1010

1111
```python
12-
from feapder.piplines import BasePipline
12+
from feapder.pipelines import BasePipeline
1313
from typing import Dict, List, Tuple
1414

1515

16-
class Pipline(BasePipline):
16+
class Pipeline(BasePipeline):
1717
"""
18-
pipline 是单线程的,批量保存数据的操作,不建议在这里写网络请求代码,如下载图片等
18+
pipeline 是单线程的,批量保存数据的操作,不建议在这里写网络请求代码,如下载图片等
1919
"""
2020

2121
def save_items(self, table, items: List[Dict]) -> bool:
@@ -30,7 +30,7 @@ class Pipline(BasePipline):
3030
3131
"""
3232

33-
print("自定义pipline, 保存数据 >>>>", table, items)
33+
print("自定义pipeline, 保存数据 >>>>", table, items)
3434

3535
return True
3636

@@ -47,23 +47,23 @@ class Pipline(BasePipline):
4747
4848
"""
4949

50-
print("自定义pipline, 更新数据 >>>>", table, items, update_keys)
50+
print("自定义pipeline, 更新数据 >>>>", table, items, update_keys)
5151

5252
return True
5353
```
5454

55-
`Pipline`需继承`BasePipline`,类名和存放位置随意,需要实现`save_items``update_items`两个接口。一定要有返回值,返回`False`表示数据没保存成功,数据不入去重库,以便再次入库
55+
`Pipeline`需继承`BasePipeline`,类名和存放位置随意,需要实现`save_items``update_items`两个接口。一定要有返回值,返回`False`表示数据没保存成功,数据不入去重库,以便再次入库
5656

5757
### 2. 编写配置文件
5858

5959
```python
60-
# 数据入库的pipline,可自定义,默认MysqlPipline
61-
ITEM_PIPLINES = [
62-
"pipline.Pipline"
60+
# 数据入库的pipeline,可自定义,默认MysqlPipeline
61+
ITEM_PIPELINES = [
62+
"pipeline.Pipeline"
6363
]
6464
```
6565

66-
将编写好的pipline配置进来,值为类的模块路径,需要指定到具体的类名
66+
将编写好的pipeline配置进来,值为类的模块路径,需要指定到具体的类名
6767

6868
## 示例
6969

docs/source_code/配置文件.md

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -64,10 +64,6 @@ WARNING_FAILED_COUNT = 1000 # 任务失败数 超过WARNING_FAILED_COUNT则报
6464
# 爬虫初始化工作
6565
# 爬虫是否自动结束,若为False,则会等待新任务下发,进程不退出
6666
AUTO_STOP_WHEN_SPIDER_DONE = True
67-
# 是否将item添加到 mysql 支持列表; 可指定添加的item,支持模糊指定,如ADD_ITEM_TO_MYSQL=["_task"], 这样只有表名包含_task的表才会入库
68-
ADD_ITEM_TO_MYSQL = True
69-
# 是否将item添加到 redis 支持列表; 可指定添加的item,支持模糊指定,如ADD_ITEM_TO_REDIS=["_task"], 这样只有表名包含_task的表才会入库
70-
ADD_ITEM_TO_REDIS = False
7167

7268
# 设置代理
7369
PROXY_EXTRACT_API = None # 代理提取API ,返回的代理分割符为\r\n

feapder/buffer/item_buffer.py

Lines changed: 29 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -17,13 +17,13 @@
1717
from feapder.db.redisdb import RedisDB
1818
from feapder.dedup import Dedup
1919
from feapder.network.item import Item, UpdateItem
20-
from feapder.piplines import BasePipline
20+
from feapder.pipelines import BasePipeline
2121
from feapder.utils.log import log
2222

2323
MAX_ITEM_COUNT = 5000 # 缓存中最大item数
2424
UPLOAD_BATCH_MAX_SIZE = 1000
2525

26-
MYSQL_PIPLINE_PATH = "feapder.piplines.mysql_pipline.MysqlPipline"
26+
MYSQL_PIPELINE_PATH = "feapder.pipelines.mysql_pipeline.MysqlPipeline"
2727

2828

2929
class Singleton(object):
@@ -59,34 +59,34 @@ def __init__(self, redis_key):
5959
# 'xxx:xxx_item': ['id', 'name'...] # 记录redis中item名与需要更新的key对应关系
6060
}
6161

62-
self._piplines = self.load_piplines()
62+
self._pipelines = self.load_pipelines()
6363

64-
self._have_mysql_pipline = MYSQL_PIPLINE_PATH in setting.ITEM_PIPLINES
65-
self._mysql_pipline = None
64+
self._have_mysql_pipeline = MYSQL_PIPELINE_PATH in setting.ITEM_PIPELINES
65+
self._mysql_pipeline = None
6666

6767
if setting.ITEM_FILTER_ENABLE and not self.__class__.dedup:
6868
self.__class__.dedup = Dedup(to_md5=False)
6969

70-
def load_piplines(self):
71-
piplines = []
72-
for pipline_path in setting.ITEM_PIPLINES:
73-
module, class_name = pipline_path.rsplit(".", 1)
74-
pipline_cls = importlib.import_module(module).__getattribute__(class_name)
75-
pipline = pipline_cls()
76-
if not isinstance(pipline, BasePipline):
77-
raise ValueError(f"{pipline_path} 需继承 feapder.piplines.BasePipline")
78-
piplines.append(pipline)
70+
def load_pipelines(self):
71+
pipelines = []
72+
for pipeline_path in setting.ITEM_PIPELINES:
73+
module, class_name = pipeline_path.rsplit(".", 1)
74+
pipeline_cls = importlib.import_module(module).__getattribute__(class_name)
75+
pipeline = pipeline_cls()
76+
if not isinstance(pipeline, BasePipeline):
77+
raise ValueError(f"{pipeline_path} 需继承 feapder.pipelines.BasePipeline")
78+
pipelines.append(pipeline)
7979

80-
return piplines
80+
return pipelines
8181

8282
@property
83-
def mysql_pipline(self):
84-
if not self._mysql_pipline:
85-
module, class_name = MYSQL_PIPLINE_PATH.rsplit(".", 1)
86-
pipline_cls = importlib.import_module(module).__getattribute__(class_name)
87-
self._mysql_pipline = pipline_cls()
83+
def mysql_pipeline(self):
84+
if not self._mysql_pipeline:
85+
module, class_name = MYSQL_PIPELINE_PATH.rsplit(".", 1)
86+
pipeline_cls = importlib.import_module(module).__getattribute__(class_name)
87+
self._mysql_pipeline = pipeline_cls()
8888

89-
return self._mysql_pipline
89+
return self._mysql_pipeline
9090

9191
def run(self):
9292
while not self._thread_stop:
@@ -244,24 +244,24 @@ def __export_to_db(self, tab_item, datas, is_update=False, update_keys=()):
244244
# 打点 校验
245245
self.check_datas(table=to_table, datas=datas)
246246

247-
for pipline in self._piplines:
247+
for pipeline in self._pipelines:
248248
if is_update:
249-
if not pipline.update_items(to_table, datas, update_keys=update_keys):
249+
if not pipeline.update_items(to_table, datas, update_keys=update_keys):
250250
log.error(
251-
f"{pipline.__class__.__name__} 更新数据失败. table: {to_table} items: {datas}"
251+
f"{pipeline.__class__.__name__} 更新数据失败. table: {to_table} items: {datas}"
252252
)
253253
return False
254254

255255
else:
256-
if not pipline.save_items(to_table, datas):
256+
if not pipeline.save_items(to_table, datas):
257257
log.error(
258-
f"{pipline.__class__.__name__} 保存数据失败. table: {to_table} items: {datas}"
258+
f"{pipeline.__class__.__name__} 保存数据失败. table: {to_table} items: {datas}"
259259
)
260260
return False
261261

262-
# 若是任务表, 且上面的pipline里没mysql,则需调用mysql更新任务
263-
if not self._have_mysql_pipline and is_update and to_table.endswith("_task"):
264-
self.mysql_pipline.update_items(to_table, datas, update_keys=update_keys)
262+
# 若是任务表, 且上面的pipeline里没mysql,则需调用mysql更新任务
263+
if not self._have_mysql_pipeline and is_update and to_table.endswith("_task"):
264+
self.mysql_pipeline.update_items(to_table, datas, update_keys=update_keys)
265265

266266
def __add_item_to_db(
267267
self, items, update_items, requests, callbacks, items_fingerprints

feapder/core/spiders/batch_spider.py

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,9 @@
2828
from feapder.utils.log import log
2929
from feapder.utils.redis_lock import RedisLock
3030

31+
CONSOLE_PIPELINE_PATH = "feapder.pipelines.console_pipeline.ConsolePipeline"
32+
MYSQL_PIPELINE_PATH = "feapder.pipelines.mysql_pipeline.MysqlPipeline"
33+
3134

3235
class BatchSpider(BatchParser, Scheduler):
3336
def __init__(
@@ -1040,7 +1043,6 @@ class DebugBatchSpider(BatchSpider):
10401043
SPIDER_TASK_COUNT=1,
10411044
SPIDER_MAX_RETRY_TIMES=10,
10421045
REQUEST_TIME_OUT=600, # 10分钟
1043-
ADD_ITEM_TO_MYSQL=False,
10441046
PROXY_ENABLE=False,
10451047
RETRY_FAILED_REQUESTS=False,
10461048
# 保存失败的request
@@ -1050,6 +1052,7 @@ class DebugBatchSpider(BatchSpider):
10501052
REQUEST_FILTER_ENABLE=False,
10511053
OSS_UPLOAD_TABLES=(),
10521054
DELETE_KEYS=True,
1055+
ITEM_PIPELINES=[CONSOLE_PIPELINE_PATH],
10531056
)
10541057

10551058
def __init__(
@@ -1077,8 +1080,10 @@ def __init__(
10771080
raise Exception("task_id 与 task 不能同时为null")
10781081

10791082
kwargs["redis_key"] = kwargs["redis_key"] + "_debug"
1080-
if save_to_db:
1081-
self.__class__.__debug_custom_setting__.update(ADD_ITEM_TO_MYSQL=True)
1083+
if save_to_db and not self.__class__.__custom_setting__.get("ITEM_PIPELINES"):
1084+
self.__class__.__debug_custom_setting__.update(
1085+
ITEM_PIPELINES=[MYSQL_PIPELINE_PATH]
1086+
)
10821087
self.__class__.__custom_setting__.update(
10831088
self.__class__.__debug_custom_setting__
10841089
)

feapder/core/spiders/spider.py

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,8 @@
2121
from feapder.network.request import Request
2222
from feapder.utils.log import log
2323

24+
CONSOLE_PIPELINE_PATH = "feapder.pipelines.console_pipeline.ConsolePipeline"
25+
2426

2527
class Spider(
2628
BaseParser, Scheduler
@@ -43,7 +45,7 @@ def __init__(
4345
auto_start_requests=None,
4446
send_run_time=False,
4547
batch_interval=0,
46-
wait_lock=True
48+
wait_lock=True,
4749
):
4850
"""
4951
@summary: 爬虫
@@ -95,9 +97,7 @@ def start_monitor_task(self, *args, **kws):
9597
while True:
9698
try:
9799
# 检查redis中是否有任务
98-
tab_requests = setting.TAB_REQUSETS.format(
99-
redis_key=self._redis_key
100-
)
100+
tab_requests = setting.TAB_REQUSETS.format(redis_key=self._redis_key)
101101
todo_task_count = redisdb.zget_count(tab_requests)
102102

103103
if todo_task_count < self._min_task_count: # 添加任务
@@ -236,7 +236,6 @@ class DebugSpider(Spider):
236236
SPIDER_TASK_COUNT=1,
237237
SPIDER_MAX_RETRY_TIMES=10,
238238
REQUEST_TIME_OUT=600, # 10分钟
239-
ADD_ITEM_TO_MYSQL=False,
240239
PROXY_ENABLE=False,
241240
RETRY_FAILED_REQUESTS=False,
242241
# 保存失败的request
@@ -246,6 +245,7 @@ class DebugSpider(Spider):
246245
REQUEST_FILTER_ENABLE=False,
247246
OSS_UPLOAD_TABLES=(),
248247
DELETE_KEYS=True,
248+
ITEM_PIPELINES=[CONSOLE_PIPELINE_PATH],
249249
)
250250

251251
def __init__(self, request=None, request_dict=None, *args, **kwargs):

feapder/db/redisdb.py

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -164,7 +164,7 @@ def sadd(self, table, values):
164164
if isinstance(values, list):
165165
pipe = self._redis.pipeline(
166166
transaction=True
167-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
167+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
168168

169169
if not self._is_redis_cluster:
170170
pipe.multi()
@@ -190,7 +190,7 @@ def sget(self, table, count=1, is_pop=True):
190190
if count > 1:
191191
pipe = self._redis.pipeline(
192192
transaction=True
193-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
193+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
194194

195195
if not self._is_redis_cluster:
196196
pipe.multi()
@@ -219,7 +219,7 @@ def srem(self, table, values):
219219
if isinstance(values, list):
220220
pipe = self._redis.pipeline(
221221
transaction=True
222-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
222+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
223223

224224
if not self._is_redis_cluster:
225225
pipe.multi()
@@ -302,7 +302,7 @@ def zget(self, table, count=1, is_pop=True):
302302

303303
pipe = self._redis.pipeline(
304304
transaction=True
305-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
305+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
306306

307307
if not self._is_redis_cluster:
308308
pipe.multi() # 标记事务的开始 参考 http://www.runoob.com/redis/redis-transactions.html
@@ -521,7 +521,7 @@ def zexists(self, table, values):
521521
if isinstance(values, list):
522522
pipe = self._redis.pipeline(
523523
transaction=True
524-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
524+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
525525
pipe.multi()
526526
for value in values:
527527
pipe.zscore(table, value)
@@ -542,7 +542,7 @@ def lpush(self, table, values):
542542
if isinstance(values, list):
543543
pipe = self._redis.pipeline(
544544
transaction=True
545-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
545+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
546546

547547
if not self._is_redis_cluster:
548548
pipe.multi()
@@ -570,7 +570,7 @@ def lpop(self, table, count=1):
570570
if count > 1:
571571
pipe = self._redis.pipeline(
572572
transaction=True
573-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
573+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
574574

575575
if not self._is_redis_cluster:
576576
pipe.multi()
@@ -712,7 +712,7 @@ def setbit(self, table, offsets, values):
712712

713713
pipe = self._redis.pipeline(
714714
transaction=True
715-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
715+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
716716
pipe.multi()
717717

718718
for offset, value in zip(offsets, values):
@@ -733,7 +733,7 @@ def getbit(self, table, offsets):
733733
if isinstance(offsets, list):
734734
pipe = self._redis.pipeline(
735735
transaction=True
736-
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipline实现一次请求指定多个命令,并且默认情况下一次pipline 是原子性操作。
736+
) # redis-py默认在执行每次请求都会创建(连接池申请连接)和断开(归还连接池)一次连接操作,如果想要在一次请求中指定多个命令,则可以使用pipeline实现一次请求指定多个命令,并且默认情况下一次pipeline 是原子性操作。
737737
pipe.multi()
738738
for offset in offsets:
739739
pipe.getbit(table, offset)
Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,9 +12,9 @@
1212
from typing import Dict, List, Tuple
1313

1414

15-
class BasePipline(metaclass=abc.ABCMeta):
15+
class BasePipeline(metaclass=abc.ABCMeta):
1616
"""
17-
pipline 是单线程的,批量保存数据的操作,不建议在这里写网络请求代码,如下载图片等
17+
pipeline 是单线程的,批量保存数据的操作,不建议在这里写网络请求代码,如下载图片等
1818
"""
1919

2020
@abc.abstractmethod

0 commit comments

Comments
 (0)