From 47bd9c0e96a66169df6dffead54612f6265bf9dd Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Sun, 14 Jul 2019 22:31:02 +0500 Subject: [PATCH] move serialization/deserialization logic to downstream queues --- scrapy/core/scheduler.py | 8 +++----- scrapy/pqueues.py | 33 +++++---------------------------- scrapy/squeues.py | 28 ++++++++++++++++++++++------ 3 files changed, 30 insertions(+), 39 deletions(-) diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 975aede0c..473cff92e 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -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}) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 6ecd1b51a..454609470 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -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): diff --git a/scrapy/squeues.py b/scrapy/squeues.py index 30cc926e5..9f78e17cc 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -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