Skip to content

Commit 1b00eb3

Browse files
committed
数据入库失败 自动重试
1 parent 21d1737 commit 1b00eb3

7 files changed

Lines changed: 137 additions & 71 deletions

File tree

feapder/buffer/item_buffer.py

Lines changed: 56 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -62,6 +62,9 @@ def __init__(self, redis_key, task_table=None):
6262
if setting.ITEM_FILTER_ENABLE and not self.__class__.dedup:
6363
self.__class__.dedup = Dedup(to_md5=False)
6464

65+
# 导出失败的次数
66+
self.export_falied_times = 0
67+
6568
@property
6669
def redis_db(self):
6770
if self.__class__.__redis_db is None:
@@ -327,22 +330,59 @@ def __add_item_to_db(
327330
tab_item, datas, is_update=True, update_keys=update_keys
328331
)
329332

330-
# 执行回调
331-
while callbacks:
332-
try:
333-
callback = callbacks.pop(0)
334-
callback()
335-
except Exception as e:
336-
log.exception(e)
337-
338-
# 删除做过的request
339-
if requests:
340-
self.redis_db.zrem(self._table_request, requests)
341-
342-
# 去重入库
343-
if export_success and setting.ITEM_FILTER_ENABLE:
344-
if items_fingerprints:
345-
self.__class__.dedup.add(items_fingerprints, skip_check=True)
333+
if export_success:
334+
# 执行回调
335+
while callbacks:
336+
try:
337+
callback = callbacks.pop(0)
338+
callback()
339+
except Exception as e:
340+
log.exception(e)
341+
342+
# 删除做过的request
343+
if requests:
344+
self.redis_db.zrem(self._table_request, requests)
345+
346+
# 去重入库
347+
if setting.ITEM_FILTER_ENABLE:
348+
if items_fingerprints:
349+
self.__class__.dedup.add(items_fingerprints, skip_check=True)
350+
else:
351+
if self.export_falied_times > setting.EXPORT_DATA_MAX_FAILED_TIMES:
352+
# 报警
353+
msg = "《{}》爬虫导出数据失败,失败次数:{},请检查爬虫是否正常".format(
354+
self._redis_key, self.export_falied_times
355+
)
356+
log.error(msg)
357+
tools.send_msg(
358+
msg=msg,
359+
level="error",
360+
message_prefix="《%s》爬虫导出数据失败" % (self._redis_key),
361+
)
362+
363+
if self.export_falied_times > setting.EXPORT_DATA_MAX_RETRY_TIMES:
364+
# 删除做过的request
365+
if requests:
366+
self.redis_db.zrem(self._table_request, requests)
367+
log.error("入库超过最大重试次数,不再重试")
368+
else:
369+
tip = ["入库不成功"]
370+
if callbacks:
371+
tip.append("不执行回调")
372+
if requests:
373+
tip.append("不删除任务")
374+
exists = self.redis_db.zexists(self._table_request, requests)
375+
for exist, request in zip(exists, requests):
376+
if exist:
377+
self.redis_db.zadd(self._table_request, requests, 300)
378+
379+
if setting.ITEM_FILTER_ENABLE:
380+
tip.append("数据不入去重库")
381+
382+
tip.append("将自动重试 (AirSpider不支持)")
383+
log.error(",".join(tip))
384+
385+
self.export_falied_times += 1
346386

347387
self._is_adding_to_db = False
348388

feapder/core/scheduler.py

