Skip to content

Commit f5f17f9

Browse files
committed
集成打点监控
1 parent 4c083d6 commit f5f17f9

15 files changed

Lines changed: 288 additions & 96 deletions

File tree

feapder/buffer/item_buffer.py

Lines changed: 4 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
@@ -351,7 +352,9 @@ def check_datas(self, table, datas):
351352
@param datas: 数据 列表
352353
@return:
353354
"""
354-
pass
355+
for data in datas:
356+
for k, v in data.items():
357+
metrics.emit_counter(k, int(bool(v)), classify=table)
355358

356359
def close(self):
357360
pass

feapder/core/parser_control.py

Lines changed: 4 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

@@ -413,7 +414,8 @@ def record_download_status(self, status, spider):
413414
记录html等文档下载状态
414415
@return:
415416
"""
416-
pass
417+
418+
metrics.emit_counter(f"{spider}:{status}", 1, classify="document")
417419

418420
def stop(self):
419421
self._thread_stop = True

feapder/core/scheduler.py

Lines changed: 21 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(

feapder/core/spiders/air_spider.py

Lines changed: 4 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):
@@ -102,3 +105,4 @@ def run(self):
102105

103106
self.end_callback()
104107
self._started.clear()
108+
metrics.close()

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

feapder/setting.py

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -115,7 +115,7 @@
115115
# 钉钉报警
116116
DINGDING_WARNING_URL = "" # 钉钉机器人api
117117
DINGDING_WARNING_PHONE = "" # 报警人 支持列表,可指定多个
118-
DINGDING_WARNING_ALL = False # 是否提示所有人, 默认为False
118+
DINGDING_WARNING_ALL = False # 是否提示所有人, 默认为False
119119
# 邮件报警
120120
EMAIL_SENDER = "" # 发件人
121121
EMAIL_PASSWORD = "" # 授权码
@@ -142,6 +142,18 @@
142142
LOG_ENCODING = "utf8" # 日志文件编码
143143
OTHERS_LOG_LEVAL = "ERROR" # 第三方库的log等级
144144

145+
# 打点监控 influxdb 配置
146+
INFLUXDB_HOST = os.getenv("INFLUXDB_HOST", "localhost")
147+
INFLUXDB_PORT = int(os.getenv("INFLUXDB_PORT", 8086))
148+
INFLUXDB_UDP_PORT = int(os.getenv("INFLUXDB_UDP_PORT", 8086))
149+
INFLUXDB_USER = os.getenv("INFLUXDB_USER", "root")
150+
INFLUXDB_PASSWORD = os.getenv("INFLUXDB_PASSWORD", "root")
151+
INFLUXDB_DATABASE = "feapder"
152+
# 监控数据存储的表名,爬虫管理系统上会以task_id命名
153+
INFLUXDB_MEASUREMENT = "task_" + os.getenv("TASK_ID") if os.getenv("TASK_ID") else None
154+
# 打点监控其他参数,若这里也配置了influxdb的参数, 则会覆盖外面的配置
155+
METRICS_OTHER_ARGS = dict(retention_policy_duration="180d", emit_interval=60)
156+
145157
############# 导入用户自定义的setting #############
146158
try:
147159
from setting import *

0 commit comments

Comments
 (0)