-
Notifications
You must be signed in to change notification settings - Fork 546
Expand file tree
/
Copy pathtask_spider.py
More file actions
733 lines (631 loc) · 27.2 KB
/
Copy pathtask_spider.py
File metadata and controls
733 lines (631 loc) · 27.2 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
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
# -*- coding: utf-8 -*-
"""
Created on 2020/4/22 12:06 AM
---------
@summary:
---------
@author: Boris
@email: boris_liu@foxmail.com
"""
import os
import time
import warnings
from collections.abc import Iterable
from typing import List, Tuple, Dict, Union
import feapder.setting as setting
import feapder.utils.tools as tools
from feapder.core.base_parser import TaskParser
from feapder.core.scheduler import Scheduler
from feapder.db.mysqldb import MysqlDB
from feapder.db.redisdb import RedisDB
from feapder.network.item import Item
from feapder.network.item import UpdateItem
from feapder.network.request import Request
from feapder.utils.log import log
from feapder.utils.perfect_dict import PerfectDict
CONSOLE_PIPELINE_PATH = "feapder.pipelines.console_pipeline.ConsolePipeline"
class TaskSpider(TaskParser, Scheduler):
def __init__(
self,
redis_key,
task_table,
task_table_type="mysql",
task_keys=None,
task_state="state",
min_task_count=10000,
check_task_interval=5,
task_limit=10000,
related_redis_key=None,
related_batch_record=None,
task_condition="",
task_order_by="",
thread_count=None,
begin_callback=None,
end_callback=None,
delete_keys=(),
keep_alive=None,
batch_interval=0,
use_mysql=True,
**kwargs,
):
"""
@summary: 任务爬虫
必要条件 需要指定任务表,可以是redis表或者mysql表作为任务种子
redis任务种子表:zset类型。值为 {"xxx":xxx, "xxx2":"xxx2"};若为集成模式,需指定parser_name字段,如{"xxx":xxx, "xxx2":"xxx2", "parser_name":"TestTaskSpider"}
mysql任务表:
任务表中必须有id及任务状态字段 如 state, 其他字段可根据爬虫需要的参数自行扩充。若为集成模式,需指定parser_name字段。
参考建表语句如下:
CREATE TABLE `table_name` (
`id` int(11) NOT NULL AUTO_INCREMENT,
`param` varchar(1000) DEFAULT NULL COMMENT '爬虫需要的抓取数据需要的参数',
`state` int(11) DEFAULT NULL COMMENT '任务状态',
`parser_name` varchar(255) DEFAULT NULL COMMENT '任务解析器的脚本类名',
PRIMARY KEY (`id`),
UNIQUE KEY `nui` (`param`) USING BTREE
) ENGINE=InnoDB AUTO_INCREMENT=1 DEFAULT CHARSET=utf8;
---------
@param task_table: mysql中的任务表 或 redis中存放任务种子的key,zset类型
@param task_table_type: 任务表类型 支持 redis 、mysql
@param task_keys: 需要获取的任务字段 列表 [] 如需指定解析的parser,则需将parser_name字段取出来。
@param task_state: mysql中任务表的任务状态字段
@param min_task_count: redis 中最少任务数, 少于这个数量会从种子表中取任务
@param check_task_interval: 检查是否还有任务的时间间隔;
@param task_limit: 每次从数据库中取任务的数量
@param redis_key: 任务等数据存放在redis中的key前缀
@param thread_count: 线程数,默认为配置文件中的线程数
@param begin_callback: 爬虫开始回调函数
@param end_callback: 爬虫结束回调函数
@param delete_keys: 爬虫启动时删除的key,类型: 元组/bool/string。 支持正则; 常用于清空任务队列,否则重启时会断点续爬
@param keep_alive: 爬虫是否常驻,默认否
@param related_redis_key: 有关联的其他爬虫任务表(redis)注意:要避免环路 如 A -> B & B -> A 。
@param related_batch_record: 有关联的其他爬虫批次表(mysql)注意:要避免环路 如 A -> B & B -> A 。
related_redis_key 与 related_batch_record 选其一配置即可;用于相关联的爬虫没结束时,本爬虫也不结束
若相关连的爬虫为批次爬虫,推荐以related_batch_record配置,
若相关连的爬虫为普通爬虫,无批次表,可以以related_redis_key配置
@param task_condition: 任务条件 用于从一个大任务表中挑选出数据自己爬虫的任务,即where后的条件语句
@param task_order_by: 取任务时的排序条件 如 id desc
@param batch_interval: 抓取时间间隔 默认为0 天为单位 多次启动时,只有当前时间与第一次抓取结束的时间间隔大于指定的时间间隔时,爬虫才启动
@param use_mysql: 是否使用mysql数据库
---------
@result:
"""
Scheduler.__init__(
self,
redis_key=redis_key,
thread_count=thread_count,
begin_callback=begin_callback,
end_callback=end_callback,
delete_keys=delete_keys,
keep_alive=keep_alive,
auto_start_requests=False,
batch_interval=batch_interval,
task_table=task_table,
**kwargs,
)
self._redisdb = RedisDB()
self._mysqldb = MysqlDB() if use_mysql else None
self._task_table = task_table # mysql中的任务表
self._task_keys = task_keys # 需要获取的任务字段
self._task_table_type = task_table_type
if self._task_table_type == "mysql" and not self._task_keys:
raise Exception("需指定任务字段 使用task_keys")
self._task_state = task_state # mysql中任务表的state字段名
self._min_task_count = min_task_count # redis 中最少任务数
self._check_task_interval = check_task_interval
self._task_limit = task_limit # mysql中一次取的任务数量
self._related_task_tables = [
setting.TAB_REQUESTS.format(redis_key=redis_key)
] # 自己的task表也需要检查是否有任务
if related_redis_key:
self._related_task_tables.append(
setting.TAB_REQUESTS.format(redis_key=related_redis_key)
)
self._related_batch_record = related_batch_record
self._task_condition = task_condition
self._task_condition_prefix_and = task_condition and " and {}".format(
task_condition
)
self._task_condition_prefix_where = task_condition and " where {}".format(
task_condition
)
self._task_order_by = task_order_by and " order by {}".format(task_order_by)
self._is_more_parsers = True # 多模版类爬虫
self.reset_task()
def add_parser(self, parser, **kwargs):
parser = parser(
self._task_table,
self._task_state,
self._mysqldb,
**kwargs,
) # parser 实例化
self._parsers.append(parser)
def start_monitor_task(self):
"""
@summary: 监控任务状态
---------
---------
@result:
"""
if not self._parsers: # 不是多模版模式, 将自己注入到parsers,自己为模版
self._is_more_parsers = False
self._parsers.append(self)
elif len(self._parsers) <= 1:
self._is_more_parsers = False
# 添加任务
for parser in self._parsers:
parser.add_task()
while True:
try:
# 检查redis中是否有任务 任务小于_min_task_count 则从mysql中取
tab_requests = setting.TAB_REQUESTS.format(redis_key=self._redis_key)
todo_task_count = self._redisdb.zget_count(tab_requests)
tasks = []
if todo_task_count < self._min_task_count:
tasks = self.get_task(todo_task_count)
if not tasks:
if not todo_task_count:
if self._keep_alive:
log.info("任务均已做完,爬虫常驻, 等待新任务")
time.sleep(self._check_task_interval)
continue
elif self.have_alive_spider():
log.info("任务均已做完,但还有爬虫在运行,等待爬虫结束")
time.sleep(self._check_task_interval)
continue
elif not self.related_spider_is_done():
continue
else:
log.info("任务均已做完,爬虫结束")
break
else:
log.info("redis 中尚有%s条积压任务,暂时不派发新任务" % todo_task_count)
if not tasks:
if todo_task_count >= self._min_task_count:
# log.info('任务正在进行 redis中剩余任务 %s' % todo_task_count)
pass
else:
log.info("无待做种子 redis中剩余任务 %s" % todo_task_count)
else:
# make start requests
self.distribute_task(tasks)
log.info(f"添加任务到redis成功 共{len(tasks)}条")
except Exception as e:
log.exception(e)
time.sleep(self._check_task_interval)
def get_task(self, todo_task_count) -> List[Union[Tuple, Dict]]:
"""
获取任务
Args:
todo_task_count: redis里剩余的任务数
Returns:
"""
tasks = []
if self._task_table_type == "mysql":
# 从mysql中取任务
log.info("redis 中剩余任务%s 数量过小 从mysql中取任务追加" % todo_task_count)
tasks = self.get_todo_task_from_mysql()
if not tasks: # 状态为0的任务已经做完,需要检查状态为2的任务是否丢失
# redis 中无待做任务,此时mysql中状态为2的任务为丢失任务。需重新做
if todo_task_count == 0:
log.info("无待做任务,尝试取丢失的任务")
tasks = self.get_doing_task_from_mysql()
elif self._task_table_type == "redis":
log.info("redis 中剩余任务%s 数量过小 从redis种子任务表中取任务追加" % todo_task_count)
tasks = self.get_task_from_redis()
else:
raise Exception(
f"task_table_type expect mysql or redis,bug got {self._task_table_type}"
)
return tasks
def distribute_task(self, tasks):
"""
@summary: 分发任务
---------
@param tasks:
---------
@result:
"""
if self._is_more_parsers: # 为多模版类爬虫,需要下发指定的parser
for task in tasks:
for parser in self._parsers: # 寻找task对应的parser
if parser.name in task:
if isinstance(task, dict):
task = PerfectDict(_dict=task)
else:
task = PerfectDict(
_dict=dict(zip(self._task_keys, task)),
_values=list(task),
)
requests = parser.start_requests(task)
if requests and not isinstance(requests, Iterable):
raise Exception(
"%s.%s返回值必须可迭代" % (parser.name, "start_requests")
)
result_type = 1
for request in requests or []:
if isinstance(request, Request):
request.parser_name = request.parser_name or parser.name
self._request_buffer.put_request(request)
result_type = 1
elif isinstance(request, Item):
self._item_buffer.put_item(request)
result_type = 2
if (
self._item_buffer.get_items_count()
>= setting.ITEM_MAX_CACHED_COUNT
):
self._item_buffer.flush()
elif callable(request): # callbale的request可能是更新数据库操作的函数
if result_type == 1:
self._request_buffer.put_request(request)
else:
self._item_buffer.put_item(request)
if (
self._item_buffer.get_items_count()
>= setting.ITEM_MAX_CACHED_COUNT
):
self._item_buffer.flush()
else:
raise TypeError(
"start_requests yield result type error, expect Request、Item、callback func, bug get type: {}".format(
type(requests)
)
)
break
else: # task没对应的parser 则将task下发到所有的parser
for task in tasks:
for parser in self._parsers:
if isinstance(task, dict):
task = PerfectDict(_dict=task)
else:
task = PerfectDict(
_dict=dict(zip(self._task_keys, task)), _values=list(task)
)
requests = parser.start_requests(task)
if requests and not isinstance(requests, Iterable):
raise Exception(
"%s.%s返回值必须可迭代" % (parser.name, "start_requests")
)
result_type = 1
for request in requests or []:
if isinstance(request, Request):
request.parser_name = request.parser_name or parser.name
self._request_buffer.put_request(request)
result_type = 1
elif isinstance(request, Item):
self._item_buffer.put_item(request)
result_type = 2
if (
self._item_buffer.get_items_count()
>= setting.ITEM_MAX_CACHED_COUNT
):
self._item_buffer.flush()
elif callable(request): # callbale的request可能是更新数据库操作的函数
if result_type == 1:
self._request_buffer.put_request(request)
else:
self._item_buffer.put_item(request)
if (
self._item_buffer.get_items_count()
>= setting.ITEM_MAX_CACHED_COUNT
):
self._item_buffer.flush()
self._request_buffer.flush()
self._item_buffer.flush()
def get_task_from_redis(self):
tasks = self._redisdb.zget(self._task_table, count=self._task_limit)
tasks = [eval(task) for task in tasks]
return tasks
def get_todo_task_from_mysql(self):
"""
@summary: 取待做的任务
---------
---------
@result:
"""
# TODO 分批取数据 每批最大取 1000000个,防止内存占用过大
# 查询任务
task_keys = ", ".join([f"`{key}`" for key in self._task_keys])
sql = "select %s from %s where %s = 0%s%s limit %s" % (
task_keys,
self._task_table,
self._task_state,
self._task_condition_prefix_and,
self._task_order_by,
self._task_limit,
)
tasks = self._mysqldb.find(sql)
if tasks:
# 更新任务状态
for i in range(0, len(tasks), 10000): # 10000 一批量更新
task_ids = str(
tuple([task[0] for task in tasks[i : i + 10000]])
).replace(",)", ")")
sql = "update %s set %s = 2 where id in %s" % (
self._task_table,
self._task_state,
task_ids,
)
self._mysqldb.update(sql)
return tasks
def get_doing_task_from_mysql(self):
"""
@summary: 取正在做的任务
---------
---------
@result:
"""
# 查询任务
task_keys = ", ".join([f"`{key}`" for key in self._task_keys])
sql = "select %s from %s where %s = 2%s%s limit %s" % (
task_keys,
self._task_table,
self._task_state,
self._task_condition_prefix_and,
self._task_order_by,
self._task_limit,
)
tasks = self._mysqldb.find(sql)
return tasks
def get_lose_task_count(self):
sql = "select count(1) from %s where %s = 2%s" % (
self._task_table,
self._task_state,
self._task_condition_prefix_and,
)
doing_count = self._mysqldb.find(sql)[0][0]
return doing_count
def reset_lose_task_from_mysql(self):
"""
@summary: 重置丢失任务为待做
---------
---------
@result:
"""
sql = "update {table} set {state} = 0 where {state} = 2{task_condition}".format(
table=self._task_table,
state=self._task_state,
task_condition=self._task_condition_prefix_and,
)
return self._mysqldb.update(sql)
def related_spider_is_done(self):
"""
相关连的爬虫是否跑完
@return: True / False / None 表示无相关的爬虫 可由自身的total_count 和 done_count 来判断
"""
for related_redis_task_table in self._related_task_tables:
if self._redisdb.exists_key(related_redis_task_table):
log.info(f"依赖的爬虫还未结束,任务表为:{related_redis_task_table}")
return False
if self._related_batch_record:
sql = "select is_done from {} order by id desc limit 1".format(
self._related_batch_record
)
is_done = self._mysqldb.find(sql)
is_done = is_done[0][0] if is_done else None
if is_done is None:
log.warning("相关联的批次表不存在或无批次信息")
return True
if not is_done:
log.info(f"依赖的爬虫还未结束,批次表为:{self._related_batch_record}")
return False
return True
# -------- 批次结束逻辑 ------------
def task_is_done(self):
"""
@summary: 检查种子表是否做完
---------
---------
@result: True / False (做完 / 未做完)
"""
is_done = False
if self._task_table_type == "mysql":
sql = "select 1 from %s where (%s = 0 or %s=2)%s limit 1" % (
self._task_table,
self._task_state,
self._task_state,
self._task_condition_prefix_and,
)
count = self._mysqldb.find(sql) # [(1,)] / []
elif self._task_table_type == "redis":
count = self._redisdb.zget_count(self._task_table)
else:
raise Exception(
f"task_table_type expect mysql or redis,bug got {self._task_table_type}"
)
if not count:
log.info("种子表中任务均已完成")
is_done = True
return is_done
def run(self):
"""
@summary: 重写run方法 检查mysql中的任务是否做完, 做完停止
---------
---------
@result:
"""
try:
if not self.is_reach_next_spider_time():
return
if not self._parsers: # 不是add_parser 模式
self._parsers.append(self)
self._start()
while True:
try:
if self._stop_spider or (
self.all_thread_is_done()
and self.task_is_done()
and self.related_spider_is_done()
): # redis全部的任务已经做完 并且mysql中的任务已经做完(检查各个线程all_thread_is_done,防止任务没做完,就更新任务状态,导致程序结束的情况)
if not self._is_notify_end:
self.spider_end()
self._is_notify_end = True
if not self._keep_alive:
self._stop_all_thread()
break
else:
log.info("常驻爬虫,等待新任务")
else:
self._is_notify_end = False
self.check_task_status()
except Exception as e:
log.exception(e)
tools.delay_time(10) # 10秒钟检查一次爬虫状态
except Exception as e:
msg = "《%s》主线程异常 爬虫结束 exception: %s" % (self.name, e)
log.error(msg)
self.send_msg(
msg, level="error", message_prefix="《%s》爬虫异常结束".format(self.name)
)
os._exit(137) # 使退出码为35072 方便爬虫管理器重启
@classmethod
def to_DebugTaskSpider(cls, *args, **kwargs):
# DebugBatchSpider 继承 cls
DebugTaskSpider.__bases__ = (cls,)
DebugTaskSpider.__name__ = cls.__name__
return DebugTaskSpider(*args, **kwargs)
class DebugTaskSpider(TaskSpider):
"""
Debug批次爬虫
"""
__debug_custom_setting__ = dict(
COLLECTOR_TASK_COUNT=1,
# SPIDER
SPIDER_THREAD_COUNT=1,
SPIDER_SLEEP_TIME=0,
SPIDER_MAX_RETRY_TIMES=10,
REQUEST_LOST_TIMEOUT=600, # 10分钟
PROXY_ENABLE=False,
RETRY_FAILED_REQUESTS=False,
# 保存失败的request
SAVE_FAILED_REQUEST=False,
# 过滤
ITEM_FILTER_ENABLE=False,
REQUEST_FILTER_ENABLE=False,
OSS_UPLOAD_TABLES=(),
DELETE_KEYS=True,
)
def __init__(
self,
task_id=None,
task=None,
save_to_db=False,
update_task=False,
*args,
**kwargs,
):
"""
@param task_id: 任务id
@param task: 任务 task 与 task_id 二者选一即可。如 task = {"url":""}
@param save_to_db: 数据是否入库 默认否
@param update_task: 是否更新任务 默认否
@param args:
@param kwargs:
"""
warnings.warn(
"您正处于debug模式下,该模式下不会更新任务状态及数据入库,仅用于调试。正式发布前请更改为正常模式", category=Warning
)
if not task and not task_id:
raise Exception("task_id 与 task 不能同时为空")
kwargs["redis_key"] = kwargs["redis_key"] + "_debug"
if not save_to_db:
self.__class__.__debug_custom_setting__["ITEM_PIPELINES"] = [
CONSOLE_PIPELINE_PATH
]
self.__class__.__custom_setting__.update(
self.__class__.__debug_custom_setting__
)
super(DebugTaskSpider, self).__init__(*args, **kwargs)
self._task_id = task_id
self._task = task
self._update_task = update_task
def start_monitor_task(self):
"""
@summary: 监控任务状态
---------
---------
@result:
"""
if not self._parsers: # 不是多模版模式, 将自己注入到parsers,自己为模版
self._is_more_parsers = False
self._parsers.append(self)
elif len(self._parsers) <= 1:
self._is_more_parsers = False
if self._task:
self.distribute_task([self._task])
else:
tasks = self.get_todo_task_from_mysql()
if not tasks:
raise Exception("未获取到任务 请检查 task_id: {} 是否存在".format(self._task_id))
self.distribute_task(tasks)
log.debug("下发任务完毕")
def get_todo_task_from_mysql(self):
"""
@summary: 取待做的任务
---------
---------
@result:
"""
# 查询任务
task_keys = ", ".join([f"`{key}`" for key in self._task_keys])
sql = "select %s from %s where id=%s" % (
task_keys,
self._task_table,
self._task_id,
)
tasks = self._mysqldb.find(sql)
return tasks
def save_cached(self, request, response, table):
pass
def update_task_state(self, task_id, state=1, *args, **kwargs):
"""
@summary: 更新任务表中任务状态,做完每个任务时代码逻辑中要主动调用。可能会重写
调用方法为 yield lambda : self.update_task_state(task_id, state)
---------
@param task_id:
@param state:
---------
@result:
"""
if self._update_task:
kwargs["id"] = task_id
kwargs[self._task_state] = state
sql = tools.make_update_sql(
self._task_table,
kwargs,
condition="id = {task_id}".format(task_id=task_id),
)
if self._mysqldb.update(sql):
log.debug("置任务%s状态成功" % task_id)
else:
log.error("置任务%s状态失败 sql=%s" % (task_id, sql))
def update_task_batch(self, task_id, state=1, *args, **kwargs):
"""
批量更新任务 多处调用,更新的字段必须一致
注意:需要 写成 yield update_task_batch(...) 否则不会更新
@param task_id:
@param state:
@param kwargs:
@return:
"""
if self._update_task:
kwargs["id"] = task_id
kwargs[self._task_state] = state
update_item = UpdateItem(**kwargs)
update_item.table_name = self._task_table
update_item.name_underline = self._task_table + "_item"
return update_item
def run(self):
self.start_monitor_task()
if not self._parsers: # 不是add_parser 模式
self._parsers.append(self)
self._start()
while True:
try:
if self.all_thread_is_done():
self._stop_all_thread()
break
except Exception as e:
log.exception(e)
tools.delay_time(1) # 1秒钟检查一次爬虫状态
self.delete_tables([self._redis_key + "*"])