Lines changed: 15 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,6 @@ def __init__(
4343
delete_keys=(),
4444
keep_alive=None,
4545
auto_start_requests=None,
46-
send_run_time=True,
4746
batch_interval=0,
4847
wait_lock=True,
4948
task_table=None,
@@ -59,7 +58,6 @@ def __init__(
5958
@param delete_keys: 爬虫启动时删除的key,类型: 元组/bool/string。 支持正则
6059
@param keep_alive: 爬虫是否常驻,默认否
6160
@param auto_start_requests: 爬虫是否自动添加任务
62-
@param send_run_time: 发送运行时间
6361
@param batch_interval: 抓取时间间隔 默认为0 天为单位 多次启动时,只有当前时间与第一次抓取结束的时间间隔大于指定的时间间隔时,爬虫才启动
6462
@param wait_lock: 下发任务时否等待锁,若不等待锁,可能会存在多进程同时在下发一样的任务,因此分布式环境下请将该值设置True
6563
@param task_table: 任务表, 批次爬虫传递
@@ -105,7 +103,6 @@ def __init__(
105103
if auto_start_requests is not None
106104
else setting.SPIDER_AUTO_START_REQUESTS
107105
)
108-
self._send_run_time = send_run_time
109106
self._batch_interval = batch_interval
110107

111108
self._begin_callback = (
@@ -394,10 +391,10 @@ def check_task_status(self):
394391
self.send_msg(
395392
msg,
396393
level="error",
397-
message_prefix="《%s》爬虫当前失败任务数预警" % (self._spider_name),
394+
message_prefix="《%s》爬虫当前失败任务数报警" % (self._spider_name),
398395
)
399396

400-
# parser_control实时统计已做任务数及失败任务数,若失败数大于10且失败任务数/已做任务数>=0.5 则报警
397+
# parser_control实时统计已做任务数及失败任务数,若成功率<0.5 则报警
401398
failed_task_count, success_task_count = PaserControl.get_task_status_count()
402399
total_count = success_task_count + failed_task_count
403400
if total_count > 0:
@@ -411,13 +408,22 @@ def check_task_status(self):
411408
task_success_rate,
412409
)
413410
log.error(msg)
414-
# 统计下上次发消息的时间,若时间大于1小时,则报警(此处为多进程,需要考虑别报重复)
415411
self.send_msg(
416412
msg,
417413
level="error",
418-
message_prefix="《%s》爬虫当前任务成功率" % (self._spider_name),
414+
message_prefix="《%s》爬虫当前任务成功率报警" % (self._spider_name),
419415
)
420416

417+
# 检查入库失败次数
418+
if self._item_buffer.export_falied_times > setting.EXPORT_DATA_MAX_FAILED_TIMES:
419+
msg = "《{}》爬虫导出数据失败,失败次数:{}, 请检查爬虫是否正常".format(
420+
self._spider_name, self._item_buffer.export_falied_times
421+
)
422+
log.error(msg)
423+
self.send_msg(
424+
msg, level="error", message_prefix="《%s》爬虫导出数据失败" % (self._spider_name)
425+
)
426+
421427
def delete_tables(self, delete_tables_list):
422428
if isinstance(delete_tables_list, bool):
423429
delete_tables_list = [self._redis_key + "*"]
@@ -447,22 +453,7 @@ def _stop_all_thread(self):
447453

448454
def send_msg(self, msg, level="debug", message_prefix=""):
449455
# log.debug("发送报警 level:{} msg{}".format(level, msg))
450-
if setting.WARNING_LEVEL == "ERROR":
451-
if level != "error":
452-
return
453-
454-
if setting.DINGDING_WARNING_URL:
455-
keyword = "feapder报警系统\n"
456-
tools.dingding_warning(keyword + msg, message_prefix=message_prefix)
457-
458-
if setting.EMAIL_RECEIVER:
459-
tools.email_warning(
460-
msg, message_prefix=message_prefix, title=self._spider_name
461-
)
462-
463-
if setting.WECHAT_WARNING_URL:
464-
keyword = "feapder报警系统\n"
465-
tools.wechat_warning(keyword + msg, message_prefix=message_prefix)
456+
tools.send_msg(msg=msg, level=level, message_prefix=message_prefix)
466457

467458
def spider_begin(self):
468459
"""
@@ -524,8 +515,7 @@ def spider_end(self):
524515
)
525516
log.info(msg)
526517

527-
if self._send_run_time:
528-
self.send_msg(msg)
518+
self.send_msg(msg)
529519

530520
if self._keep_alive:
531521
log.info("爬虫不自动结束, 等待下一轮任务...")

feapder/core/spiders/batch_spider.py

Lines changed: 26 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -54,8 +54,7 @@ def __init__(
5454
end_callback=None,
5555
delete_keys=(),
5656
keep_alive=None,
57-
send_run_time=False,
58-
**kwargs
57+
**kwargs,
5958
):
6059
"""
6160
@summary: 批次爬虫
@@ -90,7 +89,6 @@ def __init__(
9089
@param end_callback: 爬虫结束回调函数
9190
@param delete_keys: 爬虫启动时删除的key,类型: 元组/bool/string。 支持正则; 常用于清空任务队列,否则重启时会断点续爬
9291
@param keep_alive: 爬虫是否常驻,默认否
93-
@param send_run_time: 发送运行时间
9492
@param related_redis_key: 有关联的其他爬虫任务表(redis)注意:要避免环路 如 A -> B & B -> A 。
9593
@param related_batch_record: 有关联的其他爬虫批次表(mysql)注意:要避免环路 如 A -> B & B -> A 。
9694
related_redis_key 与 related_batch_record 选其一配置即可;用于相关联的爬虫没结束时,本爬虫也不结束
@@ -110,10 +108,9 @@ def __init__(
110108
delete_keys=delete_keys,
111109
keep_alive=keep_alive,
112110
auto_start_requests=False,
113-
send_run_time=send_run_time,
114111
batch_interval=batch_interval,
115112
task_table=task_table,
116-
**kwargs
113+
**kwargs,
117114
)
118115

119116
self._redisdb = RedisDB()
@@ -659,7 +656,15 @@ def check_batch(self, is_first_check=False):
659656
>= self._send_msg_interval
660657
):
661658
self._last_send_msg_time = now_date
662-
self.send_msg(msg, level="error")
659+
self.send_msg(
660+
msg,
661+
level="error",
662+
message_prefix="《{}》本批次未完成, 正在等待依赖爬虫 {} 结束".format(
663+
self._batch_name,
664+
self._related_batch_record
665+
or self._related_task_tables,
666+
),
667+
)
663668

664669
return False
665670

@@ -768,7 +773,11 @@ def check_batch(self, is_first_check=False):
768773
>= self._send_msg_interval
769774
):
770775
self._last_send_msg_time = now_date
771-
self.send_msg(msg, level="error")
776+
self.send_msg(
777+
msg,
778+
level="error",
779+
message_prefix="《{}》批次超时".format(self._batch_name),
780+
)
772781

773782
else: # 未超时
774783
remaining_time = (
@@ -823,7 +832,13 @@ def check_batch(self, is_first_check=False):
823832
>= self._send_msg_interval
824833
):
825834
self._last_send_msg_time = now_date
826-
self.send_msg(msg, level="error")
835+
self.send_msg(
836+
msg,
837+
level="error",
838+
message_prefix="《{}》批次可能超时".format(
839+
self._batch_name
840+
),
841+
)
827842

828843
elif overflow_time < 0:
829844
msg += ", 该批次预计提前 {} 完成".format(
@@ -1036,7 +1051,9 @@ def run(self):
10361051
except Exception as e:
10371052
msg = "《%s》主线程异常 爬虫结束 exception: %s" % (self._batch_name, e)
10381053
log.error(msg)
1039-
self.send_msg(msg, level="error")
1054+
self.send_msg(
1055+
msg, level="error", message_prefix="《%s》爬虫异常结束".format(self._batch_name)
1056+
)
10401057

10411058
os._exit(137) # 使退出码为35072 方便爬虫管理器重启
10421059

feapder/core/spiders/spider.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,6 @@ def __init__(
4343
delete_keys=(),
4444
keep_alive=None,
4545
auto_start_requests=None,
46-
send_run_time=False,
4746
batch_interval=0,
4847
wait_lock=True,
4948
**kwargs
@@ -60,7 +59,6 @@ def __init__(
6059
@param delete_keys: 爬虫启动时删除的key,类型: 元组/bool/string。 支持正则; 常用于清空任务队列,否则重启时会断点续爬
6160
@param keep_alive: 爬虫是否常驻
6261
@param auto_start_requests: 爬虫是否自动添加任务
63-
@param send_run_time: 发送运行时间
6462
@param batch_interval: 抓取时间间隔 默认为0 天为单位 多次启动时,只有当前时间与第一次抓取结束的时间间隔大于指定的时间间隔时,爬虫才启动
6563
@param wait_lock: 下发任务时否等待锁,若不等待锁,可能会存在多进程同时在下发一样的任务,因此分布式环境下请将该值设置True
6664
---------
@@ -74,7 +72,6 @@ def __init__(
7472
delete_keys=delete_keys,
7573
keep_alive=keep_alive,
7674
auto_start_requests=auto_start_requests,
77-
send_run_time=send_run_time,
7875
batch_interval=batch_interval,
7976
wait_lock=wait_lock,
8077
**kwargs

feapder/setting.py

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,6 @@
3232
# ip:port 多个可写为列表或者逗号隔开 如 ip1:port1,ip2:port2 或 ["ip1:port1", "ip2:port2"]
3333
REDISDB_IP_PORTS = os.getenv("REDISDB_IP_PORTS")
3434
REDISDB_USER_PASS = os.getenv("REDISDB_USER_PASS")
35-
# 默认 0 到 15 共16个数据库
3635
REDISDB_DB = int(os.getenv("REDISDB_DB", 0))
3736
# 适用于redis哨兵模式
3837
REDISDB_SERVICE_NAME = os.getenv("REDISDB_SERVICE_NAME")
@@ -42,6 +41,8 @@
4241
"feapder.pipelines.mysql_pipeline.MysqlPipeline",
4342
# "feapder.pipelines.mongo_pipeline.MongoPipeline",
4443
]
44+
EXPORT_DATA_MAX_FAILED_TIMES = 10 # 导出数据时最大的失败次数,包括保存和更新,超过这个次数报警
45+
EXPORT_DATA_MAX_RETRY_TIMES = 10 # 导出数据时最大的重试次数,包括保存和更新,超过这个次数则放弃重试
4546

4647
# 爬虫相关
4748
# COLLECTOR
@@ -100,7 +101,7 @@
100101

101102
# 随机headers
102103
RANDOM_HEADERS = True
103-
# UserAgent类型 支持 'chrome', 'opera', 'firefox', 'internetexplorer', 'safari',若不指定则随机类型
104+
# UserAgent类型 支持 'chrome', 'opera', 'firefox', 'internetexplorer', 'safari','mobile' 若不指定则随机类型
104105
USER_AGENT_TYPE = "chrome"
105106
# 默认使用的浏览器头 RANDOM_HEADERS=True时不生效
106107
DEFAULT_USERAGENT = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_14_2) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/73.0.3683.103 Safari/537.36"

feapder/templates/project_template/setting.py

Lines changed: 17 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -4,33 +4,34 @@
44
# import sys
55
#
66
# # MYSQL
7-
# MYSQL_IP = os.getenv("MYSQL_IP")
8-
# MYSQL_PORT = int(os.getenv("MYSQL_PORT", 3306))
9-
# MYSQL_DB = os.getenv("MYSQL_DB")
10-
# MYSQL_USER_NAME = os.getenv("MYSQL_USER_NAME")
11-
# MYSQL_USER_PASS = os.getenv("MYSQL_USER_PASS")
7+
# MYSQL_IP = "localhost"
8+
# MYSQL_PORT = 3306
9+
# MYSQL_DB = ""
10+
# MYSQL_USER_NAME = ""
11+
# MYSQL_USER_PASS = ""
1212
#
1313
# # MONGODB
14-
# MONGO_IP = os.getenv("MONGO_IP", "localhost")
15-
# MONGO_PORT = int(os.getenv("MONGO_PORT", 27017))
16-
# MONGO_DB = os.getenv("MONGO_DB")
17-
# MONGO_USER_NAME = os.getenv("MONGO_USER_NAME")
18-
# MONGO_USER_PASS = os.getenv("MONGO_USER_PASS")
14+
# MONGO_IP = "localhost"
15+
# MONGO_PORT = 27017
16+
# MONGO_DB = ""
17+
# MONGO_USER_NAME = ""
18+
# MONGO_USER_PASS = ""
1919
#
2020
# # REDIS
2121
# # ip:port 多个可写为列表或者逗号隔开 如 ip1:port1,ip2:port2 或 ["ip1:port1", "ip2:port2"]
22-
# REDISDB_IP_PORTS = os.getenv("REDISDB_IP_PORTS")
23-
# REDISDB_USER_PASS = os.getenv("REDISDB_USER_PASS")
24-
# # 默认 0 到 15 共16个数据库
25-
# REDISDB_DB = int(os.getenv("REDISDB_DB", 0))
22+
# REDISDB_IP_PORTS = "localhost:6379"
23+
# REDISDB_USER_PASS = ""
24+
# REDISDB_DB = 0
2625
# # 适用于redis哨兵模式
27-
# REDISDB_SERVICE_NAME = os.getenv("REDISDB_SERVICE_NAME")
26+
# REDISDB_SERVICE_NAME = ""
2827
#
2928
# # 数据入库的pipeline,可自定义,默认MysqlPipeline
3029
# ITEM_PIPELINES = [
3130
# "feapder.pipelines.mysql_pipeline.MysqlPipeline",
3231
# # "feapder.pipelines.mongo_pipeline.MongoPipeline",
3332
# ]
33+
# EXPORT_DATA_MAX_FAILED_TIMES = 10 # 导出数据时最大的失败次数,包括保存和更新,超过这个次数报警
34+
# EXPORT_DATA_MAX_RETRY_TIMES = 10 # 导出数据时最大的重试次数,包括保存和更新,超过这个次数则放弃重试
3435
#
3536
# # 爬虫相关
3637
# # COLLECTOR
@@ -79,7 +80,7 @@
7980
#
8081
# # 随机headers
8182
# RANDOM_HEADERS = True
82-
# # UserAgent类型 支持 'chrome', 'opera', 'firefox', 'internetexplorer', 'safari',若不指定则随机类型
83+
# # UserAgent类型 支持 'chrome', 'opera', 'firefox', 'internetexplorer', 'safari','mobile' 若不指定则随机类型
8384
# USER_AGENT_TYPE = "chrome"
8485
# # 默认使用的浏览器头 RANDOM_HEADERS=True时不生效
8586
# DEFAULT_USERAGENT = "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_14_2) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/73.0.3683.103 Safari/537.36"

0 commit comments

Comments
 (0)