diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 975aede0c..e184ed50e 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -119,7 +119,7 @@ class Scheduler(object): if self.dqs is None: return try: - self.dqs.push(request, -request.priority) + self.dqs.push(request) except ValueError as e: # non serializable request if self.logunser: msg = ("Unable to serialize request: %(request)s - reason:" @@ -135,35 +135,29 @@ class Scheduler(object): return True def _mqpush(self, request): - self.mqs.push(request, -request.priority) + self.mqs.push(request) def _dqpop(self): if self.dqs: return self.dqs.pop() - def _newmq(self, priority): - """ Factory for creating memory queues. """ - return self.mqclass() - - def _newdq(self, priority): - """ Factory for creating disk queues. """ - path = join(self.dqdir, 'p%s' % (priority, )) - return self.dqclass(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, + settings=None, + crawler=self.crawler, + downstream_queue_cls=self.mqclass, + key='') def _dq(self): """ Create a new priority queue instance, with disk storage """ state = self._read_dqs_state(self.dqdir) q = create_instance(self.pqclass, - None, - self.crawler, - self._newdq, - state, - serialize=True) + settings=None, + crawler=self.crawler, + downstream_queue_cls=self.dqclass, + key=self.dqdir, + startprios=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 717ed4d27..1afe58dab 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -1,11 +1,7 @@ import hashlib import logging -from collections import namedtuple - -from queuelib import PriorityQueue - -from scrapy.utils.reqser import request_to_dict, request_from_dict +from scrapy.utils.misc import create_instance logger = logging.getLogger(__name__) @@ -29,88 +25,89 @@ def _path_safe(text): return '-'.join([pathable_slot, unique_slot]) -class _Priority(namedtuple("_Priority", ["priority", "slot"])): - """ Slot-specific priority. It is a hack - ``(priority, slot)`` tuple - which can be used instead of int priorities in queues: +class ScrapyPriorityQueue: + """A priority queue implemented using multiple internal queues (typically, + FIFO queues). It uses one internal queue for each priority value. The internal + queue must implement the following methods: + + * push(obj) + * pop() + * close() + * __len__() + + ``__init__`` method of ScrapyPriorityQueue receives a downstream_queue_cls + argument, which is a class used to instantiate a new (internal) queue when + a new priority is allocated. + + Only integer priorities should be used. Lower numbers are higher + priorities. + + startprios is a sequence of priorities to start with. If the queue was + previously closed leaving some priority buckets non-empty, those priorities + should be passed in startprios. - * they are ordered in the same way - order is still by priority value, - min(prios) works; - * str(p) representation is guaranteed to be different when slots - are different - this is important because str(p) is used to create - queue files on disk; - * they have readable str(p) representation which is safe - to use as a file name. """ - __slots__ = () - def __str__(self): - return '%s_%s' % (self.priority, _path_safe(str(self.slot))) + @classmethod + def from_crawler(cls, crawler, downstream_queue_cls, key, startprios=()): + return cls(crawler, downstream_queue_cls, key, startprios) + def __init__(self, crawler, downstream_queue_cls, key, startprios=()): + self.crawler = crawler + self.downstream_queue_cls = downstream_queue_cls + self.key = key + self.queues = {} + self.curprio = None + self.init_prios(startprios) -class _SlotPriorityQueues(object): - """ Container for multiple priority queues. """ - def __init__(self, pqfactory, slot_startprios=None): - """ - ``pqfactory`` is a factory for creating new PriorityQueues. - It must be a function which accepts a single optional ``startprios`` - argument, with a list of priorities to create queues for. + def init_prios(self, startprios): + if not startprios: + return - ``slot_startprios`` is a ``{slot: startprios}`` dict. - """ - self.pqfactory = pqfactory - self.pqueues = {} # slot -> priority queue - for slot, startprios in (slot_startprios or {}).items(): - self.pqueues[slot] = self.pqfactory(startprios) + for priority in startprios: + self.queues[priority] = self.qfactory(priority) - def pop_slot(self, slot): - """ Pop an object from a priority queue for this slot """ - queue = self.pqueues[slot] - request = queue.pop() - if len(queue) == 0: - del self.pqueues[slot] - return request + self.curprio = min(startprios) - def push_slot(self, slot, obj, priority): - """ Push an object to a priority queue for this slot """ - if slot not in self.pqueues: - self.pqueues[slot] = self.pqfactory() - queue = self.pqueues[slot] - queue.push(obj, priority) + def qfactory(self, key): + return create_instance(self.downstream_queue_cls, + None, + self.crawler, + self.key + '/' + str(key)) + + def priority(self, request): + return -request.priority + + def push(self, request): + priority = self.priority(request) + if priority not in self.queues: + self.queues[priority] = self.qfactory(priority) + q = self.queues[priority] + q.push(request) # this may fail (eg. serialization error) + if self.curprio is None or priority < self.curprio: + self.curprio = priority + + def pop(self): + if self.curprio is None: + return + q = self.queues[self.curprio] + m = q.pop() + if not q: + del self.queues[self.curprio] + q.close() + prios = [p for p, q in self.queues.items() if q] + self.curprio = min(prios) if prios else None + return m def close(self): - active = {slot: queue.close() - for slot, queue in self.pqueues.items()} - self.pqueues.clear() + active = [] + for p, q in self.queues.items(): + active.append(p) + q.close() return active def __len__(self): - return sum(len(x) for x in self.pqueues.values()) if self.pqueues else 0 - - -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 + return sum(len(x) for x in self.queues.values()) if self.queues else 0 class DownloaderInterface(object): @@ -133,16 +130,16 @@ class DownloaderInterface(object): class DownloaderAwarePriorityQueue(object): - """ PriorityQueue which takes Downlaoder activity in account: + """ PriorityQueue which takes Downloader activity in account: domains (slots) with the least amount of active downloads are dequeued first. """ @classmethod - def from_crawler(cls, crawler, qfactory, slot_startprios=None, serialize=False): - return cls(crawler, qfactory, slot_startprios, serialize) + def from_crawler(cls, crawler, downstream_queue_cls, key, startprios=()): + return cls(crawler, downstream_queue_cls, key, startprios) - def __init__(self, crawler, qfactory, slot_startprios=None, serialize=False): + def __init__(self, crawler, downstream_queue_cls, key, slot_startprios=()): if crawler.settings.getint('CONCURRENT_REQUESTS_PER_IP') != 0: raise ValueError('"%s" does not support CONCURRENT_REQUESTS_PER_IP' % (self.__class__,)) @@ -156,35 +153,49 @@ class DownloaderAwarePriorityQueue(object): "queue class can be resumed." % slot_startprios.__class__) - slot_startprios = { - slot: [_Priority(p, slot) for p in startprios] - for slot, startprios in (slot_startprios or {}).items()} - - def pqfactory(startprios=()): - return ScrapyPriorityQueue(crawler, qfactory, startprios, serialize) - self._slot_pqueues = _SlotPriorityQueues(pqfactory, slot_startprios) - self.serialize = serialize self._downloader_interface = DownloaderInterface(crawler) + self.downstream_queue_cls = downstream_queue_cls + self.key = key + self.crawler = crawler + + self.pqueues = {} # slot -> priority queue + for slot, startprios in (slot_startprios or {}).items(): + self.pqueues[slot] = self.pqfactory(slot, startprios) + + def pqfactory(self, slot, startprios=()): + return ScrapyPriorityQueue(self.crawler, + self.downstream_queue_cls, + self.key + '/' + _path_safe(slot), + startprios) def pop(self): - stats = self._downloader_interface.stats(self._slot_pqueues.pqueues) + stats = self._downloader_interface.stats(self.pqueues) if not stats: return slot = min(stats)[1] - request = self._slot_pqueues.pop_slot(slot) + queue = self.pqueues[slot] + request = queue.pop() + if len(queue) == 0: + del self.pqueues[slot] return request - def push(self, request, priority): + def push(self, request): slot = self._downloader_interface.get_slot_key(request) - priority_slot = _Priority(priority=priority, slot=slot) - self._slot_pqueues.push_slot(slot, request, priority_slot) + if slot not in self.pqueues: + self.pqueues[slot] = self.pqfactory(slot) + queue = self.pqueues[slot] + queue.push(request) def close(self): - active = self._slot_pqueues.close() - return {slot: [p.priority for p in startprios] - for slot, startprios in active.items()} + active = {slot: queue.close() + for slot, queue in self.pqueues.items()} + self.pqueues.clear() + return active def __len__(self): - return len(self._slot_pqueues) + return sum(len(x) for x in self.pqueues.values()) if self.pqueues else 0 + + def __contains__(self, slot): + return slot in self.pqueues diff --git a/scrapy/squeues.py b/scrapy/squeues.py index d5d3be67e..d0686dac3 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -3,10 +3,27 @@ Scheduler queues """ import marshal +import os import pickle from queuelib import queue +from scrapy.utils.reqser import request_to_dict, request_from_dict + + +def _with_mkdir(queue_class): + + class DirectoriesCreated(queue_class): + + def __init__(self, path, *args, **kwargs): + dirname = os.path.dirname(path) + if not os.path.exists(dirname): + os.makedirs(dirname, exist_ok=True) + + super(DirectoriesCreated, self).__init__(path, *args, **kwargs) + + return DirectoriesCreated + def _serializable_queue(queue_class, serialize, deserialize): @@ -24,6 +41,44 @@ def _serializable_queue(queue_class, serialize, deserialize): return SerializableQueue +def _scrapy_serialization_queue(queue_class): + + class ScrapyRequestQueue(queue_class): + + def __init__(self, crawler, key): + self.spider = crawler.spider + super(ScrapyRequestQueue, self).__init__(key) + + @classmethod + def from_crawler(cls, crawler, key, *args, **kwargs): + return cls(crawler, key) + + def push(self, request): + request = request_to_dict(request, self.spider) + return super(ScrapyRequestQueue, self).push(request) + + def pop(self): + request = super(ScrapyRequestQueue, self).pop() + + if not request: + return None + + request = request_from_dict(request, self.spider) + return request + + return ScrapyRequestQueue + + +def _scrapy_non_serialization_queue(queue_class): + + class ScrapyRequestQueue(queue_class): + @classmethod + def from_crawler(cls, crawler, *args, **kwargs): + return cls() + + return ScrapyRequestQueue + + def _pickle_serialize(obj): try: return pickle.dumps(obj, protocol=2) @@ -34,13 +89,38 @@ def _pickle_serialize(obj): raise ValueError(str(e)) -PickleFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, - _pickle_serialize, pickle.loads) -PickleLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, - _pickle_serialize, pickle.loads) -MarshalFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, - marshal.dumps, marshal.loads) -MarshalLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, - marshal.dumps, marshal.loads) -FifoMemoryQueue = queue.FifoMemoryQueue -LifoMemoryQueue = queue.LifoMemoryQueue +PickleFifoDiskQueueNonRequest = _serializable_queue( + _with_mkdir(queue.FifoDiskQueue), + _pickle_serialize, + pickle.loads +) +PickleLifoDiskQueueNonRequest = _serializable_queue( + _with_mkdir(queue.LifoDiskQueue), + _pickle_serialize, + pickle.loads +) +MarshalFifoDiskQueueNonRequest = _serializable_queue( + _with_mkdir(queue.FifoDiskQueue), + marshal.dumps, + marshal.loads +) +MarshalLifoDiskQueueNonRequest = _serializable_queue( + _with_mkdir(queue.LifoDiskQueue), + marshal.dumps, + marshal.loads +) + +PickleFifoDiskQueue = _scrapy_serialization_queue( + PickleFifoDiskQueueNonRequest +) +PickleLifoDiskQueue = _scrapy_serialization_queue( + PickleLifoDiskQueueNonRequest +) +MarshalFifoDiskQueue = _scrapy_serialization_queue( + MarshalFifoDiskQueueNonRequest +) +MarshalLifoDiskQueue = _scrapy_serialization_queue( + MarshalLifoDiskQueueNonRequest +) +FifoMemoryQueue = _scrapy_non_serialization_queue(queue.FifoMemoryQueue) +LifoMemoryQueue = _scrapy_non_serialization_queue(queue.LifoMemoryQueue) diff --git a/tests/test_squeues.py b/tests/test_squeues.py index d5fcf2f7f..5c626fbcb 100644 --- a/tests/test_squeues.py +++ b/tests/test_squeues.py @@ -1,7 +1,12 @@ import pickle from queuelib.tests import test_queue as t -from scrapy.squeues import MarshalFifoDiskQueue, MarshalLifoDiskQueue, PickleFifoDiskQueue, PickleLifoDiskQueue +from scrapy.squeues import ( + MarshalFifoDiskQueueNonRequest as MarshalFifoDiskQueue, + MarshalLifoDiskQueueNonRequest as MarshalLifoDiskQueue, + PickleFifoDiskQueueNonRequest as PickleFifoDiskQueue, + PickleLifoDiskQueueNonRequest as PickleLifoDiskQueue +) from scrapy.item import Item, Field from scrapy.http import Request from scrapy.loader import ItemLoader