From a25cf5c82f99f7ae11346a2e565d6255835c3814 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Tue, 20 Nov 2018 16:13:09 +0000 Subject: [PATCH 01/38] function to get unique file queues for any type of base queue --- scrapy/core/queues.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) create mode 100644 scrapy/core/queues.py diff --git a/scrapy/core/queues.py b/scrapy/core/queues.py new file mode 100644 index 000000000..96d582fc7 --- /dev/null +++ b/scrapy/core/queues.py @@ -0,0 +1,15 @@ +import uuid +import os.path + + +def unique_files_queue(queue_class): + + class UniqueFilesQueue(queue_class): + def __init__(self, path): + path = path + "-" + uuid.uuid4().hex + while os.path.exists(path): + path = path + "-" + uuid.uuid4().hex + + super().__init__(path) + + return UniqueFilesQueue From 821f5bb26077d7f9a6b2b1a72f210f81779f5393 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Mon, 3 Dec 2018 11:00:03 +0000 Subject: [PATCH 02/38] First implementation handle exception use O(N) instead of O(NlogN) here we have request as struct additional check for meptiness small performance improvement do not consume another request test number of responses mark requests back to 3 slots test case raise exceptions in case of missed meta add marks to requests and work only with your own requests only disk queue should obtain signals separate functions for slot rasd/write use signlas without variable stop crawler get signals in correct place logic test for download-aware priority queue update comment for structure ensure text type transform slot name to path use implicit structure use unicode type implicitly use real crawler add signals more slot accounting simple implementation of pop small slot accounting code no need for custom len function ability to call super in py27 add slots generic tests for downloader aware queue dummy implementation of crawler aware priority queue move common logic to base class rename class pass crawler to pqclass constructor do not copy quelib.PriorityQueue code add comment about new class remove obsolete function modify behaviour of queuelib.PriorityQueue to dodge very complex priority better way to get name remove obsolete commentary check boundaries function for priority convertion with known limits correct import path move file do not switch on by deffault as ip concurrency not supported set scheduler slot in case of empty slot use constant single place for added urls single place for constants use as default queue correct format for error text test migration from old version with on disk queue in these tests we have only two inflection points - jobdir and priority_queue_cls we do not need separate mock spider, use usual one do not rely on order of dict elements, imply order of list test round robiness of priority queue add comments and requirements for our magick function remove debug logging put queues into slot as we fabricate priorities we do not need special types anymore fabricate priority for priority queue more versatile priorities Scheduler class is not inflection point wrap correct types check for emptinees before initialization tests for new priority queue correct default type for startprios use exact values put common settings to base class test priorities for disk scheduler test dequeue for disk scheduler test length for disk scheduler setUp/tearDown methods for on disk schedulers new methods remove excessive line base class to handle scheduler creation correct method names test priorities deque test close scheduler on test end enqueue some requests test template for scheduler use downloader slot I/O implementation for RoundRobin queue round-robin implementation without I/O and slot detection wrappers for every disk queue class --- scrapy/core/downloader/__init__.py | 8 +- scrapy/core/queues.py | 15 -- scrapy/core/scheduler.py | 17 +- scrapy/pqueues.py | 246 ++++++++++++++++++++++ tests/test_scheduler.py | 315 +++++++++++++++++++++++++++++ 5 files changed, 578 insertions(+), 23 deletions(-) delete mode 100644 scrapy/core/queues.py create mode 100644 scrapy/pqueues.py create mode 100644 tests/test_scheduler.py diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index 59c3ad074..4695d75f4 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -75,6 +75,8 @@ def _get_concurrency_delay(concurrency, spider, settings): class Downloader(object): + DOWNLOAD_SLOT = 'download_slot' + def __init__(self, crawler): self.settings = crawler.settings self.signals = crawler.signals @@ -111,8 +113,8 @@ class Downloader(object): return key, self.slots[key] def _get_slot_key(self, request, spider): - if 'download_slot' in request.meta: - return request.meta['download_slot'] + if self.DOWNLOAD_SLOT in request.meta: + return request.meta[self.DOWNLOAD_SLOT] key = urlparse_cached(request).hostname or '' if self.ip_concurrency: @@ -122,7 +124,7 @@ class Downloader(object): def _enqueue_request(self, request, spider): key, slot = self._get_slot(request, spider) - request.meta['download_slot'] = key + request.meta[self.DOWNLOAD_SLOT] = key def _deactivate(response): slot.active.remove(request) diff --git a/scrapy/core/queues.py b/scrapy/core/queues.py deleted file mode 100644 index 96d582fc7..000000000 --- a/scrapy/core/queues.py +++ /dev/null @@ -1,15 +0,0 @@ -import uuid -import os.path - - -def unique_files_queue(queue_class): - - class UniqueFilesQueue(queue_class): - def __init__(self, path): - path = path + "-" + uuid.uuid4().hex - while os.path.exists(path): - path = path + "-" + uuid.uuid4().hex - - super().__init__(path) - - return UniqueFilesQueue diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index eb790a67e..d40f3aa0c 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -13,7 +13,7 @@ logger = logging.getLogger(__name__) class Scheduler(object): def __init__(self, dupefilter, jobdir=None, dqclass=None, mqclass=None, - logunser=False, stats=None, pqclass=None): + logunser=False, stats=None, pqclass=None, crawler=None): self.df = dupefilter self.dqdir = self._dqdir(jobdir) self.pqclass = pqclass @@ -21,6 +21,7 @@ class Scheduler(object): self.mqclass = mqclass self.logunser = logunser self.stats = stats + self.crawler = crawler @classmethod def from_crawler(cls, crawler): @@ -32,14 +33,15 @@ class Scheduler(object): mqclass = load_object(settings['SCHEDULER_MEMORY_QUEUE']) 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) + stats=crawler.stats, pqclass=pqclass, dqclass=dqclass, + mqclass=mqclass, crawler=crawler) def has_pending_requests(self): return len(self) > 0 def open(self, spider): self.spider = spider - self.mqs = self.pqclass(self._newmq) + self.mqs = create_instance(self.pqclass, None, self.crawler, self._newmq) self.dqs = self._dq() if self.dqdir else None return self.df.open() @@ -111,7 +113,7 @@ class Scheduler(object): return self.mqclass() def _newdq(self, priority): - return self.dqclass(join(self.dqdir, 'p%s' % priority)) + return self.dqclass(join(self.dqdir, 'p%s' % (priority, ))) def _dq(self): activef = join(self.dqdir, 'active.json') @@ -120,7 +122,12 @@ class Scheduler(object): prios = json.load(f) else: prios = () - q = self.pqclass(self._newdq, startprios=prios) + + q = create_instance(self.pqclass, + None, + self.crawler, + self._newdq, + startprios=prios) 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 new file mode 100644 index 000000000..75073b7a4 --- /dev/null +++ b/scrapy/pqueues.py @@ -0,0 +1,246 @@ +from collections import deque +import hashlib +import logging +from six import text_type +from six.moves.urllib.parse import urlparse + +from queuelib import PriorityQueue + +from scrapy.core.downloader import Downloader +from scrapy.http import Request +from scrapy.signals import request_reached_downloader, response_downloaded + + +logger = logging.getLogger(__name__) + + +SCHEDULER_SLOT_META_KEY = Downloader.DOWNLOAD_SLOT + + +def _get_from_request(request, key, default=None): + if isinstance(request, dict): + return request.get(key, default) + + if isinstance(request, Request): + return getattr(request, key, default) + + raise ValueError('Bad type of request "%s"' % (request.__class__, )) + + +def _scheduler_slot_read(request, default=None): + meta = _get_from_request(request, 'meta', dict()) + slot = meta.get(SCHEDULER_SLOT_META_KEY, default) + return slot + + +def _scheduler_slot_write(request, slot): + meta = _get_from_request(request, 'meta', None) + if not isinstance(meta, dict): + raise ValueError('No meta attribute in %s' % (request, )) + meta[SCHEDULER_SLOT_META_KEY] = slot + + +def _scheduler_slot(request): + + slot = _scheduler_slot_read(request, None) + if slot is None: + url = _get_from_request(request, 'url') + slot = urlparse(url).hostname or '' + _scheduler_slot_write(request, slot) + + return slot + + +def _pathable(x): + pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' for c in x]) + + """ + as we replace some letters we can get collision for different slots + add we add unique part + """ + unique_slot = hashlib.md5(x.encode('utf8')).hexdigest() + + return '-'.join([pathable_slot, unique_slot]) + + +class PrioritySlot: + __slots__ = ('priority', 'slot') + + def __init__(self, priority=0, slot=None): + self.priority = priority + self.slot = slot + + def __hash__(self): + return hash((self.priority, self.slot)) + + def __eq__(self, other): + return (self.priority, self.slot) == (other.priority, other.slot) + + def __lt__(self, other): + return (self.priority, self.slot) < (other.priority, other.slot) + + def __str__(self): + return '_'.join([text_type(self.priority), _pathable(text_type(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=()): + + super(PriorityAsTupleQueue, self).__init__( + qfactory, + [PrioritySlot(priority=p[0], slot=p[1]) for p in startprios] + ) + + def close(self): + startprios = super(PriorityAsTupleQueue, self).close() + return [(s.priority, s.slot) for s in startprios] + + def is_empty(self): + return not self.queues or len(self) == 0 + + +class SlotBasedPriorityQueue(object): + + def __init__(self, qfactory, startprios={}): + self.pqueues = dict() # slot -> priority queue + self.qfactory = qfactory # factory for creating new internal queues + + if not startprios: + return + + if not isinstance(startprios, dict): + raise ValueError("Looks like your priorities file malforfemed. " + "Possible reason: You run scrapy with previous " + "version. Interrupted it. Updated scrapy. And " + "run again.") + + for slot, prios in startprios.items(): + self.pqueues[slot] = PriorityAsTupleQueue(self.qfactory, prios) + + def pop_slot(self, slot): + queue = self.pqueues[slot] + request = queue.pop() + is_empty = queue.is_empty() + if is_empty: + del self.pqueues[slot] + + return request, is_empty + + def push_slot(self, request, priority): + slot = _scheduler_slot(request) + is_new = False + if slot not in self.pqueues: + is_new = True + self.pqueues[slot] = PriorityAsTupleQueue(self.qfactory) + self.pqueues[slot].push(request, PrioritySlot(priority=priority, slot=slot)) + return slot, is_new + + def close(self): + startprios = dict() + for slot, queue in self.pqueues.items(): + prios = queue.close() + startprios[slot] = prios + self.pqueues.clear() + return startprios + + def __len__(self): + return sum(len(x) for x in self.pqueues.values()) if self.pqueues else 0 + + +class RoundRobinPriorityQueue(SlotBasedPriorityQueue): + + def __init__(self, qfactory, startprios={}): + super(RoundRobinPriorityQueue, self).__init__(qfactory, startprios) + self._slots = deque() + for slot in self.pqueues: + self._slots.append(slot) + + def push(self, request, priority): + slot, is_new = self.push_slot(request, priority) + if is_new: + self._slots.append(slot) + + def pop(self): + if not self._slots: + return + + slot = self._slots.popleft() + request, is_empty = self.pop_slot(slot) + + if not is_empty: + self._slots.append(slot) + + return request + + def close(self): + self._slots.clear() + return super(RoundRobinPriorityQueue, self).close() + + +class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): + + _DOWNLOADER_AWARE_PQ_ID = 'DOWNLOADER_AWARE_PQ_ID' + + @classmethod + def from_crawler(cls, crawler, qfactory, startprios={}): + return cls(crawler, qfactory, startprios) + + def __init__(self, crawler, qfactory, startprios={}): + super(DownloaderAwarePriorityQueue, self).__init__(qfactory, startprios) + self._slots = {slot: 0 for slot in self.pqueues} + crawler.signals.connect(self.on_response_download, + signal=response_downloaded) + crawler.signals.connect(self.on_request_reached_downloader, + signal=request_reached_downloader) + + def mark(self, request): + meta = _get_from_request(request, 'meta', None) + if not isinstance(meta, dict): + raise ValueError('No meta attribute in %s' % (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) + + def pop(self): + slots = [(d, s) for s,d in self._slots.items() if s in self.pqueues] + + if not slots: + return + + slot = min(slots)[1] + request, _ = self.pop_slot(slot) + self.mark(request) + return request + + def push(self, request, priority): + slot, _ = self.push_slot(request, priority) + if slot not in self._slots: + self._slots[slot] = 0 + + def on_response_download(self, response, request, spider): + if not self.check_mark(request): + return + + slot = _scheduler_slot_read(request) + if slot not in self._slots or self._slots[slot] <= 0: + raise ValueError('Get response for wrong slot "%s"' % (slot, )) + self._slots[slot] = self._slots[slot] - 1 + if self._slots[slot] == 0 and slot not in self.pqueues: + del self._slots[slot] + + def on_request_reached_downloader(self, request, spider): + if not self.check_mark(request): + return + + slot = _scheduler_slot_read(request) + self._slots[slot] = self._slots.get(slot, 0) + 1 + + def close(self): + self._slots.clear() + return super(DownloaderAwarePriorityQueue, self).close() diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py new file mode 100644 index 000000000..fd86e8d8c --- /dev/null +++ b/tests/test_scheduler.py @@ -0,0 +1,315 @@ +import contextlib +import shutil +import tempfile +import unittest + +from scrapy.crawler import Crawler +from scrapy.core.scheduler import Scheduler +from scrapy.http import Request +from scrapy.pqueues import _scheduler_slot_read, _scheduler_slot_write +from scrapy.signals import request_reached_downloader, response_downloaded +from scrapy.spiders import Spider + +class MockCrawler(Crawler): + def __init__(self, priority_queue_cls, jobdir): + + settings = dict(LOG_UNSERIALIZABLE_REQUESTS=False, + SCHEDULER_DISK_QUEUE='scrapy.squeues.PickleLifoDiskQueue', + SCHEDULER_MEMORY_QUEUE='scrapy.squeues.LifoMemoryQueue', + SCHEDULER_PRIORITY_QUEUE=priority_queue_cls, + JOBDIR=jobdir, + DUPEFILTER_CLASS='scrapy.dupefilters.BaseDupeFilter') + super(MockCrawler, self).__init__(Spider, settings) + + +class SchedulerHandler: + priority_queue_cls = None + jobdir = None + + def create_scheduler(self): + self.mock_crawler = MockCrawler(self.priority_queue_cls, self.jobdir) + self.scheduler = Scheduler.from_crawler(self.mock_crawler) + self.spider = Spider(name='spider') + self.scheduler.open(self.spider) + + def close_scheduler(self): + self.scheduler.close('finished') + self.mock_crawler.stop() + + def setUp(self): + self.create_scheduler() + + def tearDown(self): + self.close_scheduler() + + +_PRIORITIES = [("http://foo.com/a", -2), + ("http://foo.com/d", 1), + ("http://foo.com/b", -1), + ("http://foo.com/c", 0), + ("http://foo.com/e", 2)] + + +_URLS = {"http://foo.com/a", "http://foo.com/b", "http://foo.com/c"} + + +class BaseSchedulerInMemoryTester(SchedulerHandler): + def test_length(self): + self.assertFalse(self.scheduler.has_pending_requests()) + self.assertEqual(len(self.scheduler), 0) + + for url in _URLS: + self.scheduler.enqueue_request(Request(url)) + + self.assertTrue(self.scheduler.has_pending_requests()) + self.assertEqual(len(self.scheduler), len(_URLS)) + + def test_dequeue(self): + for url in _URLS: + self.scheduler.enqueue_request(Request(url)) + + urls = set() + while self.scheduler.has_pending_requests(): + urls.add(self.scheduler.next_request().url) + + self.assertEqual(urls, _URLS) + + def test_dequeue_priorities(self): + for url, priority in _PRIORITIES: + self.scheduler.enqueue_request(Request(url, priority=priority)) + + priorities = list() + while self.scheduler.has_pending_requests(): + priorities.append(self.scheduler.next_request().priority) + + self.assertEqual(priorities, sorted([x[1] for x in _PRIORITIES], key=lambda x: -x)) + + +class BaseSchedulerOnDiskTester(SchedulerHandler): + + def setUp(self): + self.jobdir = tempfile.mkdtemp() + self.create_scheduler() + + def tearDown(self): + self.close_scheduler() + + shutil.rmtree(self.jobdir) + self.jobdir = None + + def test_length(self): + self.assertFalse(self.scheduler.has_pending_requests()) + self.assertEqual(len(self.scheduler), 0) + + for url in _URLS: + self.scheduler.enqueue_request(Request(url)) + + self.close_scheduler() + self.create_scheduler() + + self.assertTrue(self.scheduler.has_pending_requests()) + self.assertEqual(len(self.scheduler), len(_URLS)) + + def test_dequeue(self): + for url in _URLS: + self.scheduler.enqueue_request(Request(url)) + + self.close_scheduler() + self.create_scheduler() + + urls = set() + while self.scheduler.has_pending_requests(): + urls.add(self.scheduler.next_request().url) + + self.assertEqual(urls, _URLS) + + def test_dequeue_priorities(self): + for url, priority in _PRIORITIES: + self.scheduler.enqueue_request(Request(url, priority=priority)) + + self.close_scheduler() + self.create_scheduler() + + priorities = list() + while self.scheduler.has_pending_requests(): + priorities.append(self.scheduler.next_request().priority) + + self.assertEqual(priorities, sorted([x[1] for x in _PRIORITIES], key=lambda x: -x)) + + +class TestSchedulerInMemory(BaseSchedulerInMemoryTester, unittest.TestCase): + priority_queue_cls = 'queuelib.PriorityQueue' + + +class TestSchedulerOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): + priority_queue_cls = 'queuelib.PriorityQueue' + + +_SLOTS = [("http://foo.com/a", 'a'), + ("http://foo.com/b", 'a'), + ("http://foo.com/c", 'b'), + ("http://foo.com/d", 'b'), + ("http://foo.com/e", 'c'), + ("http://foo.com/f", 'c')] + + +class TestSchedulerWithRoundRobinInMemory(BaseSchedulerInMemoryTester, unittest.TestCase): + priority_queue_cls = 'scrapy.pqueues.RoundRobinPriorityQueue' + + def test_round_robin(self): + for url, slot in _SLOTS: + request = Request(url) + _scheduler_slot_write(request, slot) + self.scheduler.enqueue_request(request) + + slots = list() + while self.scheduler.has_pending_requests(): + slots.append(_scheduler_slot_read(self.scheduler.next_request())) + + for i in range(0, len(_SLOTS), 2): + self.assertNotEqual(slots[i], slots[i+1]) + + def test_is_meta_set(self): + url = "http://foo.com/a" + request = Request(url) + if _scheduler_slot_read(request): + _scheduler_slot_write(request, None) + self.scheduler.enqueue_request(request) + self.assertIsNotNone(_scheduler_slot_read(request, None), None) + + +class TestSchedulerWithRoundRobinOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): + priority_queue_cls = 'scrapy.pqueues.RoundRobinPriorityQueue' + + def test_round_robin(self): + for url, slot in _SLOTS: + request = Request(url) + _scheduler_slot_write(request, slot) + self.scheduler.enqueue_request(request) + + self.close_scheduler() + self.create_scheduler() + + slots = list() + while self.scheduler.has_pending_requests(): + slots.append(_scheduler_slot_read(self.scheduler.next_request())) + + for i in range(0, len(_SLOTS), 2): + self.assertNotEqual(slots[i], slots[i+1]) + + def test_is_meta_set(self): + url = "http://foo.com/a" + request = Request(url) + if _scheduler_slot_read(request): + _scheduler_slot_write(request, None) + self.scheduler.enqueue_request(request) + + self.close_scheduler() + self.create_scheduler() + + self.assertIsNotNone(_scheduler_slot_read(request, None), None) + + +@contextlib.contextmanager +def mkdtemp(): + dir = tempfile.mkdtemp() + try: + yield dir + finally: + shutil.rmtree(dir) + + +def _migration(): + + with mkdtemp() as tmp_dir: + prev_scheduler_handler = SchedulerHandler() + prev_scheduler_handler.priority_queue_cls = 'queuelib.PriorityQueue' + prev_scheduler_handler.jobdir = tmp_dir + + prev_scheduler_handler.create_scheduler() + for url in _URLS: + prev_scheduler_handler.scheduler.enqueue_request(Request(url)) + prev_scheduler_handler.close_scheduler() + + next_scheduler_handler = SchedulerHandler() + next_scheduler_handler.priority_queue_cls = 'scrapy.pqueues.RoundRobinPriorityQueue' + next_scheduler_handler.jobdir = tmp_dir + + next_scheduler_handler.create_scheduler() + + +class TestMigration(unittest.TestCase): + def test_migration(self): + self.assertRaises(ValueError, _migration) + + +class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, unittest.TestCase): + priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' + + def test_logic(self): + for url, slot in _SLOTS: + request = Request(url) + _scheduler_slot_write(request, slot) + self.scheduler.enqueue_request(request) + + slots = list() + requests = list() + while self.scheduler.has_pending_requests(): + request = self.scheduler.next_request() + slots.append(_scheduler_slot_read(request)) + self.mock_crawler.signals.send_catch_log( + signal=request_reached_downloader, + request=request, + spider=self.spider + ) + requests.append(request) + self.assertEqual(len(slots), len(_SLOTS)) + + for request in requests: + self.mock_crawler.signals.send_catch_log(signal=response_downloaded, + request=request, + response=None, + spider=self.spider) + + unique_slots = len(set(s for _, s in _SLOTS)) + for i in range(0, len(_SLOTS), unique_slots): + part = slots[i:i + unique_slots] + self.assertEqual(len(part), len(set(part))) + + +class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): + priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' + def test_logic(self): + for url, slot in _SLOTS: + request = Request(url) + _scheduler_slot_write(request, slot) + self.scheduler.enqueue_request(request) + + self.close_scheduler() + self.create_scheduler() + + slots = list() + requests = list() + while self.scheduler.has_pending_requests(): + request = self.scheduler.next_request() + slots.append(_scheduler_slot_read(request)) + self.mock_crawler.signals.send_catch_log( + signal=request_reached_downloader, + request=request, + spider=self.spider + ) + requests.append(request) + + self.assertEqual(self.scheduler.mqs._slots, {}) + self.assertEqual(len(slots), len(_SLOTS)) + + for request in requests: + self.mock_crawler.signals.send_catch_log(signal=response_downloaded, + request=request, + response=None, + spider=self.spider) + + unique_slots = len(set(s for _, s in _SLOTS)) + for i in range(0, len(_SLOTS), unique_slots): + part = slots[i:i + unique_slots] + self.assertEqual(len(part), len(set(part))) From afdb69ea6daac8bd4f580d6c20bf9e93b741957b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Mon, 3 Dec 2018 16:36:05 +0100 Subject: [PATCH 03/38] Add a troubleshooting section to the installation instructions Its initial content covers the workaround for #2473. --- docs/intro/install.rst | 31 ++++++++++++++++++++++++++++++- 1 file changed, 30 insertions(+), 1 deletion(-) diff --git a/docs/intro/install.rst b/docs/intro/install.rst index 4a9aa3cfb..daec7fcb7 100644 --- a/docs/intro/install.rst +++ b/docs/intro/install.rst @@ -30,7 +30,8 @@ dependencies depending on your operating system, so be sure to check the We strongly recommend that you install Scrapy in :ref:`a dedicated virtualenv `, to avoid conflicting with your system packages. -For more detailed and platform specifics instructions, read on. +For more detailed and platform specifics instructions, as well as +troubleshooting information, read on. Things that are good to know @@ -247,6 +248,34 @@ that setuptools was unable to pick up one PyPy-specific dependency. To fix this issue, run ``pip install 'PyPyDispatcher>=2.1.0'``. +.. _intro-install-troubleshooting: + +Troubleshooting +=============== + +AttributeError: 'module' object has no attribute 'OP_NO_TLSv1_1' +---------------------------------------------------------------- + +After you install or upgrade Scrapy, Twisted or pyOpenSSL, you may get an +exception with the following traceback:: + + […] + File "[…]/site-packages/twisted/protocols/tls.py", line 63, in + from twisted.internet._sslverify import _setAcceptableProtocols + File "[…]/site-packages/twisted/internet/_sslverify.py", line 38, in + TLSVersion.TLSv1_1: SSL.OP_NO_TLSv1_1, + AttributeError: 'module' object has no attribute 'OP_NO_TLSv1_1' + +The reason you get this exception is that your system or virtual environment +has a version of pyOpenSSL that your version of Twisted does not support. + +To install a version of pyOpenSSL that your version of Twisted supports, +reinstall Twisted with the :code:`tls` extra option:: + + pip install twisted[tls] + +For details, see `Issue #2473 `_. + .. _Python: https://www.python.org/ .. _pip: https://pip.pypa.io/en/latest/installing/ .. _lxml: http://lxml.de/ From 9c314800e4b195df41e5c0aba0d9ffe4bcffec8e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Mon, 3 Dec 2018 17:14:10 +0100 Subject: [PATCH 04/38] Document the SCRAPY_PROJECT environment variable Fixes #1109 --- docs/topics/commands.rst | 29 ++++++++++++++++++++++++++++- 1 file changed, 28 insertions(+), 1 deletion(-) diff --git a/docs/topics/commands.rst b/docs/topics/commands.rst index ef9c45196..97f8311de 100644 --- a/docs/topics/commands.rst +++ b/docs/topics/commands.rst @@ -37,7 +37,7 @@ Scrapy also understands, and can be configured through, a number of environment variables. Currently these are: * ``SCRAPY_SETTINGS_MODULE`` (see :ref:`topics-settings-module-envvar`) -* ``SCRAPY_PROJECT`` +* ``SCRAPY_PROJECT`` (see :ref:`topics-project-envvar`) * ``SCRAPY_PYTHON_SHELL`` (see :ref:`topics-shell`) .. _topics-project-structure: @@ -71,6 +71,33 @@ the project settings. Here is an example:: [settings] default = myproject.settings +.. _topics-project-envvar: + +Sharing the root directory between projects +=========================================== + +A project root directory, the one that contains the ``scrapy.cfg``, may be +shared by multiple Scrapy projects, each with its own settings module. + +In that case, you must define one or more aliases for those settings modules +under ``[settings]`` in your ``scrapy.cfg`` file:: + + [settings] + default = myproject1.settings + project1 = myproject1.settings + project2 = myproject2.settings + +By default, the ``scrapy`` command-line tool will use the ``default`` settings. +Use the ``SCRAPY_PROJECT`` environment variable to specify a different project +for ``scrapy`` to use:: + + $ scrapy settings --get BOT_NAME + Project 1 Bot + $ export SCRAPY_PROJECT=project2 + $ scrapy settings --get BOT_NAME + Project 2 Bot + + Using the ``scrapy`` tool ========================= From f56079f6c71a77c1f70510cf291cd808617933cd Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Wed, 5 Dec 2018 10:02:42 +0000 Subject: [PATCH 05/38] Test cleanups PEP8 fixes no need to close implicitly do not use pytest need to put it into class remove round-robin queue additional check for empty queue use pytest tmpdir fixture --- scrapy/pqueues.py | 50 ++++------------ tests/test_scheduler.py | 128 ++++++++++++---------------------------- 2 files changed, 50 insertions(+), 128 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 75073b7a4..287a8de35 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -1,4 +1,3 @@ -from collections import deque import hashlib import logging from six import text_type @@ -71,16 +70,17 @@ class PrioritySlot: self.slot = slot def __hash__(self): - return hash((self.priority, self.slot)) + return hash((self.priority, self.slot)) def __eq__(self, other): - return (self.priority, self.slot) == (other.priority, other.slot) + return (self.priority, self.slot) == (other.priority, other.slot) def __lt__(self, other): - return (self.priority, self.slot) < (other.priority, other.slot) + return (self.priority, self.slot) < (other.priority, other.slot) def __str__(self): - return '_'.join([text_type(self.priority), _pathable(text_type(self.slot))]) + return '_'.join([text_type(self.priority), + _pathable(text_type(self.slot))]) class PriorityAsTupleQueue(PriorityQueue): @@ -135,9 +135,10 @@ class SlotBasedPriorityQueue(object): slot = _scheduler_slot(request) is_new = False if slot not in self.pqueues: - is_new = True self.pqueues[slot] = PriorityAsTupleQueue(self.qfactory) - self.pqueues[slot].push(request, PrioritySlot(priority=priority, slot=slot)) + queue = self.pqueues[slot] + is_new = queue.is_empty() + queue.push(request, PrioritySlot(priority=priority, slot=slot)) return slot, is_new def close(self): @@ -152,36 +153,6 @@ class SlotBasedPriorityQueue(object): return sum(len(x) for x in self.pqueues.values()) if self.pqueues else 0 -class RoundRobinPriorityQueue(SlotBasedPriorityQueue): - - def __init__(self, qfactory, startprios={}): - super(RoundRobinPriorityQueue, self).__init__(qfactory, startprios) - self._slots = deque() - for slot in self.pqueues: - self._slots.append(slot) - - def push(self, request, priority): - slot, is_new = self.push_slot(request, priority) - if is_new: - self._slots.append(slot) - - def pop(self): - if not self._slots: - return - - slot = self._slots.popleft() - request, is_empty = self.pop_slot(slot) - - if not is_empty: - self._slots.append(slot) - - return request - - def close(self): - self._slots.clear() - return super(RoundRobinPriorityQueue, self).close() - - class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): _DOWNLOADER_AWARE_PQ_ID = 'DOWNLOADER_AWARE_PQ_ID' @@ -191,7 +162,8 @@ class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): return cls(crawler, qfactory, startprios) def __init__(self, crawler, qfactory, startprios={}): - super(DownloaderAwarePriorityQueue, self).__init__(qfactory, startprios) + super(DownloaderAwarePriorityQueue, self).__init__(qfactory, + startprios) self._slots = {slot: 0 for slot in self.pqueues} crawler.signals.connect(self.on_response_download, signal=response_downloaded) @@ -208,7 +180,7 @@ class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): return request.meta.get(self._DOWNLOADER_AWARE_PQ_ID, None) == id(self) def pop(self): - slots = [(d, s) for s,d in self._slots.items() if s in self.pqueues] + slots = [(d, s) for s, d in self._slots.items() if s in self.pqueues] if not slots: return diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index fd86e8d8c..e1cf5842d 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -1,4 +1,3 @@ -import contextlib import shutil import tempfile import unittest @@ -10,15 +9,18 @@ from scrapy.pqueues import _scheduler_slot_read, _scheduler_slot_write from scrapy.signals import request_reached_downloader, response_downloaded from scrapy.spiders import Spider + class MockCrawler(Crawler): def __init__(self, priority_queue_cls, jobdir): - settings = dict(LOG_UNSERIALIZABLE_REQUESTS=False, - SCHEDULER_DISK_QUEUE='scrapy.squeues.PickleLifoDiskQueue', - SCHEDULER_MEMORY_QUEUE='scrapy.squeues.LifoMemoryQueue', - SCHEDULER_PRIORITY_QUEUE=priority_queue_cls, - JOBDIR=jobdir, - DUPEFILTER_CLASS='scrapy.dupefilters.BaseDupeFilter') + settings = dict( + LOG_UNSERIALIZABLE_REQUESTS=False, + SCHEDULER_DISK_QUEUE='scrapy.squeues.PickleLifoDiskQueue', + SCHEDULER_MEMORY_QUEUE='scrapy.squeues.LifoMemoryQueue', + SCHEDULER_PRIORITY_QUEUE=priority_queue_cls, + JOBDIR=jobdir, + DUPEFILTER_CLASS='scrapy.dupefilters.BaseDupeFilter' + ) super(MockCrawler, self).__init__(Spider, settings) @@ -82,7 +84,8 @@ class BaseSchedulerInMemoryTester(SchedulerHandler): while self.scheduler.has_pending_requests(): priorities.append(self.scheduler.next_request().priority) - self.assertEqual(priorities, sorted([x[1] for x in _PRIORITIES], key=lambda x: -x)) + self.assertEqual(priorities, + sorted([x[1] for x in _PRIORITIES], key=lambda x: -x)) class BaseSchedulerOnDiskTester(SchedulerHandler): @@ -134,7 +137,8 @@ class BaseSchedulerOnDiskTester(SchedulerHandler): while self.scheduler.has_pending_requests(): priorities.append(self.scheduler.next_request().priority) - self.assertEqual(priorities, sorted([x[1] for x in _PRIORITIES], key=lambda x: -x)) + self.assertEqual(priorities, + sorted([x[1] for x in _PRIORITIES], key=lambda x: -x)) class TestSchedulerInMemory(BaseSchedulerInMemoryTester, unittest.TestCase): @@ -153,75 +157,15 @@ _SLOTS = [("http://foo.com/a", 'a'), ("http://foo.com/f", 'c')] -class TestSchedulerWithRoundRobinInMemory(BaseSchedulerInMemoryTester, unittest.TestCase): - priority_queue_cls = 'scrapy.pqueues.RoundRobinPriorityQueue' +class TestMigration(unittest.TestCase): - def test_round_robin(self): - for url, slot in _SLOTS: - request = Request(url) - _scheduler_slot_write(request, slot) - self.scheduler.enqueue_request(request) + def setUp(self): + self.tmpdir = tempfile.mkdtemp() - slots = list() - while self.scheduler.has_pending_requests(): - slots.append(_scheduler_slot_read(self.scheduler.next_request())) + def tearDown(self): + shutil.rmtree(self.tmpdir) - for i in range(0, len(_SLOTS), 2): - self.assertNotEqual(slots[i], slots[i+1]) - - def test_is_meta_set(self): - url = "http://foo.com/a" - request = Request(url) - if _scheduler_slot_read(request): - _scheduler_slot_write(request, None) - self.scheduler.enqueue_request(request) - self.assertIsNotNone(_scheduler_slot_read(request, None), None) - - -class TestSchedulerWithRoundRobinOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): - priority_queue_cls = 'scrapy.pqueues.RoundRobinPriorityQueue' - - def test_round_robin(self): - for url, slot in _SLOTS: - request = Request(url) - _scheduler_slot_write(request, slot) - self.scheduler.enqueue_request(request) - - self.close_scheduler() - self.create_scheduler() - - slots = list() - while self.scheduler.has_pending_requests(): - slots.append(_scheduler_slot_read(self.scheduler.next_request())) - - for i in range(0, len(_SLOTS), 2): - self.assertNotEqual(slots[i], slots[i+1]) - - def test_is_meta_set(self): - url = "http://foo.com/a" - request = Request(url) - if _scheduler_slot_read(request): - _scheduler_slot_write(request, None) - self.scheduler.enqueue_request(request) - - self.close_scheduler() - self.create_scheduler() - - self.assertIsNotNone(_scheduler_slot_read(request, None), None) - - -@contextlib.contextmanager -def mkdtemp(): - dir = tempfile.mkdtemp() - try: - yield dir - finally: - shutil.rmtree(dir) - - -def _migration(): - - with mkdtemp() as tmp_dir: + def _migration(self, tmp_dir): prev_scheduler_handler = SchedulerHandler() prev_scheduler_handler.priority_queue_cls = 'queuelib.PriorityQueue' prev_scheduler_handler.jobdir = tmp_dir @@ -232,18 +176,18 @@ def _migration(): prev_scheduler_handler.close_scheduler() next_scheduler_handler = SchedulerHandler() - next_scheduler_handler.priority_queue_cls = 'scrapy.pqueues.RoundRobinPriorityQueue' + next_scheduler_handler.priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' next_scheduler_handler.jobdir = tmp_dir next_scheduler_handler.create_scheduler() - -class TestMigration(unittest.TestCase): def test_migration(self): - self.assertRaises(ValueError, _migration) + with self.assertRaises(ValueError): + self._migration(self.tmpdir) -class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, unittest.TestCase): +class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, + unittest.TestCase): priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' def test_logic(self): @@ -266,10 +210,12 @@ class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, unit self.assertEqual(len(slots), len(_SLOTS)) for request in requests: - self.mock_crawler.signals.send_catch_log(signal=response_downloaded, - request=request, - response=None, - spider=self.spider) + self.mock_crawler.signals.send_catch_log( + signal=response_downloaded, + request=request, + response=None, + spider=self.spider + ) unique_slots = len(set(s for _, s in _SLOTS)) for i in range(0, len(_SLOTS), unique_slots): @@ -277,8 +223,10 @@ class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, unit self.assertEqual(len(part), len(set(part))) -class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): +class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, + unittest.TestCase): priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' + def test_logic(self): for url, slot in _SLOTS: request = Request(url) @@ -304,10 +252,12 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, unittest self.assertEqual(len(slots), len(_SLOTS)) for request in requests: - self.mock_crawler.signals.send_catch_log(signal=response_downloaded, - request=request, - response=None, - spider=self.spider) + self.mock_crawler.signals.send_catch_log( + signal=response_downloaded, + request=request, + response=None, + spider=self.spider + ) unique_slots = len(set(s for _, s in _SLOTS)) for i in range(0, len(_SLOTS), unique_slots): From 7efba101946af93397ec3c2323b920644e20ce04 Mon Sep 17 00:00:00 2001 From: Lucy Wang Date: Mon, 10 Dec 2018 14:44:15 +0800 Subject: [PATCH 06/38] remove "sudo: false" now that travis no longer supports it https://changelog.travis-ci.com/deprecation-container-based-linux-build-environment-82037 --- .travis.yml | 1 - 1 file changed, 1 deletion(-) diff --git a/.travis.yml b/.travis.yml index 4218d13bf..08b0bf119 100644 --- a/.travis.yml +++ b/.travis.yml @@ -1,5 +1,4 @@ language: python -sudo: false branches: only: - master From 0e06b9a81672ec432d2fccc3cbacc823ea47b656 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 14 Dec 2018 14:35:18 +0000 Subject: [PATCH 07/38] use urlparse_cached where it is possible --- scrapy/pqueues.py | 24 ++++++++++++++++++++---- 1 file changed, 20 insertions(+), 4 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 287a8de35..ff7ec8c8a 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -8,6 +8,7 @@ from queuelib import PriorityQueue from scrapy.core.downloader import Downloader from scrapy.http import Request from scrapy.signals import request_reached_downloader, response_downloaded +from scrapy.utils.httpobj import urlparse_cached logger = logging.getLogger(__name__) @@ -41,11 +42,26 @@ def _scheduler_slot_write(request, slot): def _scheduler_slot(request): - slot = _scheduler_slot_read(request, None) - if slot is None: - url = _get_from_request(request, 'url') + if isinstance(request, dict): + meta = request.get('meta', dict()) + elif isinstance(request, Request): + meta = request.meta + else: + raise ValueError('Bad type of request "%s"' % (request.__class__, )) + + slot = meta.get(SCHEDULER_SLOT_META_KEY, None) + + if slot is not None: + return slot + + if isinstance(request, dict): + url = request.get('url', None) slot = urlparse(url).hostname or '' - _scheduler_slot_write(request, slot) + elif isinstance(request, Request): + url = request.url + slot = urlparse_cached(request).hostname or '' + + meta[SCHEDULER_SLOT_META_KEY] = slot return slot From 484927b08caff66ea622f8553468c831154df30a Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 14 Dec 2018 14:38:28 +0000 Subject: [PATCH 08/38] less complex implementation --- scrapy/pqueues.py | 9 ++------- 1 file changed, 2 insertions(+), 7 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index ff7ec8c8a..538678345 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -28,16 +28,11 @@ def _get_from_request(request, key, default=None): def _scheduler_slot_read(request, default=None): - meta = _get_from_request(request, 'meta', dict()) - slot = meta.get(SCHEDULER_SLOT_META_KEY, default) - return slot + return request.meta.get(SCHEDULER_SLOT_META_KEY, default) def _scheduler_slot_write(request, slot): - meta = _get_from_request(request, 'meta', None) - if not isinstance(meta, dict): - raise ValueError('No meta attribute in %s' % (request, )) - meta[SCHEDULER_SLOT_META_KEY] = slot + request.meta[SCHEDULER_SLOT_META_KEY] = slot def _scheduler_slot(request): From 6af964cc0b47c570e035a3486b9f8aebd349bd84 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 14 Dec 2018 14:54:24 +0000 Subject: [PATCH 09/38] common indentation for comment --- scrapy/pqueues.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 538678345..31e90ff12 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -96,9 +96,9 @@ class PrioritySlot: 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 + 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=()): From a46613afa8acd136f4ba62df2ced2f3c87679512 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 14 Dec 2018 14:55:06 +0000 Subject: [PATCH 10/38] use regular comments --- scrapy/pqueues.py | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 31e90ff12..75fc198d0 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -64,10 +64,8 @@ def _scheduler_slot(request): def _pathable(x): pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' for c in x]) - """ - as we replace some letters we can get collision for different slots - add we add unique part - """ + # as we replace some letters we can get collision for different slots + # add we add unique part unique_slot = hashlib.md5(x.encode('utf8')).hexdigest() return '-'.join([pathable_slot, unique_slot]) From a23e1894b3a09e1daf49dd9592546b2d21bc9a72 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 14 Dec 2018 16:18:34 +0000 Subject: [PATCH 11/38] Fix boto problem another way to fix boto problem Revert "fix for travis ci based on https://github.com/boto/boto/issues/3717" This reverts commit 150d2564ff0ea994652da7f5be333d72e0b38d93. fix for travis ci based on https://github.com/boto/boto/issues/3717 --- tests/requirements-py2.txt | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/requirements-py2.txt b/tests/requirements-py2.txt index 790f29d34..f5bcfda60 100644 --- a/tests/requirements-py2.txt +++ b/tests/requirements-py2.txt @@ -11,3 +11,4 @@ testfixtures # optional for shell wrapper tests bpython ipython<6.0 +google-compute-engine From d970be64cc47c382bd615cd547e7e94c17e27b48 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Mon, 17 Dec 2018 13:52:11 +0000 Subject: [PATCH 12/38] Integration test integration testing only everything is working, not logic of PQ use method create slot attribute in constructor corect class for test case stop crawler in teardown method use class correct entity naming python 2 adaptation integration test with crawler and spider --- tests/test_scheduler.py | 46 +++++++++++++++++++++++++++++++++++++---- 1 file changed, 42 insertions(+), 4 deletions(-) diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index e1cf5842d..9bdc82b30 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -2,12 +2,17 @@ import shutil import tempfile import unittest +from twisted.internet import defer +from twisted.trial.unittest import TestCase + from scrapy.crawler import Crawler from scrapy.core.scheduler import Scheduler from scrapy.http import Request from scrapy.pqueues import _scheduler_slot_read, _scheduler_slot_write from scrapy.signals import request_reached_downloader, response_downloaded from scrapy.spiders import Spider +from scrapy.utils.test import get_crawler +from tests.mockserver import MockServer class MockCrawler(Crawler): @@ -223,6 +228,13 @@ class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, self.assertEqual(len(part), len(set(part))) +def _is_slots_unique(base_slots, result_slots): + unique_slots = len(set(s for _, s in base_slots)) + for i in range(0, len(result_slots), unique_slots): + part = result_slots[i:i + unique_slots] + assert len(part) == len(set(part)) + + class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' @@ -259,7 +271,33 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, spider=self.spider ) - unique_slots = len(set(s for _, s in _SLOTS)) - for i in range(0, len(_SLOTS), unique_slots): - part = slots[i:i + unique_slots] - self.assertEqual(len(part), len(set(part))) + _is_slots_unique(_SLOTS, slots) + + +class StartUrlsSpider(Spider): + + def __init__(self, start_urls): + self.start_urls = start_urls + + +class TestIntegrationWithDownloaderAwareOnDisk(TestCase): + def setUp(self): + self.crawler = get_crawler( + StartUrlsSpider, + {'SCHEDULER_PRIORITY_QUEUE': 'scrapy.pqueues.DownloaderAwarePriorityQueue', + 'DUPEFILTER_CLASS': 'scrapy.dupefilters.BaseDupeFilter'} + ) + + @defer.inlineCallbacks + def tearDown(self): + yield self.crawler.stop() + + @defer.inlineCallbacks + def test_integration_downloader_aware_priority_queue(self): + with MockServer() as mockserver: + + url = mockserver.url("/status?n=200", is_secure=False) + slots = [url] * 6 + yield self.crawler.crawl(slots) + self.assertEqual(self.crawler.stats.get_value('downloader/response_count'), + len(slots)) From 7d3175ac8433f964ebbb80ebd67f9899cf059100 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Thu, 20 Dec 2018 19:23:23 -0300 Subject: [PATCH 13/38] Fix boto import error under Jessie testing environment --- .travis.yml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/.travis.yml b/.travis.yml index 08b0bf119..a201f97b1 100644 --- a/.travis.yml +++ b/.travis.yml @@ -42,6 +42,11 @@ install: virtualenv --python="$PYPY_VERSION/bin/pypy3" "$HOME/virtualenvs/$PYPY_VERSION" source "$HOME/virtualenvs/$PYPY_VERSION/bin/activate" fi + if [ "$TOXENV" = "jessie" ]; then + # Not used directly but allows boto GCE plugins to load. + # https://github.com/GoogleCloudPlatform/compute-image-packages/issues/262 + pip install google-compute-engine + fi - pip install -U tox twine wheel codecov script: tox From 6ff2574c277ba1eda31fb43f86f43d5b7b4bef09 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Thu, 20 Dec 2018 19:39:29 -0300 Subject: [PATCH 14/38] Needs to be installed within tox env --- .travis.yml | 5 ----- tox.ini | 3 +++ 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/.travis.yml b/.travis.yml index a201f97b1..08b0bf119 100644 --- a/.travis.yml +++ b/.travis.yml @@ -42,11 +42,6 @@ install: virtualenv --python="$PYPY_VERSION/bin/pypy3" "$HOME/virtualenvs/$PYPY_VERSION" source "$HOME/virtualenvs/$PYPY_VERSION/bin/activate" fi - if [ "$TOXENV" = "jessie" ]; then - # Not used directly but allows boto GCE plugins to load. - # https://github.com/GoogleCloudPlatform/compute-image-packages/issues/262 - pip install google-compute-engine - fi - pip install -U tox twine wheel codecov script: tox diff --git a/tox.ini b/tox.ini index e5543fe2a..0c0f8f7b7 100644 --- a/tox.ini +++ b/tox.ini @@ -51,6 +51,9 @@ deps = cssselect==0.9.1 zope.interface==4.1.1 -rtests/requirements-py2.txt +# Not used directly but allows boto GCE plugins to load. +# https://github.com/GoogleCloudPlatform/compute-image-packages/issues/262 + google-compute-engine==2.8.12 [testenv:trunk] basepython = python2.7 From 4163a7a1c7ac11c8d4db70f371c26181b90d8dfd Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 21 Dec 2018 09:10:32 +0000 Subject: [PATCH 15/38] no need for this --- tests/requirements-py2.txt | 1 - 1 file changed, 1 deletion(-) diff --git a/tests/requirements-py2.txt b/tests/requirements-py2.txt index f5bcfda60..790f29d34 100644 --- a/tests/requirements-py2.txt +++ b/tests/requirements-py2.txt @@ -11,4 +11,3 @@ testfixtures # optional for shell wrapper tests bpython ipython<6.0 -google-compute-engine From 987c2ae4a964e45120c245235c9b0c49dc36b71f Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Tue, 25 Dec 2018 09:13:09 +0000 Subject: [PATCH 16/38] test ip concurrency incompatibility with DAPQ --- tests/test_scheduler.py | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 9bdc82b30..17b706bd7 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -301,3 +301,20 @@ class TestIntegrationWithDownloaderAwareOnDisk(TestCase): yield self.crawler.crawl(slots) self.assertEqual(self.crawler.stats.get_value('downloader/response_count'), len(slots)) + + +class TestIncompatibility(unittest.TestCase): + + def _incompatible(self): + settings = dict( + SCHEDULER_PRIORITY_QUEUE='scrapy.pqueues.DownloaderAwarePriorityQueue', + CONCURRENT_REQUESTS_PER_IP=1 + ) + crawler = Crawler(Spider, settings) + scheduler = Scheduler.from_crawler(crawler) + spider = Spider(name='spider') + scheduler.open(spider) + + def test_incompatibility(self): + with self.assertRaises(ValueError): + self._incompatible() From 8e8ce301b1a56e40f7e9c322a7b73b8dcfcefc43 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Tue, 25 Dec 2018 09:14:09 +0000 Subject: [PATCH 17/38] check CONCURRENT_REQUESTS_PER_IP is not set --- scrapy/pqueues.py | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 75fc198d0..d9effc9d1 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -171,6 +171,14 @@ class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): return cls(crawler, qfactory, startprios) def __init__(self, crawler, qfactory, startprios={}): + ip_concurrency_key = 'CONCURRENT_REQUESTS_PER_IP' + ip_concurrency = crawler.settings.getint(ip_concurrency_key, 0) + + if ip_concurrency > 0: + raise ValueError('"%s" does not support %s=%d' % (self.__class__, + ip_concurrency_key, + ip_concurrency)) + super(DownloaderAwarePriorityQueue, self).__init__(qfactory, startprios) self._slots = {slot: 0 for slot in self.pqueues} From 338b78d796de6c93af0f4bcb762f82f5a14b87cd Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Tue, 25 Dec 2018 09:44:20 +0000 Subject: [PATCH 18/38] Add documentation add section to broad-crawl topic reword in accord with broad-crawl topic add documentation for new priority queue --- docs/topics/broad-crawls.rst | 11 +++++++++++ docs/topics/settings.rst | 7 ++++++- 2 files changed, 17 insertions(+), 1 deletion(-) diff --git a/docs/topics/broad-crawls.rst b/docs/topics/broad-crawls.rst index eb02086dc..37f7a8748 100644 --- a/docs/topics/broad-crawls.rst +++ b/docs/topics/broad-crawls.rst @@ -39,6 +39,17 @@ you need to keep in mind when using Scrapy for doing broad crawls, along with concrete suggestions of Scrapy settings to tune in order to achieve an efficient broad crawl. +Use proper :setting:`SCHEDULER_PRIORITY_QUEUE` +============================================== + +Default scrapy's scheduler priority queue is ``'queuelib.PriorityQueue'``. +It works best during single domain crawl. And it does not work well with crawling +many different domains in parallel + +To apply recommended priority queue use:: + + SCHEDULER_PRIORITY_QUEUE = 'scrapy.pqueues.DownloaderAwarePriorityQueue' + Increase concurrency ==================== diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 47b6cf13d..7b9ff7e39 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1144,7 +1144,12 @@ SCHEDULER_PRIORITY_QUEUE ------------------------ Default: ``'queuelib.PriorityQueue'`` -Type of priority queue used by scheduler. +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`` +does not work together with :setting:`CONCURRENT_REQUESTS_PER_IP`. .. setting:: SPIDER_CONTRACTS From 90934959d07db881aec5933fd7b77bcd2dccfa4f Mon Sep 17 00:00:00 2001 From: Mikhail Korobov Date: Thu, 27 Dec 2018 17:12:24 +0500 Subject: [PATCH 19/38] actually apply __slots__ suggestion [wip] refactoring * SlotPriorityQueues doesn't care about objects inside, it is now just a container for multiple priority queues * assorted variable renames * don't inherit DownloaderAwarePriorityQueue from SlotBasedPriorityQueue * apply @whalebot-helmsman's suggestions for __slots__ and meta issues more bike-shedding * remove mutable default arguments * more verbose variable names remove unneeded code * PriorityAsTupleQueue.is_empty does the same as len(self) == 0 * custom PriorityAsTupleQueue.close is not needed after a switch to namedtuples * is_new and is_empty return values are unused * "url" local variable is unused PrioritySlot.__str__ shouldn't return unicode in Python 2 also, do some bike-shedding: _pathable -> _path_safe use namedtuple for PrioritySlot cleanup: _get_from_request does the same here Request.meta is always a dict --- scrapy/pqueues.py | 180 ++++++++++++++++++---------------------- tests/test_scheduler.py | 2 +- 2 files changed, 82 insertions(+), 100 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index d9effc9d1..3ef896b99 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -1,6 +1,6 @@ import hashlib import logging -from six import text_type +from collections import namedtuple from six.moves.urllib.parse import urlparse from queuelib import PriorityQueue @@ -17,12 +17,12 @@ logger = logging.getLogger(__name__) SCHEDULER_SLOT_META_KEY = Downloader.DOWNLOAD_SLOT -def _get_from_request(request, key, default=None): +def _get_request_meta(request): if isinstance(request, dict): - return request.get(key, default) + return request.setdefault('meta', {}) if isinstance(request, Request): - return getattr(request, key, default) + return request.meta raise ValueError('Bad type of request "%s"' % (request.__class__, )) @@ -35,15 +35,8 @@ def _scheduler_slot_write(request, slot): request.meta[SCHEDULER_SLOT_META_KEY] = slot -def _scheduler_slot(request): - - if isinstance(request, dict): - meta = request.get('meta', dict()) - elif isinstance(request, Request): - meta = request.meta - else: - raise ValueError('Bad type of request "%s"' % (request.__class__, )) - +def _set_scheduler_slot(request): + meta = _get_request_meta(request) slot = meta.get(SCHEDULER_SLOT_META_KEY, None) if slot is not None: @@ -53,43 +46,29 @@ def _scheduler_slot(request): url = request.get('url', None) slot = urlparse(url).hostname or '' elif isinstance(request, Request): - url = request.url slot = urlparse_cached(request).hostname or '' meta[SCHEDULER_SLOT_META_KEY] = slot - return slot -def _pathable(x): - pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' for c in x]) - +def _path_safe(text): + """ Return a filesystem-safe version of a string ``text`` """ + pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' + for c in text]) # as we replace some letters we can get collision for different slots # add we add unique part - unique_slot = hashlib.md5(x.encode('utf8')).hexdigest() - + unique_slot = hashlib.md5(text.encode('utf8')).hexdigest() return '-'.join([pathable_slot, unique_slot]) -class PrioritySlot: - __slots__ = ('priority', 'slot') - - def __init__(self, priority=0, slot=None): - self.priority = priority - self.slot = slot - - def __hash__(self): - return hash((self.priority, self.slot)) - - def __eq__(self, other): - return (self.priority, self.slot) == (other.priority, other.slot) - - def __lt__(self, other): - return (self.priority, self.slot) < (other.priority, other.slot) +class PrioritySlot(namedtuple("PrioritySlot", ["priority", "slot"])): + """ ``(priority, slot)`` tuple which uses a path-safe slot name + when converting to str """ + __slots__ = () def __str__(self): - return '_'.join([text_type(self.priority), - _pathable(text_type(self.slot))]) + return '%s_%s' % (self.priority, _path_safe(str(self.slot))) class PriorityAsTupleQueue(PriorityQueue): @@ -99,78 +78,65 @@ class PriorityAsTupleQueue(PriorityQueue): 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, - [PrioritySlot(priority=p[0], slot=p[1]) for p in startprios] - ) - - def close(self): - startprios = super(PriorityAsTupleQueue, self).close() - return [(s.priority, s.slot) for s in startprios] - - def is_empty(self): - return not self.queues or len(self) == 0 + qfactory=qfactory, + startprios=startprios) -class SlotBasedPriorityQueue(object): +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__(self, qfactory, startprios={}): - self.pqueues = dict() # slot -> priority queue - self.qfactory = qfactory # factory for creating new internal queues - - if not startprios: - return - - if not isinstance(startprios, dict): - raise ValueError("Looks like your priorities file malforfemed. " - "Possible reason: You run scrapy with previous " - "version. Interrupted it. Updated scrapy. And " - "run again.") - - for slot, prios in startprios.items(): - self.pqueues[slot] = PriorityAsTupleQueue(self.qfactory, prios) + ``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) def pop_slot(self, slot): + """ Pop an object from a priority queue for this slot """ queue = self.pqueues[slot] request = queue.pop() - is_empty = queue.is_empty() - if is_empty: + if len(queue) == 0: del self.pqueues[slot] + return request - return request, is_empty - - def push_slot(self, request, priority): - slot = _scheduler_slot(request) - is_new = False + 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] = PriorityAsTupleQueue(self.qfactory) + self.pqueues[slot] = self.pqfactory() queue = self.pqueues[slot] - is_new = queue.is_empty() - queue.push(request, PrioritySlot(priority=priority, slot=slot)) - return slot, is_new + queue.push(obj, priority) def close(self): - startprios = dict() - for slot, queue in self.pqueues.items(): - prios = queue.close() - startprios[slot] = prios + active = {slot: queue.close() + for slot, queue in self.pqueues.items()} self.pqueues.clear() - return startprios + return active def __len__(self): return sum(len(x) for x in self.pqueues.values()) if self.pqueues else 0 + def __contains__(self, slot): + return slot in self.pqueues -class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): + +class DownloaderAwarePriorityQueue(object): _DOWNLOADER_AWARE_PQ_ID = 'DOWNLOADER_AWARE_PQ_ID' @classmethod - def from_crawler(cls, crawler, qfactory, startprios={}): + def from_crawler(cls, crawler, qfactory, startprios=None): return cls(crawler, qfactory, startprios) - def __init__(self, crawler, qfactory, startprios={}): + def __init__(self, crawler, qfactory, startprios=None): ip_concurrency_key = 'CONCURRENT_REQUESTS_PER_IP' ip_concurrency = crawler.settings.getint(ip_concurrency_key, 0) @@ -179,16 +145,25 @@ class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): ip_concurrency_key, ip_concurrency)) - super(DownloaderAwarePriorityQueue, self).__init__(qfactory, - startprios) - self._slots = {slot: 0 for slot in self.pqueues} + 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) + + self._active_downloads = {slot: 0 for slot in self._slot_pqueues.pqueues} crawler.signals.connect(self.on_response_download, signal=response_downloaded) crawler.signals.connect(self.on_request_reached_downloader, signal=request_reached_downloader) def mark(self, request): - meta = _get_from_request(request, 'meta', None) + 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) @@ -197,39 +172,46 @@ class DownloaderAwarePriorityQueue(SlotBasedPriorityQueue): return request.meta.get(self._DOWNLOADER_AWARE_PQ_ID, None) == id(self) def pop(self): - slots = [(d, s) for s, d in self._slots.items() if s in self.pqueues] + slots = [(active_downloads, slot) + for slot, active_downloads in self._active_downloads.items() + if slot in self._slot_pqueues] if not slots: return slot = min(slots)[1] - request, _ = self.pop_slot(slot) + request = self._slot_pqueues.pop_slot(slot) self.mark(request) return request def push(self, request, priority): - slot, _ = self.push_slot(request, priority) - if slot not in self._slots: - self._slots[slot] = 0 + slot = _set_scheduler_slot(request) + priority_slot = PrioritySlot(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 def on_response_download(self, response, request, spider): if not self.check_mark(request): return slot = _scheduler_slot_read(request) - if slot not in self._slots or self._slots[slot] <= 0: + if slot not in self._active_downloads or self._active_downloads[slot] <= 0: raise ValueError('Get response for wrong slot "%s"' % (slot, )) - self._slots[slot] = self._slots[slot] - 1 - if self._slots[slot] == 0 and slot not in self.pqueues: - del self._slots[slot] + self._active_downloads[slot] = self._active_downloads[slot] - 1 + if self._active_downloads[slot] == 0 and slot not in self._slot_pqueues: + del self._active_downloads[slot] def on_request_reached_downloader(self, request, spider): if not self.check_mark(request): return slot = _scheduler_slot_read(request) - self._slots[slot] = self._slots.get(slot, 0) + 1 + self._active_downloads[slot] = self._active_downloads.get(slot, 0) + 1 def close(self): - self._slots.clear() - return super(DownloaderAwarePriorityQueue, self).close() + self._active_downloads.clear() + return self._slot_pqueues.close() + + def __len__(self): + return len(self._slot_pqueues) diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 17b706bd7..5dd35f45c 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -260,7 +260,7 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, ) requests.append(request) - self.assertEqual(self.scheduler.mqs._slots, {}) + self.assertEqual(self.scheduler.mqs._active_downloads, {}) self.assertEqual(len(slots), len(_SLOTS)) for request in requests: From 757f53a32461ef0c3d2fe4caf64197f67271b5f3 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Wed, 9 Jan 2019 10:00:13 +0000 Subject: [PATCH 20/38] Address Lucy's comments add tests to check correctness of slot setermination unmark requests after downloading shorter better exception message --- scrapy/pqueues.py | 15 ++++++++++++--- tests/test_scheduler.py | 4 ++-- 2 files changed, 14 insertions(+), 5 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 3ef896b99..d8eed010f 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -36,6 +36,12 @@ 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 + """ meta = _get_request_meta(request) slot = meta.get(SCHEDULER_SLOT_META_KEY, None) @@ -141,9 +147,8 @@ class DownloaderAwarePriorityQueue(object): ip_concurrency = crawler.settings.getint(ip_concurrency_key, 0) if ip_concurrency > 0: - raise ValueError('"%s" does not support %s=%d' % (self.__class__, - ip_concurrency_key, - ip_concurrency)) + raise ValueError('"%s" does not support setting %s' % (self.__class__, + ip_concurrency_key)) def pqfactory(startprios=()): return PriorityAsTupleQueue(qfactory, startprios) @@ -171,6 +176,9 @@ class DownloaderAwarePriorityQueue(object): def check_mark(self, request): return request.meta.get(self._DOWNLOADER_AWARE_PQ_ID, None) == id(self) + def unmark(self, request): + del request.meta[self._DOWNLOADER_AWARE_PQ_ID] + def pop(self): slots = [(active_downloads, slot) for slot, active_downloads in self._active_downloads.items() @@ -194,6 +202,7 @@ class DownloaderAwarePriorityQueue(object): def on_response_download(self, response, request, spider): if not self.check_mark(request): return + self.unmark(request) slot = _scheduler_slot_read(request) if slot not in self._active_downloads or self._active_downloads[slot] <= 0: diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 5dd35f45c..3fb70a110 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -248,8 +248,8 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, self.close_scheduler() self.create_scheduler() - slots = list() - requests = list() + slots = [] + requests = [] while self.scheduler.has_pending_requests(): request = self.scheduler.next_request() slots.append(_scheduler_slot_read(request)) From 3b1db71dac8716878ff1b94ee0d1095e5c80795f Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Wed, 9 Jan 2019 12:14:40 +0000 Subject: [PATCH 21/38] New signal update signature documentation for new signal utilize new signal correct signal handler signature emit new signal test another signal new signal rename test file faster test rename test case tests for signal emitting in bad cases --- docs/topics/signals.rst | 17 +++++++++ scrapy/core/downloader/__init__.py | 3 ++ scrapy/pqueues.py | 6 +-- scrapy/signals.py | 1 + tests/test_request_left.py | 59 ++++++++++++++++++++++++++++++ tests/test_scheduler.py | 8 ++-- 6 files changed, 86 insertions(+), 8 deletions(-) create mode 100644 tests/test_request_left.py diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index ff07b9d55..f13e8270c 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -295,6 +295,23 @@ request_reached_downloader :param spider: the spider that yielded the request :type spider: :class:`~scrapy.spiders.Spider` object +request_left_downloader +--------------------------- + +.. signal:: request_left_downloader +.. function:: request_left_downloader(request, spider) + + Sent when a :class:`~scrapy.http.Request` left downloader even in case of + failure. + + The signal does not support returning deferreds from their handlers. + + :param request: the request that reached downloader + :type request: :class:`~scrapy.http.Request` object + + :param spider: the spider that yielded the request + :type spider: :class:`~scrapy.spiders.Spider` object + response_received ----------------- diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index 4695d75f4..d856a2f37 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -188,6 +188,9 @@ class Downloader(object): def finish_transferring(_): slot.transferring.remove(request) self._process_queue(spider, slot) + self.signals.send_catch_log(signal=signals.request_left_downloader, + request=request, + spider=spider) return _ return dfd.addBoth(finish_transferring) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index d8eed010f..6a9feb599 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -7,7 +7,7 @@ from queuelib import PriorityQueue from scrapy.core.downloader import Downloader from scrapy.http import Request -from scrapy.signals import request_reached_downloader, response_downloaded +from scrapy.signals import request_reached_downloader, request_left_downloader from scrapy.utils.httpobj import urlparse_cached @@ -163,7 +163,7 @@ class DownloaderAwarePriorityQueue(object): self._active_downloads = {slot: 0 for slot in self._slot_pqueues.pqueues} crawler.signals.connect(self.on_response_download, - signal=response_downloaded) + signal=request_left_downloader) crawler.signals.connect(self.on_request_reached_downloader, signal=request_reached_downloader) @@ -199,7 +199,7 @@ class DownloaderAwarePriorityQueue(object): if slot not in self._active_downloads: self._active_downloads[slot] = 0 - def on_response_download(self, response, request, spider): + def on_response_download(self, request, spider): if not self.check_mark(request): return self.unmark(request) diff --git a/scrapy/signals.py b/scrapy/signals.py index c0e4bb74e..2ea986b8c 100644 --- a/scrapy/signals.py +++ b/scrapy/signals.py @@ -14,6 +14,7 @@ spider_error = object() request_scheduled = object() request_dropped = object() request_reached_downloader = object() +request_left_downloader = object() response_received = object() response_downloaded = object() item_scraped = object() diff --git a/tests/test_request_left.py b/tests/test_request_left.py new file mode 100644 index 000000000..ddeca0499 --- /dev/null +++ b/tests/test_request_left.py @@ -0,0 +1,59 @@ +from twisted.internet import defer +from twisted.trial.unittest import TestCase +from scrapy.signals import request_left_downloader +from scrapy.spiders import Spider +from scrapy.utils.test import get_crawler +from tests.mockserver import MockServer + +class SignalCatcherSpider(Spider): + name = 'signal_catcher' + + def __init__(self, crawler, url, *args, **kwargs): + super(SignalCatcherSpider, self).__init__(*args, **kwargs) + crawler.signals.connect(self.on_response_download, + signal=request_left_downloader) + self.catched_times = 0 + self.start_urls = [url] + + @classmethod + def from_crawler(cls, crawler, *args, **kwargs): + spider = cls(crawler, *args, **kwargs) + return spider + + def on_response_download(self, request, spider): + self.catched_times = self.catched_times + 1 + + +class TestCatching(TestCase): + + def setUp(self): + self.mockserver = MockServer() + self.mockserver.__enter__() + + def tearDown(self): + self.mockserver.__exit__(None, None, None) + + @defer.inlineCallbacks + def test_success(self): + crawler = get_crawler(SignalCatcherSpider) + yield crawler.crawl(self.mockserver.url("/status?n=200")) + self.assertEqual(crawler.spider.catched_times, 1) + + @defer.inlineCallbacks + def test_timeout(self): + crawler = get_crawler(SignalCatcherSpider, + {'DOWNLOAD_TIMEOUT': 0.1}) + yield crawler.crawl(self.mockserver.url("/delay?n=0.2")) + self.assertEqual(crawler.spider.catched_times, 1) + + @defer.inlineCallbacks + def test_disconnect(self): + crawler = get_crawler(SignalCatcherSpider) + yield crawler.crawl(self.mockserver.url("/drop")) + self.assertEqual(crawler.spider.catched_times, 1) + + @defer.inlineCallbacks + def test_noconnect(self): + crawler = get_crawler(SignalCatcherSpider) + yield crawler.crawl('http://thereisdefinetelynosuchdomain.com') + self.assertEqual(crawler.spider.catched_times, 1) diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 3fb70a110..1bcc1e5a8 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -9,7 +9,7 @@ from scrapy.crawler import Crawler from scrapy.core.scheduler import Scheduler from scrapy.http import Request from scrapy.pqueues import _scheduler_slot_read, _scheduler_slot_write -from scrapy.signals import request_reached_downloader, response_downloaded +from scrapy.signals import request_reached_downloader, request_left_downloader from scrapy.spiders import Spider from scrapy.utils.test import get_crawler from tests.mockserver import MockServer @@ -216,9 +216,8 @@ class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, for request in requests: self.mock_crawler.signals.send_catch_log( - signal=response_downloaded, + signal=request_left_downloader, request=request, - response=None, spider=self.spider ) @@ -265,9 +264,8 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, for request in requests: self.mock_crawler.signals.send_catch_log( - signal=response_downloaded, + signal=request_left_downloader, request=request, - response=None, spider=self.spider ) From 83eb5376458ce1d444e8ad7911730ee0c58c8544 Mon Sep 17 00:00:00 2001 From: Mikhail Korobov Date: Thu, 17 Jan 2019 07:38:15 +0500 Subject: [PATCH 22/38] assorted cleanups: comments, docstrings, etc scheduler cleanup Scheduler no longer converts requests to dicts; PriorityQueue instances always work with Request instances; converting Requests to dicts is now Priority Queue responsibility. minor cleanup --- docs/topics/settings.rst | 6 +- scrapy/core/scheduler.py | 99 +++++++++++++---- scrapy/pqueues.py | 158 +++++++++++++++------------- scrapy/settings/default_settings.py | 2 +- scrapy/squeues.py | 11 +- 5 files changed, 175 insertions(+), 101 deletions(-) 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 From 443fb98a4776f4196662bb48918f1471758b7ae7 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Tue, 5 Mar 2019 12:44:07 +0000 Subject: [PATCH 23/38] Use downloader directly rename variable remove old write function remove unused imports remove old read function remove unused function use mock methods mock downloader close downloader add parse method use new PQ class create mock downloader use downloader directly remove mark/unmark mechanism --- scrapy/pqueues.py | 103 ++++++++++------------------------------ tests/test_scheduler.py | 87 +++++++++++++++++++++------------ 2 files changed, 81 insertions(+), 109 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 622f6bbc5..0681e6729 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -5,42 +5,11 @@ from collections import namedtuple 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 -from scrapy.utils.httpobj import urlparse_cached logger = logging.getLogger(__name__) -SCHEDULER_SLOT_META_KEY = Downloader.DOWNLOAD_SLOT - - -def _scheduler_slot_read(request, default=None): - return request.meta.get(SCHEDULER_SLOT_META_KEY, default) - - -def _scheduler_slot_write(request, slot): - request.meta[SCHEDULER_SLOT_META_KEY] = slot - - -def _set_scheduler_slot(request): - """ - >>> request = Request('http://example.com') - >>> _set_scheduler_slot(request) - 'example.com' - >>> _scheduler_slot_read(request) - 'example.com' - """ - slot = _scheduler_slot_read(request, None) - if slot is not None: - return slot - slot = urlparse_cached(request).hostname or '' - _scheduler_slot_write(request, slot) - return slot - - def _path_safe(text): """ Return a filesystem-safe version of a string ``text`` """ pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' @@ -138,6 +107,25 @@ class ScrapyPriorityQueue(PriorityQueue): return request +class DownloaderInterface(object): + + def __init__(self, crawler): + self.downloader = crawler.engine.downloader + + def stats(self, possible_slots): + return [(self._active_downloads(slot), slot) + for slot in possible_slots] + + def get_slot_key(self, request): + return self.downloader._get_slot_key(request, None) + + def _active_downloads(self, slot): + """ Return a number of requests in a Downloader for a given slot """ + if slot not in self.downloader.slots: + return 0 + return len(self.downloader.slots[slot].active) + + class DownloaderAwarePriorityQueue(object): """ PriorityQueue which takes Downlaoder activity in account: domains (slots) with the least amount of active downloads are dequeued @@ -170,68 +158,25 @@ class DownloaderAwarePriorityQueue(object): def pqfactory(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): - 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) - - def unmark(self, request): - del request.meta[self._DOWNLOADER_AWARE_PQ_ID] + self._downloader_interface = DownloaderInterface(crawler) def pop(self): - slots = [(active_downloads, slot) - for slot, active_downloads in self._active_downloads.items() - if slot in self._slot_pqueues] + stats = self._downloader_interface.stats(self._slot_pqueues.pqueues) - if not slots: + if not stats: return - slot = min(slots)[1] + slot = min(stats)[1] request = self._slot_pqueues.pop_slot(slot) - self.mark(request) return request def push(self, request, priority): - slot = _set_scheduler_slot(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._active_downloads: - self._active_downloads[slot] = 0 - - def on_response_download(self, request, spider): - if not self.check_mark(request): - return - self.unmark(request) - - slot = _scheduler_slot_read(request) - if slot not in self._active_downloads or self._active_downloads[slot] <= 0: - 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] - - def on_request_reached_downloader(self, request, spider): - if not self.check_mark(request): - return - - slot = _scheduler_slot_read(request) - self._active_downloads[slot] = self._active_downloads.get(slot, 0) + 1 def close(self): - self._active_downloads.clear() active = self._slot_pqueues.close() return {slot: [p.priority for p in startprios] for slot, startprios in active.items()} diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 1bcc1e5a8..75c0b7530 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -1,20 +1,50 @@ import shutil import tempfile import unittest +import collections from twisted.internet import defer from twisted.trial.unittest import TestCase from scrapy.crawler import Crawler +from scrapy.core.downloader import Downloader from scrapy.core.scheduler import Scheduler from scrapy.http import Request -from scrapy.pqueues import _scheduler_slot_read, _scheduler_slot_write -from scrapy.signals import request_reached_downloader, request_left_downloader from scrapy.spiders import Spider +from scrapy.utils.httpobj import urlparse_cached from scrapy.utils.test import get_crawler from tests.mockserver import MockServer +MockEngine = collections.namedtuple('MockEngine', ['downloader']) +MockSlot = collections.namedtuple('MockSlot', ['active']) + + +class MockDownloader: + def __init__(self): + self.slots = dict() + + def _set_slot_key(self, slot, request, spider): + request.meta[Downloader.DOWNLOAD_SLOT] = slot + + def _get_slot_key(self, request, spider): + if Downloader.DOWNLOAD_SLOT in request.meta: + return request.meta[Downloader.DOWNLOAD_SLOT] + + return urlparse_cached(request).hostname or '' + + def increment(self, slot_key): + slot = self.slots.setdefault(slot_key, MockSlot(active=list())) + slot.active.append(1) + + def decrement(self, slot_key): + slot = self.slots.get(slot_key) + slot.active.pop() + + def close(self): + pass + + class MockCrawler(Crawler): def __init__(self, priority_queue_cls, jobdir): @@ -27,6 +57,7 @@ class MockCrawler(Crawler): DUPEFILTER_CLASS='scrapy.dupefilters.BaseDupeFilter' ) super(MockCrawler, self).__init__(Spider, settings) + self.engine = MockEngine(downloader=MockDownloader()) class SchedulerHandler: @@ -42,6 +73,7 @@ class SchedulerHandler: def close_scheduler(self): self.scheduler.close('finished') self.mock_crawler.stop() + self.mock_crawler.engine.downloader.close() def setUp(self): self.create_scheduler() @@ -147,11 +179,11 @@ class BaseSchedulerOnDiskTester(SchedulerHandler): class TestSchedulerInMemory(BaseSchedulerInMemoryTester, unittest.TestCase): - priority_queue_cls = 'queuelib.PriorityQueue' + priority_queue_cls = 'scrapy.pqueues.ScrapyPriorityQueue' class TestSchedulerOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): - priority_queue_cls = 'queuelib.PriorityQueue' + priority_queue_cls = 'scrapy.pqueues.ScrapyPriorityQueue' _SLOTS = [("http://foo.com/a", 'a'), @@ -172,7 +204,7 @@ class TestMigration(unittest.TestCase): def _migration(self, tmp_dir): prev_scheduler_handler = SchedulerHandler() - prev_scheduler_handler.priority_queue_cls = 'queuelib.PriorityQueue' + prev_scheduler_handler.priority_queue_cls = 'scrapy.pqueues.ScrapyPriorityQueue' prev_scheduler_handler.jobdir = tmp_dir prev_scheduler_handler.create_scheduler() @@ -196,30 +228,25 @@ class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' def test_logic(self): + downloader = self.mock_crawler.engine.downloader for url, slot in _SLOTS: request = Request(url) - _scheduler_slot_write(request, slot) + downloader._set_slot_key(slot, request, None) self.scheduler.enqueue_request(request) slots = list() requests = list() while self.scheduler.has_pending_requests(): request = self.scheduler.next_request() - slots.append(_scheduler_slot_read(request)) - self.mock_crawler.signals.send_catch_log( - signal=request_reached_downloader, - request=request, - spider=self.spider - ) + slot = downloader._get_slot_key(request, None) + slots.append(slot) + downloader.increment(slot) requests.append(request) self.assertEqual(len(slots), len(_SLOTS)) for request in requests: - self.mock_crawler.signals.send_catch_log( - signal=request_left_downloader, - request=request, - spider=self.spider - ) + slot = downloader._get_slot_key(request, None) + self.mock_crawler.engine.downloader.decrement(slot) unique_slots = len(set(s for _, s in _SLOTS)) for i in range(0, len(_SLOTS), unique_slots): @@ -239,9 +266,11 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' def test_logic(self): + downloader = self.mock_crawler.engine.downloader + for url, slot in _SLOTS: request = Request(url) - _scheduler_slot_write(request, slot) + downloader._set_slot_key(slot, request, None) self.scheduler.enqueue_request(request) self.close_scheduler() @@ -249,27 +278,22 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, slots = [] requests = [] + downloader = self.mock_crawler.engine.downloader while self.scheduler.has_pending_requests(): request = self.scheduler.next_request() - slots.append(_scheduler_slot_read(request)) - self.mock_crawler.signals.send_catch_log( - signal=request_reached_downloader, - request=request, - spider=self.spider - ) + slot = downloader._get_slot_key(request, None) + slots.append(slot) + downloader.increment(slot) requests.append(request) - self.assertEqual(self.scheduler.mqs._active_downloads, {}) self.assertEqual(len(slots), len(_SLOTS)) for request in requests: - self.mock_crawler.signals.send_catch_log( - signal=request_left_downloader, - request=request, - spider=self.spider - ) + slot = downloader._get_slot_key(request, None) + downloader.decrement(slot) _is_slots_unique(_SLOTS, slots) + self.assertEqual(sum(len(s.active) for s in downloader.slots.values()), 0) class StartUrlsSpider(Spider): @@ -277,6 +301,9 @@ class StartUrlsSpider(Spider): def __init__(self, start_urls): self.start_urls = start_urls + def parse(self, response): + pass + class TestIntegrationWithDownloaderAwareOnDisk(TestCase): def setUp(self): From 989bba6cb340fcc1ddb32e75ade567864d8b3884 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Thu, 7 Mar 2019 09:00:14 +0000 Subject: [PATCH 24/38] Revert "new signal" This reverts commit 646164fd7d6dd52061804d2df7424cff929bf739. remove tests Revert "emit new signal" This reverts commit fcde0c6880678957a76af6083b6248f430a00fcf. Revert "documentation for new signal" This reverts commit 8aeb9f696ece95c16499a96767a7afa3d9c4abf4. --- docs/topics/signals.rst | 17 --------- scrapy/core/downloader/__init__.py | 3 -- scrapy/signals.py | 1 - tests/test_request_left.py | 59 ------------------------------ 4 files changed, 80 deletions(-) delete mode 100644 tests/test_request_left.py diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index f13e8270c..ff07b9d55 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -295,23 +295,6 @@ request_reached_downloader :param spider: the spider that yielded the request :type spider: :class:`~scrapy.spiders.Spider` object -request_left_downloader ---------------------------- - -.. signal:: request_left_downloader -.. function:: request_left_downloader(request, spider) - - Sent when a :class:`~scrapy.http.Request` left downloader even in case of - failure. - - The signal does not support returning deferreds from their handlers. - - :param request: the request that reached downloader - :type request: :class:`~scrapy.http.Request` object - - :param spider: the spider that yielded the request - :type spider: :class:`~scrapy.spiders.Spider` object - response_received ----------------- diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index d856a2f37..4695d75f4 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -188,9 +188,6 @@ class Downloader(object): def finish_transferring(_): slot.transferring.remove(request) self._process_queue(spider, slot) - self.signals.send_catch_log(signal=signals.request_left_downloader, - request=request, - spider=spider) return _ return dfd.addBoth(finish_transferring) diff --git a/scrapy/signals.py b/scrapy/signals.py index 2ea986b8c..c0e4bb74e 100644 --- a/scrapy/signals.py +++ b/scrapy/signals.py @@ -14,7 +14,6 @@ spider_error = object() request_scheduled = object() request_dropped = object() request_reached_downloader = object() -request_left_downloader = object() response_received = object() response_downloaded = object() item_scraped = object() diff --git a/tests/test_request_left.py b/tests/test_request_left.py deleted file mode 100644 index ddeca0499..000000000 --- a/tests/test_request_left.py +++ /dev/null @@ -1,59 +0,0 @@ -from twisted.internet import defer -from twisted.trial.unittest import TestCase -from scrapy.signals import request_left_downloader -from scrapy.spiders import Spider -from scrapy.utils.test import get_crawler -from tests.mockserver import MockServer - -class SignalCatcherSpider(Spider): - name = 'signal_catcher' - - def __init__(self, crawler, url, *args, **kwargs): - super(SignalCatcherSpider, self).__init__(*args, **kwargs) - crawler.signals.connect(self.on_response_download, - signal=request_left_downloader) - self.catched_times = 0 - self.start_urls = [url] - - @classmethod - def from_crawler(cls, crawler, *args, **kwargs): - spider = cls(crawler, *args, **kwargs) - return spider - - def on_response_download(self, request, spider): - self.catched_times = self.catched_times + 1 - - -class TestCatching(TestCase): - - def setUp(self): - self.mockserver = MockServer() - self.mockserver.__enter__() - - def tearDown(self): - self.mockserver.__exit__(None, None, None) - - @defer.inlineCallbacks - def test_success(self): - crawler = get_crawler(SignalCatcherSpider) - yield crawler.crawl(self.mockserver.url("/status?n=200")) - self.assertEqual(crawler.spider.catched_times, 1) - - @defer.inlineCallbacks - def test_timeout(self): - crawler = get_crawler(SignalCatcherSpider, - {'DOWNLOAD_TIMEOUT': 0.1}) - yield crawler.crawl(self.mockserver.url("/delay?n=0.2")) - self.assertEqual(crawler.spider.catched_times, 1) - - @defer.inlineCallbacks - def test_disconnect(self): - crawler = get_crawler(SignalCatcherSpider) - yield crawler.crawl(self.mockserver.url("/drop")) - self.assertEqual(crawler.spider.catched_times, 1) - - @defer.inlineCallbacks - def test_noconnect(self): - crawler = get_crawler(SignalCatcherSpider) - yield crawler.crawl('http://thereisdefinetelynosuchdomain.com') - self.assertEqual(crawler.spider.catched_times, 1) From 8afffb7234b282dd8bd28eec2e4eb8e3f86b5723 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 22 Mar 2019 09:12:23 +0000 Subject: [PATCH 25/38] Tests Cleanup add doctest for function no need in this variables move common assertion inside function rename variable rename variables rename function use function this is not a method of public API correct name for test Update docs/topics/settings.rst Co-Authored-By: whalebot-helmsman --- docs/topics/settings.rst | 4 +- tests/test_scheduler.py | 82 ++++++++++++++++++++++------------------ 2 files changed, 48 insertions(+), 38 deletions(-) diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 6e13e64d6..cf454f4ec 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1146,9 +1146,9 @@ 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 +``scrapy.pqueues.DownloaderAwarePriorityQueue`` works better than ``scrapy.pqueues.ScrapyPriorityQueue`` when you crawl many different -domains in parallel. But ``scrapy.pqueues.DownloaderAwarePriorityQueue`` +domains in parallel. But currently ``scrapy.pqueues.DownloaderAwarePriorityQueue`` does not work together with :setting:`CONCURRENT_REQUESTS_PER_IP`. .. setting:: SPIDER_CONTRACTS diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 75c0b7530..eaf748d35 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -24,9 +24,6 @@ class MockDownloader: def __init__(self): self.slots = dict() - def _set_slot_key(self, slot, request, spider): - request.meta[Downloader.DOWNLOAD_SLOT] = slot - def _get_slot_key(self, request, spider): if Downloader.DOWNLOAD_SLOT in request.meta: return request.meta[Downloader.DOWNLOAD_SLOT] @@ -186,12 +183,12 @@ class TestSchedulerOnDisk(BaseSchedulerOnDiskTester, unittest.TestCase): priority_queue_cls = 'scrapy.pqueues.ScrapyPriorityQueue' -_SLOTS = [("http://foo.com/a", 'a'), - ("http://foo.com/b", 'a'), - ("http://foo.com/c", 'b'), - ("http://foo.com/d", 'b'), - ("http://foo.com/e", 'c'), - ("http://foo.com/f", 'c')] +_URLS_WITH_SLOTS = [("http://foo.com/a", 'a'), + ("http://foo.com/b", 'a'), + ("http://foo.com/c", 'b'), + ("http://foo.com/d", 'b'), + ("http://foo.com/e", 'c'), + ("http://foo.com/f", 'c')] class TestMigration(unittest.TestCase): @@ -228,37 +225,52 @@ class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' def test_logic(self): - downloader = self.mock_crawler.engine.downloader - for url, slot in _SLOTS: + for url, slot in _URLS_WITH_SLOTS: request = Request(url) - downloader._set_slot_key(slot, request, None) + request.meta[Downloader.DOWNLOAD_SLOT] = slot self.scheduler.enqueue_request(request) - slots = list() + downloader = self.mock_crawler.engine.downloader + dequeued_slots = list() requests = list() while self.scheduler.has_pending_requests(): request = self.scheduler.next_request() slot = downloader._get_slot_key(request, None) - slots.append(slot) + dequeued_slots.append(slot) downloader.increment(slot) requests.append(request) - self.assertEqual(len(slots), len(_SLOTS)) for request in requests: slot = downloader._get_slot_key(request, None) self.mock_crawler.engine.downloader.decrement(slot) - unique_slots = len(set(s for _, s in _SLOTS)) - for i in range(0, len(_SLOTS), unique_slots): - part = slots[i:i + unique_slots] - self.assertEqual(len(part), len(set(part))) + self.assertTrue(_is_scheduling_fair(list(s for u, s in _URLS_WITH_SLOTS), + dequeued_slots)) -def _is_slots_unique(base_slots, result_slots): - unique_slots = len(set(s for _, s in base_slots)) - for i in range(0, len(result_slots), unique_slots): - part = result_slots[i:i + unique_slots] - assert len(part) == len(set(part)) +def _is_scheduling_fair(enqueued_slots, dequeued_slots): + """ + We enqueued same number of requests for every slot. + Assert correct order, e.g. + + >>> enqueued = ['a', 'b', 'c'] * 2 + >>> correct = ['a', 'c', 'b', 'b', 'a', 'c'] + >>> incorrect = ['a', 'a', 'b', 'c', 'c', 'b'] + >>> _is_scheduling_fair(enqueued, correct) + True + >>> _is_scheduling_fair(enqueued, incorrect) + False + """ + if len(dequeued_slots) != len(enqueued_slots): + return False + + slots_number = len(set(enqueued_slots)) + for i in range(0, len(dequeued_slots), slots_number): + part = dequeued_slots[i:i + slots_number] + if len(part) != len(set(part)): + return False + + return True class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, @@ -266,33 +278,31 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' def test_logic(self): - downloader = self.mock_crawler.engine.downloader - for url, slot in _SLOTS: + for url, slot in _URLS_WITH_SLOTS: request = Request(url) - downloader._set_slot_key(slot, request, None) + request.meta[Downloader.DOWNLOAD_SLOT] = slot self.scheduler.enqueue_request(request) self.close_scheduler() self.create_scheduler() - slots = [] + dequeued_slots = list() requests = [] downloader = self.mock_crawler.engine.downloader while self.scheduler.has_pending_requests(): request = self.scheduler.next_request() slot = downloader._get_slot_key(request, None) - slots.append(slot) + dequeued_slots.append(slot) downloader.increment(slot) requests.append(request) - self.assertEqual(len(slots), len(_SLOTS)) - for request in requests: slot = downloader._get_slot_key(request, None) downloader.decrement(slot) - _is_slots_unique(_SLOTS, slots) + self.assertTrue(_is_scheduling_fair(list(s for u, s in _URLS_WITH_SLOTS), + dequeued_slots)) self.assertEqual(sum(len(s.active) for s in downloader.slots.values()), 0) @@ -305,7 +315,7 @@ class StartUrlsSpider(Spider): pass -class TestIntegrationWithDownloaderAwareOnDisk(TestCase): +class TestIntegrationWithDownloaderAwareInMemory(TestCase): def setUp(self): self.crawler = get_crawler( StartUrlsSpider, @@ -322,10 +332,10 @@ class TestIntegrationWithDownloaderAwareOnDisk(TestCase): with MockServer() as mockserver: url = mockserver.url("/status?n=200", is_secure=False) - slots = [url] * 6 - yield self.crawler.crawl(slots) + start_urls = [url] * 6 + yield self.crawler.crawl(start_urls) self.assertEqual(self.crawler.stats.get_value('downloader/response_count'), - len(slots)) + len(start_urls)) class TestIncompatibility(unittest.TestCase): From df574de8cc5c58618f6075ca3afb14059a9e30ed Mon Sep 17 00:00:00 2001 From: Lucy Wang Date: Sat, 23 Mar 2019 00:54:39 +0800 Subject: [PATCH 26/38] improve tests and fix some lint warnings (#6) * refactor downloader-aware test cases * fix lint * add doctest for _path_safe * remove unused code * better doctest --- scrapy/pqueues.py | 12 +++++++-- tests/test_scheduler.py | 57 ++++++++++++++++------------------------- 2 files changed, 32 insertions(+), 37 deletions(-) diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index 0681e6729..6ecd1b51a 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -11,7 +11,16 @@ logger = logging.getLogger(__name__) def _path_safe(text): - """ Return a filesystem-safe version of a string ``text`` """ + """ + Return a filesystem-safe version of a string ``text`` + + >>> _path_safe('simple.org').startswith('simple.org') + True + >>> _path_safe('dash-underscore_.org').startswith('dash-underscore_.org') + True + >>> _path_safe('some@symbol?').startswith('some_symbol_') + True + """ pathable_slot = "".join([c if c.isalnum() or c in '-._' else '_' for c in text]) # as we replace some letters we can get collision for different slots @@ -131,7 +140,6 @@ class DownloaderAwarePriorityQueue(object): 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): diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index eaf748d35..e0e3600e5 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -20,7 +20,7 @@ MockEngine = collections.namedtuple('MockEngine', ['downloader']) MockSlot = collections.namedtuple('MockSlot', ['active']) -class MockDownloader: +class MockDownloader(object): def __init__(self): self.slots = dict() @@ -57,7 +57,7 @@ class MockCrawler(Crawler): self.engine = MockEngine(downloader=MockDownloader()) -class SchedulerHandler: +class SchedulerHandler(object): priority_queue_cls = None jobdir = None @@ -220,34 +220,6 @@ class TestMigration(unittest.TestCase): self._migration(self.tmpdir) -class TestSchedulerWithDownloaderAwareInMemory(BaseSchedulerInMemoryTester, - unittest.TestCase): - priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' - - def test_logic(self): - for url, slot in _URLS_WITH_SLOTS: - request = Request(url) - request.meta[Downloader.DOWNLOAD_SLOT] = slot - self.scheduler.enqueue_request(request) - - downloader = self.mock_crawler.engine.downloader - dequeued_slots = list() - requests = list() - while self.scheduler.has_pending_requests(): - request = self.scheduler.next_request() - slot = downloader._get_slot_key(request, None) - dequeued_slots.append(slot) - downloader.increment(slot) - requests.append(request) - - for request in requests: - slot = downloader._get_slot_key(request, None) - self.mock_crawler.engine.downloader.decrement(slot) - - self.assertTrue(_is_scheduling_fair(list(s for u, s in _URLS_WITH_SLOTS), - dequeued_slots)) - - def _is_scheduling_fair(enqueued_slots, dequeued_slots): """ We enqueued same number of requests for every slot. @@ -273,31 +245,33 @@ def _is_scheduling_fair(enqueued_slots, dequeued_slots): return True -class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, - unittest.TestCase): +class DownloaderAwareSchedulerTestMixin(object): priority_queue_cls = 'scrapy.pqueues.DownloaderAwarePriorityQueue' + reopen = False def test_logic(self): - for url, slot in _URLS_WITH_SLOTS: request = Request(url) request.meta[Downloader.DOWNLOAD_SLOT] = slot self.scheduler.enqueue_request(request) - self.close_scheduler() - self.create_scheduler() + if self.reopen: + self.close_scheduler() + self.create_scheduler() dequeued_slots = list() requests = [] downloader = self.mock_crawler.engine.downloader while self.scheduler.has_pending_requests(): request = self.scheduler.next_request() + # pylint: disable=protected-access slot = downloader._get_slot_key(request, None) dequeued_slots.append(slot) downloader.increment(slot) requests.append(request) for request in requests: + # pylint: disable=protected-access slot = downloader._get_slot_key(request, None) downloader.decrement(slot) @@ -306,10 +280,23 @@ class TestSchedulerWithDownloaderAwareOnDisk(BaseSchedulerOnDiskTester, self.assertEqual(sum(len(s.active) for s in downloader.slots.values()), 0) +class TestSchedulerWithDownloaderAwareInMemory(DownloaderAwareSchedulerTestMixin, + BaseSchedulerInMemoryTester, + unittest.TestCase): + pass + + +class TestSchedulerWithDownloaderAwareOnDisk(DownloaderAwareSchedulerTestMixin, + BaseSchedulerOnDiskTester, + unittest.TestCase): + reopen = True + + class StartUrlsSpider(Spider): def __init__(self, start_urls): self.start_urls = start_urls + super(StartUrlsSpider, self).__init__(start_urls) def parse(self, response): pass From 31b8a6b33aed9e77a4d37a5c83b1545202207cad Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Mon, 25 Mar 2019 08:53:15 +0000 Subject: [PATCH 27/38] report warnings --- tests/test_crawler.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 268948a70..d9ec9ee8d 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -1,5 +1,4 @@ import logging -import tempfile import warnings from twisted.internet import defer @@ -37,7 +36,11 @@ class CrawlerTestCase(BaseCrawlerTest): self.assertIsInstance(spiders, sl_cls) self.crawler.spiders - self.assertEqual(len(w), 1, "Warn deprecated access only once") + is_one_warning = len(w) == 1 + if not is_one_warning: + for warning in w: + print(warning) + self.assertTrue(is_one_warning, "Warn deprecated access only once") def test_populate_spidercls_settings(self): spider_settings = {'TEST1': 'spider', 'TEST2': 'spider'} From 73e4ff5304d273404a147d06726a8ae8cae1c925 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Mon, 25 Mar 2019 13:48:58 +0000 Subject: [PATCH 28/38] report warnings --- tests/test_crawler.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 8c4bbe0d9..e811c5757 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -182,8 +182,12 @@ class CrawlerRunnerTestCase(BaseCrawlerTest): 'SPIDER_MANAGER_CLASS': 'tests.test_crawler.CustomSpiderLoader' }) self.assertIsInstance(runner.spider_loader, CustomSpiderLoader) - self.assertEqual(len(w), 1) + is_one_warning = len(w) == 1 + if not is_one_warning: + for warning in w: + print(warning) self.assertIn('Please use SPIDER_LOADER_CLASS', str(w[0].message)) + self.assertTrue(is_one_warning) def test_crawl_rejects_spider_objects(self): with raises(ValueError): From 845bae6637239c859c9952c23f42902e36d10f6b Mon Sep 17 00:00:00 2001 From: Mikhail Korobov Date: Wed, 27 Mar 2019 08:49:19 +0000 Subject: [PATCH 29/38] Update docs/topics/broad-crawls.rst Co-Authored-By: whalebot-helmsman --- docs/topics/broad-crawls.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/topics/broad-crawls.rst b/docs/topics/broad-crawls.rst index 37f7a8748..64c8883b1 100644 --- a/docs/topics/broad-crawls.rst +++ b/docs/topics/broad-crawls.rst @@ -42,7 +42,7 @@ efficient broad crawl. Use proper :setting:`SCHEDULER_PRIORITY_QUEUE` ============================================== -Default scrapy's scheduler priority queue is ``'queuelib.PriorityQueue'``. +Default scrapy's scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQueue'``. It works best during single domain crawl. And it does not work well with crawling many different domains in parallel From 46b9ab0c58354deb1045c20f3bc061526d69f356 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Fri, 29 Mar 2019 10:28:36 +0000 Subject: [PATCH 30/38] Update docs/topics/broad-crawls.rst Co-Authored-By: whalebot-helmsman --- docs/topics/broad-crawls.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/topics/broad-crawls.rst b/docs/topics/broad-crawls.rst index 64c8883b1..68a24a4d2 100644 --- a/docs/topics/broad-crawls.rst +++ b/docs/topics/broad-crawls.rst @@ -42,7 +42,7 @@ efficient broad crawl. Use proper :setting:`SCHEDULER_PRIORITY_QUEUE` ============================================== -Default scrapy's scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQueue'``. +Scrapy’s default scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQueue'``. It works best during single domain crawl. And it does not work well with crawling many different domains in parallel From e3df6be360a58f016e31d5bfa2e04cd2e5d1965b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Fri, 29 Mar 2019 10:28:52 +0000 Subject: [PATCH 31/38] Update docs/topics/broad-crawls.rst Co-Authored-By: whalebot-helmsman --- docs/topics/broad-crawls.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/topics/broad-crawls.rst b/docs/topics/broad-crawls.rst index 68a24a4d2..b149d7f4a 100644 --- a/docs/topics/broad-crawls.rst +++ b/docs/topics/broad-crawls.rst @@ -43,7 +43,7 @@ Use proper :setting:`SCHEDULER_PRIORITY_QUEUE` ============================================== Scrapy’s default scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQueue'``. -It works best during single domain crawl. And it does not work well with crawling +It works best during single-domain crawl. It does not work well with crawling many different domains in parallel To apply recommended priority queue use:: From bd228f1d962c7f4759536d8cda278857de7d5234 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Fri, 29 Mar 2019 10:29:04 +0000 Subject: [PATCH 32/38] Update docs/topics/broad-crawls.rst Co-Authored-By: whalebot-helmsman --- docs/topics/broad-crawls.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/topics/broad-crawls.rst b/docs/topics/broad-crawls.rst index b149d7f4a..a01f28248 100644 --- a/docs/topics/broad-crawls.rst +++ b/docs/topics/broad-crawls.rst @@ -46,7 +46,7 @@ Scrapy’s default scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQ It works best during single-domain crawl. It does not work well with crawling many different domains in parallel -To apply recommended priority queue use:: +To apply the recommended priority queue use:: SCHEDULER_PRIORITY_QUEUE = 'scrapy.pqueues.DownloaderAwarePriorityQueue' From 1ee99e1f4240af6a7a72fe7c58b89d7bce1cd09e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Fri, 29 Mar 2019 10:29:15 +0000 Subject: [PATCH 33/38] Update docs/topics/settings.rst Co-Authored-By: whalebot-helmsman --- docs/topics/settings.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index ed94146f4..4a5439bfc 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1157,7 +1157,7 @@ SCHEDULER_PRIORITY_QUEUE ------------------------ Default: ``'scrapy.pqueues.ScrapyPriorityQueue'`` -Type of priority queue used by scheduler. Another available type is +Type of priority queue used by the scheduler. Another available type is ``scrapy.pqueues.DownloaderAwarePriorityQueue``. ``scrapy.pqueues.DownloaderAwarePriorityQueue`` works better than ``scrapy.pqueues.ScrapyPriorityQueue`` when you crawl many different From 2b4bcfaf494073520e84bbf301d5141a2e19a3e6 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 29 Mar 2019 10:30:26 +0000 Subject: [PATCH 34/38] remove comment --- scrapy/core/scheduler.py | 1 - 1 file changed, 1 deletion(-) diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index c385fafe1..9d0258db2 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -57,7 +57,6 @@ class Scheduler(object): 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'", From 554d8728227a9ea96e5ea3a8a4fd782d42fdbd66 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 29 Mar 2019 10:31:15 +0000 Subject: [PATCH 35/38] remove spacing --- scrapy/core/scheduler.py | 5 ----- 1 file changed, 5 deletions(-) diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 9d0258db2..d87d2ffdc 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -77,13 +77,8 @@ class Scheduler(object): def open(self, spider): self.spider = spider - - # 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): From f08f841d0bebd097358889c0c98f83f051828f15 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 29 Mar 2019 10:35:49 +0000 Subject: [PATCH 36/38] remove small single use method --- scrapy/core/scheduler.py | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index d87d2ffdc..975aede0c 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -101,7 +101,7 @@ class Scheduler(object): return True def next_request(self): - request = self._mqpop() + request = self.mqs.pop() if request: self.stats.inc_value('scheduler/dequeued/memory', spider=self.spider) else: @@ -141,9 +141,6 @@ class Scheduler(object): if self.dqs: return self.dqs.pop() - def _mqpop(self): - return self.mqs.pop() - def _newmq(self, priority): """ Factory for creating memory queues. """ return self.mqclass() From ef743983a98ae0891abf9aca4c9b19cb44861c49 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 29 Mar 2019 10:38:13 +0000 Subject: [PATCH 37/38] change wording --- docs/topics/broad-crawls.rst | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/topics/broad-crawls.rst b/docs/topics/broad-crawls.rst index a01f28248..6e50c0bc7 100644 --- a/docs/topics/broad-crawls.rst +++ b/docs/topics/broad-crawls.rst @@ -39,7 +39,7 @@ you need to keep in mind when using Scrapy for doing broad crawls, along with concrete suggestions of Scrapy settings to tune in order to achieve an efficient broad crawl. -Use proper :setting:`SCHEDULER_PRIORITY_QUEUE` +Use the right :setting:`SCHEDULER_PRIORITY_QUEUE` ============================================== Scrapy’s default scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQueue'``. @@ -96,7 +96,7 @@ When doing broad crawls you are often only interested in the crawl rates you get and any errors found. These stats are reported by Scrapy when using the ``INFO`` log level. In order to save CPU (and log storage requirements) you should not use ``DEBUG`` log level when preforming large broad crawls in -production. Using ``DEBUG`` level when developing your (broad) crawler may be +production. Using ``DEBUG`` level when developing your (broad) crawler may be fine though. To set the log level use:: From 1c6733454e14a3c237ed602b65ae5e0a8a78dee5 Mon Sep 17 00:00:00 2001 From: Vostretsov Nikita Date: Fri, 29 Mar 2019 10:44:55 +0000 Subject: [PATCH 38/38] added underlines --- docs/topics/broad-crawls.rst | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/topics/broad-crawls.rst b/docs/topics/broad-crawls.rst index 6e50c0bc7..b887b98af 100644 --- a/docs/topics/broad-crawls.rst +++ b/docs/topics/broad-crawls.rst @@ -40,7 +40,7 @@ concrete suggestions of Scrapy settings to tune in order to achieve an efficient broad crawl. Use the right :setting:`SCHEDULER_PRIORITY_QUEUE` -============================================== +================================================= Scrapy’s default scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQueue'``. It works best during single-domain crawl. It does not work well with crawling