forked from Boris-code/feapder
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathair_spider.py
More file actions
125 lines (96 loc) · 3.69 KB
/
Copy pathair_spider.py
File metadata and controls
125 lines (96 loc) · 3.69 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
# -*- coding: utf-8 -*-
"""
Created on 2020/4/22 12:05 AM
---------
@summary: 基于内存队列的爬虫,不支持分布式
---------
@author: Boris
@email: boris_liu@foxmail.com
"""
from threading import Thread
import feapder.setting as setting
import feapder.utils.tools as tools
from feapder.buffer.item_buffer import ItemBuffer
from feapder.core.base_parser import BaseParser
from feapder.core.parser_control import AirSpiderParserControl
from feapder.db.memory_db import MemoryDB
from feapder.network.request import Request
from feapder.utils.log import log
from feapder.utils import metrics
class AirSpider(BaseParser, Thread):
__custom_setting__ = {}
def __init__(self, thread_count=None):
"""
基于内存队列的爬虫,不支持分布式
:param thread_count: 线程数
"""
super(AirSpider, self).__init__()
for key, value in self.__class__.__custom_setting__.items():
setattr(setting, key, value)
self._thread_count = (
setting.SPIDER_THREAD_COUNT if not thread_count else thread_count
)
self._memory_db = MemoryDB()
self._parser_controls = []
self._item_buffer = ItemBuffer(redis_key="air_spider")
metrics.init(**setting.METRICS_OTHER_ARGS)
def distribute_task(self):
for request in self.start_requests():
if not isinstance(request, Request):
raise ValueError("仅支持 yield Request")
request.parser_name = request.parser_name or self.name
self._memory_db.add(request)
def all_thread_is_done(self):
for i in range(3): # 降低偶然性, 因为各个环节不是并发的,很有可能当时状态为假,但检测下一条时该状态为真。一次检测很有可能遇到这种偶然性
# 检测 parser_control 状态
for parser_control in self._parser_controls:
if not parser_control.is_not_task():
return False
# 检测 任务队列 状态
if not self._memory_db.empty():
return False
# 检测 item_buffer 状态
if (
self._item_buffer.get_items_count() > 0
or self._item_buffer.is_adding_to_db()
):
return False
tools.delay_time(1)
return True
def run(self):
self.start_callback()
for i in range(self._thread_count):
parser_control = AirSpiderParserControl(self._memory_db, self._item_buffer)
parser_control.add_parser(self)
parser_control.start()
self._parser_controls.append(parser_control)
self._item_buffer.start()
self.distribute_task()
while True:
try:
if self.all_thread_is_done():
# 停止 parser_controls
for parser_control in self._parser_controls:
parser_control.stop()
# 关闭item_buffer
self._item_buffer.stop()
# 关闭webdirver
if Request.webdriver_pool:
Request.webdriver_pool.close()
log.info("无任务,爬虫结束")
break
except Exception as e:
log.exception(e)
tools.delay_time(1) # 1秒钟检查一次爬虫状态
self.end_callback()
# 为了线程可重复start
self._started.clear()
# 关闭打点
metrics.close()
def join(self, timeout=None):
"""
重写线程的join
"""
if not self._started.is_set():
return
super().join()