Skip to content

Commit be57e55

Browse files
committed
Merge branch 'develop'
2 parents 0ec612d + ca59283 commit be57e55

18 files changed

Lines changed: 322 additions & 97 deletions

feapder/VERSION

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1 @@
1-
1.6.1
1+
1.6.2-beta1

feapder/buffer/item_buffer.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
from feapder.pipelines import BasePipeline
2121
from feapder.pipelines.mysql_pipeline import MysqlPipeline
2222
from feapder.utils.log import log
23+
from feapder.utils import metrics
2324

2425
MAX_ITEM_COUNT = 5000 # 缓存中最大item数
2526
UPLOAD_BATCH_MAX_SIZE = 1000
@@ -90,6 +91,7 @@ def mysql_pipeline(self):
9091
return self._mysql_pipeline
9192

9293
def run(self):
94+
self._thread_stop = False
9395
while not self._thread_stop:
9496
self.flush()
9597
tools.delay_time(0.5)
@@ -351,7 +353,9 @@ def check_datas(self, table, datas):
351353
@param datas: 数据 列表
352354
@return:
353355
"""
354-
pass
356+
for data in datas:
357+
for k, v in data.items():
358+
metrics.emit_counter(k, int(bool(v)), classify=table)
355359

356360
def close(self):
357361
pass

feapder/buffer/request_buffer.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ def __init__(self, redis_key):
4848
) # 过期时间为一个月
4949

5050
def run(self):
51+
self._thread_stop = False
5152
while not self._thread_stop:
5253
try:
5354
self.__add_request_to_db()

feapder/core/collector.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ def __init__(self, redis_key):
4949
self.__delete_dead_node()
5050

5151
def run(self):
52+
self._thread_stop = False
5253
while not self._thread_stop:
5354
try:
5455
self.__report_node_heartbeat()

feapder/core/parser_control.py

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,14 +18,15 @@
1818
from feapder.db.memory_db import MemoryDB
1919
from feapder.network.item import Item
2020
from feapder.network.request import Request
21+
from feapder.utils import metrics
2122
from feapder.utils.log import log
2223

2324

2425
class PaserControl(threading.Thread):
2526
DOWNLOAD_EXCEPTION = "download_exception"
2627
DOWNLOAD_SUCCESS = "download_success"
2728
DOWNLOAD_TOTAL = "download_total"
28-
PAESERS_EXCEPTION = "parsers_exception"
29+
PAESERS_EXCEPTION = "parser_exception"
2930

3031
is_show_tip = False
3132

@@ -46,6 +47,7 @@ def __init__(self, collector, redis_key, request_buffer, item_buffer):
4647
self._wait_task_time = 0
4748

4849
def run(self):
50+
self._thread_stop = False
4951
while not self._thread_stop:
5052
try:
5153
requests = self._collector.get_requests(setting.SPIDER_TASK_COUNT)
@@ -413,7 +415,8 @@ def record_download_status(self, status, spider):
413415
记录html等文档下载状态
414416
@return:
415417
"""
416-
pass
418+
419+
metrics.emit_counter(f"{spider}:{status}", 1, classify="document")
417420

418421
def stop(self):
419422
self._thread_stop = True

feapder/core/scheduler.py

Lines changed: 30 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
from feapder.network.request import Request
2525
from feapder.utils.log import log
2626
from feapder.utils.redis_lock import RedisLock
27+
from feapder.utils import metrics
2728

2829
SPIDER_START_TIME_KEY = "spider_start_time"
2930
SPIDER_END_TIME_KEY = "spider_end_time"
@@ -144,6 +145,14 @@ def __init__(
144145
self._last_check_task_status_time = 0
145146
self.wait_lock = wait_lock
146147

148+
self.init_metrics()
149+
150+
def init_metrics(self):
151+
"""
152+
初始化打点系统
153+
"""
154+
metrics.init(**setting.METRICS_OTHER_ARGS)
155+
147156
def add_parser(self, parser):
148157
parser = parser() # parser 实例化
149158
if isinstance(parser, BaseParser):
@@ -473,19 +482,26 @@ def spider_begin(self):
473482
# 发送消息
474483
self.send_msg("《%s》爬虫开始" % self._spider_name)
475484

476-
def spider_end(self):
485+
def spider_end(self, close=True):
477486
self.record_end_time()
478487

479488
if self._end_callback:
480489
self._end_callback()
481490

482491
for parser in self._parsers:
483-
parser.close()
492+
if close:
493+
parser.close()
484494
parser.end_callback()
485495

486-
# 关闭webdirver
487-
if Request.webdriver_pool:
488-
Request.webdriver_pool.close()
496+
if close:
497+
# 关闭webdirver
498+
if Request.webdriver_pool:
499+
Request.webdriver_pool.close()
500+
501+
# 关闭打点
502+
metrics.close()
503+
else:
504+
metrics.flush()
489505

490506
# 计算抓取时长
491507
data = self._redisdb.hget(
@@ -552,3 +568,12 @@ def record_spider_state(
552568
batch_interval=None,
553569
):
554570
pass
571+
572+
def join(self, timeout=None):
573+
"""
574+
重写线程的join
575+
"""
576+
if not self._started.is_set():
577+
return
578+
579+
super().join()

feapder/core/spiders/air_spider.py

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@
1818
from feapder.db.memory_db import MemoryDB
1919
from feapder.network.request import Request
2020
from feapder.utils.log import log
21+
from feapder.utils import metrics
2122

2223

2324
class AirSpider(BaseParser, Thread):
@@ -41,6 +42,8 @@ def __init__(self, thread_count=None):
4142
self._parser_controls = []
4243
self._item_buffer = ItemBuffer(redis_key="air_spider")
4344

45+
metrics.init(**setting.METRICS_OTHER_ARGS)
46+
4447
def distribute_task(self):
4548
for request in self.start_requests():
4649
if not isinstance(request, Request):
@@ -101,4 +104,15 @@ def run(self):
101104
break
102105

103106
self.end_callback()
107+
# 为了线程可重复start
104108
self._started.clear()
109+
metrics.close()
110+
111+
def join(self, timeout=None):
112+
"""
113+
重写线程的join
114+
"""
115+
if not self._started.is_set():
116+
return
117+
118+
super().join()

feapder/core/spiders/batch_spider.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1011,7 +1011,7 @@ def run(self):
10111011
self.task_is_done() and self.all_thread_is_done()
10121012
): # redis全部的任务已经做完 并且mysql中的任务已经做完(检查各个线程all_thread_is_done,防止任务没做完,就更新任务状态,导致程序结束的情况)
10131013
if not self._is_notify_end:
1014-
self.spider_end()
1014+
self.spider_end(close=self._auto_stop_when_spider_done)
10151015
self.record_spider_state(
10161016
spider_type=2,
10171017
state=1,

feapder/core/spiders/spider.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -193,7 +193,7 @@ def run(self):
193193
while True:
194194
if self.all_thread_is_done():
195195
if not self._is_notify_end:
196-
self.spider_end() # 跑完一轮
196+
self.spider_end(close=self._auto_stop_when_spider_done) # 跑完一轮
197197
self.record_spider_state(
198198
spider_type=1,
199199
state=1,

feapder/requirements.txt

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,4 +13,5 @@ bitarray>=1.5.3
1313
redis-py-cluster>=2.1.0
1414
cryptography>=3.3.2
1515
urllib3>=1.25.8
16-
loguru>=0.5.3
16+
loguru>=0.5.3
17+
influxdb>=5.3.1

0 commit comments

Comments
 (0)