move serialization/deserialization logic to downstream queues

This commit is contained in:
Vostretsov Nikita 2019-07-14 22:31:02 +05:00
parent fa6a0d799b
commit 47bd9c0e96
3 changed files with 30 additions and 39 deletions

View File

@ -148,12 +148,11 @@ class Scheduler(object):
def _newdq(self, priority):
""" Factory for creating disk queues. """
path = join(self.dqdir, 'p%s' % (priority, ))
return self.dqclass(path)
return self.dqclass(self.crawler, path)
def _mq(self):
""" Create a new priority queue instance, with in-memory storage """
return create_instance(self.pqclass, None, self.crawler, self._newmq,
serialize=False)
return create_instance(self.pqclass, None, self.crawler, self._newmq)
def _dq(self):
""" Create a new priority queue instance, with disk storage """
@ -162,8 +161,7 @@ class Scheduler(object):
None,
self.crawler,
self._newdq,
state,
serialize=True)
state)
if q:
logger.info("Resuming crawl (%(queuesize)d requests scheduled)",
{'queuesize': len(q)}, extra={'spider': self.spider})

View File

@ -91,29 +91,7 @@ class _SlotPriorityQueues(object):
class ScrapyPriorityQueue(PriorityQueue):
"""
PriorityQueue which works with scrapy.Request instances and
can optionally convert them to/from dicts before/after putting to a queue.
"""
def __init__(self, crawler, qfactory, startprios=(), serialize=False):
super(ScrapyPriorityQueue, self).__init__(qfactory, startprios)
self.serialize = serialize
self.spider = crawler.spider
@classmethod
def from_crawler(cls, crawler, qfactory, startprios=(), serialize=False):
return cls(crawler, qfactory, startprios, serialize)
def push(self, request, priority=0):
if self.serialize:
request = request_to_dict(request, self.spider)
super(ScrapyPriorityQueue, self).push(request, priority)
def pop(self):
request = super(ScrapyPriorityQueue, self).pop()
if request and self.serialize:
request = request_from_dict(request, self.spider)
return request
pass
class DownloaderInterface(object):
@ -142,10 +120,10 @@ class DownloaderAwarePriorityQueue(object):
"""
@classmethod
def from_crawler(cls, crawler, qfactory, slot_startprios=None, serialize=False):
return cls(crawler, qfactory, slot_startprios, serialize)
def from_crawler(cls, crawler, qfactory, slot_startprios=None):
return cls(crawler, qfactory, slot_startprios)
def __init__(self, crawler, qfactory, slot_startprios=None, serialize=False):
def __init__(self, crawler, qfactory, slot_startprios=None):
if crawler.settings.getint('CONCURRENT_REQUESTS_PER_IP') != 0:
raise ValueError('"%s" does not support CONCURRENT_REQUESTS_PER_IP'
% (self.__class__,))
@ -164,9 +142,8 @@ class DownloaderAwarePriorityQueue(object):
for slot, startprios in (slot_startprios or {}).items()}
def pqfactory(startprios=()):
return ScrapyPriorityQueue(crawler, qfactory, startprios, serialize)
return ScrapyPriorityQueue(qfactory, startprios)
self._slot_pqueues = _SlotPriorityQueues(pqfactory, slot_startprios)
self.serialize = serialize
self._downloader_interface = DownloaderInterface(crawler)
def pop(self):

View File

@ -7,19 +7,35 @@ from six.moves import cPickle as pickle
from queuelib import queue
from scrapy.utils.reqser import request_to_dict, request_from_dict
def _serializable_queue(queue_class, serialize, deserialize):
class SerializableQueue(queue_class):
def push(self, obj):
s = serialize(obj)
super(SerializableQueue, self).push(s)
def __init__(self, crawler, key):
self.spider = crawler.spider
super(SerializableQueue, self).__init__(key)
def push(self, request):
if serialize:
request = request_to_dict(request, self.spider)
request = serialize(request)
return super(SerializableQueue, self).push(request)
def pop(self):
s = super(SerializableQueue, self).pop()
if s:
return deserialize(s)
request = super(SerializableQueue, self).pop()
if not request:
return None
if deserialize:
request = deserialize(request)
request = request_from_dict(request, self.spider)
return request
return SerializableQueue