diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 7b9ff7e39..6e13e64d6 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1142,13 +1142,13 @@ Type of in-memory queue used by scheduler. Other available type is: SCHEDULER_PRIORITY_QUEUE ------------------------ -Default: ``'queuelib.PriorityQueue'`` +Default: ``'scrapy.pqueues.ScrapyPriorityQueue'`` Type of priority queue used by scheduler. Another available type is ``scrapy.pqueues.DownloaderAwarePriorityQueue``. ``scrapy.pqueues.DownloaderAwarePriorityQueue`` is works better than -``'queuelib.PriorityQueue'`` when you crawl many different domains in parallel. -But ``scrapy.pqueues.DownloaderAwarePriorityQueue`` +``scrapy.pqueues.ScrapyPriorityQueue`` when you crawl many different +domains in parallel. But ``scrapy.pqueues.DownloaderAwarePriorityQueue`` does not work together with :setting:`CONCURRENT_REQUESTS_PER_IP`. .. setting:: SPIDER_CONTRACTS diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index d40f3aa0c..c385fafe1 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -1,17 +1,44 @@ import os import json import logging +import warnings from os.path import join, exists -from scrapy.utils.reqser import request_to_dict, request_from_dict +from queuelib import PriorityQueue + from scrapy.utils.misc import load_object, create_instance from scrapy.utils.job import job_dir +from scrapy.utils.deprecate import ScrapyDeprecationWarning + logger = logging.getLogger(__name__) class Scheduler(object): + """ + Scrapy Scheduler. It allows to enqueue requests and then get + a next request to download. Scheduler is also handling duplication + filtering, via dupefilter. + Prioritization and queueing is not performed by the Scheduler. + User sets ``priority`` field for each Request, and a PriorityQueue + (defined by :setting:`SCHEDULER_PRIORITY_QUEUE`) uses these priorities + to dequeue requests in a desired order. + + Scheduler uses two PriorityQueue instances, configured to work in-memory + and on-disk (optional). When on-disk queue is present, it is used by + default, and an in-memory queue is used as a fallback for cases where + a disk queue can't handle a request (can't serialize it). + + :setting:`SCHEDULER_MEMORY_QUEUE` and + :setting:`SCHEDULER_DISK_QUEUE` allow to specify lower-level queue classes + which PriorityQueue instances would be instantiated with, to keep requests + on disk and in memory respectively. + + Overall, Scheduler is an object which holds several PriorityQueue instances + (in-memory and on-disk) and implements fallback logic for them. + Also, it handles dupefilters. + """ def __init__(self, dupefilter, jobdir=None, dqclass=None, mqclass=None, logunser=False, stats=None, pqclass=None, crawler=None): self.df = dupefilter @@ -29,9 +56,19 @@ class Scheduler(object): dupefilter_cls = load_object(settings['DUPEFILTER_CLASS']) dupefilter = create_instance(dupefilter_cls, settings, crawler) pqclass = load_object(settings['SCHEDULER_PRIORITY_QUEUE']) + if pqclass is PriorityQueue: + # backwards compatibility + warnings.warn("SCHEDULER_PRIORITY_QUEUE='queuelib.PriorityQueue'" + " is no longer supported because of API changes; " + "please use 'scrapy.pqueues.ScrapyPriorityQueue'", + ScrapyDeprecationWarning) + from scrapy.pqueues import ScrapyPriorityQueue + pqclass = ScrapyPriorityQueue + dqclass = load_object(settings['SCHEDULER_DISK_QUEUE']) mqclass = load_object(settings['SCHEDULER_MEMORY_QUEUE']) - logunser = settings.getbool('LOG_UNSERIALIZABLE_REQUESTS', settings.getbool('SCHEDULER_DEBUG')) + logunser = settings.getbool('LOG_UNSERIALIZABLE_REQUESTS', + settings.getbool('SCHEDULER_DEBUG')) return cls(dupefilter, jobdir=job_dir(settings), logunser=logunser, stats=crawler.stats, pqclass=pqclass, dqclass=dqclass, mqclass=mqclass, crawler=crawler) @@ -41,15 +78,19 @@ class Scheduler(object): def open(self, spider): self.spider = spider - self.mqs = create_instance(self.pqclass, None, self.crawler, self._newmq) + + # in-memory PriorityQueue instance + self.mqs = self._mq() + + # on-disk PriorityQueue instance self.dqs = self._dq() if self.dqdir else None + return self.df.open() def close(self, reason): if self.dqs: - prios = self.dqs.close() - with open(join(self.dqdir, 'active.json'), 'w') as f: - json.dump(prios, f) + state = self.dqs.close() + self._write_dqs_state(self.dqdir, state) return self.df.close(reason) def enqueue_request(self, request): @@ -66,7 +107,7 @@ class Scheduler(object): return True def next_request(self): - request = self.mqs.pop() + request = self._mqpop() if request: self.stats.inc_value('scheduler/dequeued/memory', spider=self.spider) else: @@ -84,8 +125,7 @@ class Scheduler(object): if self.dqs is None: return try: - reqd = request_to_dict(request, self.spider) - self.dqs.push(reqd, -request.priority) + self.dqs.push(request, -request.priority) except ValueError as e: # non serializable request if self.logunser: msg = ("Unable to serialize request: %(request)s - reason:" @@ -105,37 +145,54 @@ class Scheduler(object): def _dqpop(self): if self.dqs: - d = self.dqs.pop() - if d: - return request_from_dict(d, self.spider) + return self.dqs.pop() + + def _mqpop(self): + return self.mqs.pop() def _newmq(self, priority): + """ Factory for creating memory queues. """ return self.mqclass() def _newdq(self, priority): - return self.dqclass(join(self.dqdir, 'p%s' % (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) def _dq(self): - activef = join(self.dqdir, 'active.json') - if exists(activef): - with open(activef) as f: - prios = json.load(f) - else: - prios = () - + """ 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, - startprios=prios) + state, + serialize=True) if q: logger.info("Resuming crawl (%(queuesize)d requests scheduled)", {'queuesize': len(q)}, extra={'spider': self.spider}) return q def _dqdir(self, jobdir): + """ Return a folder name to keep disk queue state at """ if jobdir: dqdir = join(jobdir, 'requests.queue') if not exists(dqdir): os.makedirs(dqdir) return dqdir + + def _read_dqs_state(self, dqdir): + path = join(dqdir, 'active.json') + if not exists(path): + return () + with open(path) as f: + return json.load(f) + + def _write_dqs_state(self, dqdir, state): + with open(join(dqdir, 'active.json'), 'w') as f: + json.dump(state, f) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 6a9feb599..622f6bbc5 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -1,10 +1,10 @@ import hashlib import logging from collections import namedtuple -from six.moves.urllib.parse import urlparse from queuelib import PriorityQueue +from scrapy.utils.reqser import request_to_dict, request_from_dict from scrapy.core.downloader import Downloader from scrapy.http import Request from scrapy.signals import request_reached_downloader, request_left_downloader @@ -17,16 +17,6 @@ logger = logging.getLogger(__name__) SCHEDULER_SLOT_META_KEY = Downloader.DOWNLOAD_SLOT -def _get_request_meta(request): - if isinstance(request, dict): - return request.setdefault('meta', {}) - - if isinstance(request, Request): - return request.meta - - raise ValueError('Bad type of request "%s"' % (request.__class__, )) - - def _scheduler_slot_read(request, default=None): return request.meta.get(SCHEDULER_SLOT_META_KEY, default) @@ -37,24 +27,17 @@ def _scheduler_slot_write(request, slot): def _set_scheduler_slot(request): """ - >>> _set_scheduler_slot({'url':'http://foo.com'}) == _set_scheduler_slot({'url':'http://bar.com'}) - False - >>> _set_scheduler_slot({'url':'http://foo.com'}) == _set_scheduler_slot({'url':'http://foo.com'}) - True + >>> request = Request('http://example.com') + >>> _set_scheduler_slot(request) + 'example.com' + >>> _scheduler_slot_read(request) + 'example.com' """ - meta = _get_request_meta(request) - slot = meta.get(SCHEDULER_SLOT_META_KEY, None) - + slot = _scheduler_slot_read(request, None) if slot is not None: return slot - - if isinstance(request, dict): - url = request.get('url', None) - slot = urlparse(url).hostname or '' - elif isinstance(request, Request): - slot = urlparse_cached(request).hostname or '' - - meta[SCHEDULER_SLOT_META_KEY] = slot + slot = urlparse_cached(request).hostname or '' + _scheduler_slot_write(request, slot) return slot @@ -68,30 +51,25 @@ def _path_safe(text): return '-'.join([pathable_slot, unique_slot]) -class PrioritySlot(namedtuple("PrioritySlot", ["priority", "slot"])): - """ ``(priority, slot)`` tuple which uses a path-safe slot name - when converting to str """ +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: + + * 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))) -class PriorityAsTupleQueue(PriorityQueue): - """ - Python structures is not directly (de)serialized (to)from json. - We need this modified queue to transform custom structure (from)to - json serializable structures - """ - def __init__(self, qfactory, startprios=()): - startprios = [PrioritySlot(priority=p[0], slot=p[1]) - for p in startprios] - super(PriorityAsTupleQueue, self).__init__( - qfactory=qfactory, - startprios=startprios) - - -class SlotPriorityQueues(object): +class _SlotPriorityQueues(object): """ Container for multiple priority queues. """ def __init__(self, pqfactory, slot_startprios=None): """ @@ -134,44 +112,78 @@ class SlotPriorityQueues(object): return slot in self.pqueues -class DownloaderAwarePriorityQueue(object): - - _DOWNLOADER_AWARE_PQ_ID = 'DOWNLOADER_AWARE_PQ_ID' +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=None): - return cls(crawler, qfactory, startprios) + def from_crawler(cls, crawler, qfactory, startprios=(), serialize=False): + return cls(crawler, qfactory, startprios, serialize) - def __init__(self, crawler, qfactory, startprios=None): - ip_concurrency_key = 'CONCURRENT_REQUESTS_PER_IP' - ip_concurrency = crawler.settings.getint(ip_concurrency_key, 0) + def push(self, request, priority=0): + if self.serialize: + request = request_to_dict(request, self.spider) + super(ScrapyPriorityQueue, self).push(request, priority) - if ip_concurrency > 0: - raise ValueError('"%s" does not support setting %s' % (self.__class__, - ip_concurrency_key)) + def pop(self): + request = super(ScrapyPriorityQueue, self).pop() + if request and self.serialize: + request = request_from_dict(request, self.spider) + return request + + +class DownloaderAwarePriorityQueue(object): + """ PriorityQueue which takes Downlaoder activity in account: + domains (slots) with the least amount of active downloads are dequeued + first. + """ + _DOWNLOADER_AWARE_PQ_ID = '_DOWNLOADER_AWARE_PQ_ID' + + @classmethod + def from_crawler(cls, crawler, qfactory, slot_startprios=None, serialize=False): + return cls(crawler, qfactory, slot_startprios, serialize) + + def __init__(self, crawler, qfactory, slot_startprios=None, serialize=False): + if crawler.settings.getint('CONCURRENT_REQUESTS_PER_IP') != 0: + raise ValueError('"%s" does not support CONCURRENT_REQUESTS_PER_IP' + % (self.__class__,)) + + if slot_startprios and not isinstance(slot_startprios, dict): + raise ValueError("DownloaderAwarePriorityQueue accepts " + "``slot_startprios`` as a dict; %r instance " + "is passed. Most likely, it means the state is" + "created by an incompatible priority queue. " + "Only a crawl started with the same priority " + "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 PriorityAsTupleQueue(qfactory, startprios) - - if startprios and not isinstance(startprios, dict): - raise ValueError("DownloaderAwarePriorityQueue accepts " - "``startprios`` as a dict; %r instance is passed." - " Only a crawl started with the same priority " - "queue class can be resumed." % startprios.__class__) - self._slot_pqueues = SlotPriorityQueues(pqfactory, - slot_startprios=startprios) + return ScrapyPriorityQueue(crawler, qfactory, startprios, serialize) + self._slot_pqueues = _SlotPriorityQueues(pqfactory, slot_startprios) self._active_downloads = {slot: 0 for slot in self._slot_pqueues.pqueues} crawler.signals.connect(self.on_response_download, signal=request_left_downloader) crawler.signals.connect(self.on_request_reached_downloader, signal=request_reached_downloader) + self.serialize = serialize + # There are two PriorityQueues at the same time (memory and disk-based), + # and they both listen to Downloader signals. To filter out signals + # coming from the other queue, each queue keeps track of its own + # requests using mark / unmark / check_mark methods. def mark(self, request): - meta = _get_request_meta(request) - if not isinstance(meta, dict): - raise ValueError('No meta attribute in %s' % (request, )) - meta[self._DOWNLOADER_AWARE_PQ_ID] = id(self) + request.meta[self._DOWNLOADER_AWARE_PQ_ID] = id(self) def check_mark(self, request): return request.meta.get(self._DOWNLOADER_AWARE_PQ_ID, None) == id(self) @@ -194,7 +206,7 @@ class DownloaderAwarePriorityQueue(object): def push(self, request, priority): slot = _set_scheduler_slot(request) - priority_slot = PrioritySlot(priority=priority, slot=slot) + priority_slot = _Priority(priority=priority, slot=slot) self._slot_pqueues.push_slot(slot, request, priority_slot) if slot not in self._active_downloads: self._active_downloads[slot] = 0 @@ -206,8 +218,8 @@ class DownloaderAwarePriorityQueue(object): slot = _scheduler_slot_read(request) if slot not in self._active_downloads or self._active_downloads[slot] <= 0: - raise ValueError('Get response for wrong slot "%s"' % (slot, )) - self._active_downloads[slot] = self._active_downloads[slot] - 1 + raise ValueError('Got response for a wrong slot "%s"' % (slot, )) + self._active_downloads[slot] -= 1 if self._active_downloads[slot] == 0 and slot not in self._slot_pqueues: del self._active_downloads[slot] @@ -220,7 +232,9 @@ class DownloaderAwarePriorityQueue(object): def close(self): self._active_downloads.clear() - return self._slot_pqueues.close() + active = self._slot_pqueues.close() + return {slot: [p.priority for p in startprios] + for slot, startprios in active.items()} def __len__(self): return len(self._slot_pqueues) diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index ca004aedd..365b405cb 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -244,7 +244,7 @@ ROBOTSTXT_OBEY = False SCHEDULER = 'scrapy.core.scheduler.Scheduler' SCHEDULER_DISK_QUEUE = 'scrapy.squeues.PickleLifoDiskQueue' SCHEDULER_MEMORY_QUEUE = 'scrapy.squeues.LifoMemoryQueue' -SCHEDULER_PRIORITY_QUEUE = 'queuelib.PriorityQueue' +SCHEDULER_PRIORITY_QUEUE = 'scrapy.pqueues.ScrapyPriorityQueue' SPIDER_LOADER_CLASS = 'scrapy.spiderloader.SpiderLoader' SPIDER_LOADER_WARN_ONLY = False diff --git a/scrapy/squeues.py b/scrapy/squeues.py index d2074a457..30cc926e5 100644 --- a/scrapy/squeues.py +++ b/scrapy/squeues.py @@ -7,6 +7,7 @@ from six.moves import cPickle as pickle from queuelib import queue + def _serializable_queue(queue_class, serialize, deserialize): class SerializableQueue(queue_class): @@ -22,6 +23,7 @@ def _serializable_queue(queue_class, serialize, deserialize): return SerializableQueue + def _pickle_serialize(obj): try: return pickle.dumps(obj, protocol=2) @@ -31,13 +33,14 @@ def _pickle_serialize(obj): except (pickle.PicklingError, AttributeError, TypeError) as e: raise ValueError(str(e)) -PickleFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, \ + +PickleFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, _pickle_serialize, pickle.loads) -PickleLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, \ +PickleLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, _pickle_serialize, pickle.loads) -MarshalFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, \ +MarshalFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, marshal.dumps, marshal.loads) -MarshalLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, \ +MarshalLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, marshal.dumps, marshal.loads) FifoMemoryQueue = queue.FifoMemoryQueue LifoMemoryQueue = queue.LifoMemoryQueue