mirror of https://github.com/scrapy/scrapy.git
Scheduler: minimal interface, API docs (#3559)
This commit is contained in:
parent
6837919798
commit
ddea6b7bfa
|
|
@ -227,6 +227,7 @@ Extending Scrapy
|
|||
topics/extensions
|
||||
topics/api
|
||||
topics/signals
|
||||
topics/scheduler
|
||||
topics/exporters
|
||||
|
||||
|
||||
|
|
@ -248,6 +249,9 @@ Extending Scrapy
|
|||
:doc:`topics/signals`
|
||||
See all available signals and how to work with them.
|
||||
|
||||
:doc:`topics/scheduler`
|
||||
Understand the scheduler component.
|
||||
|
||||
:doc:`topics/exporters`
|
||||
Quickly export your scraped items to a file (XML, CSV, etc).
|
||||
|
||||
|
|
|
|||
|
|
@ -87,8 +87,9 @@ of the system, and triggering events when certain actions occur. See the
|
|||
Scheduler
|
||||
---------
|
||||
|
||||
The Scheduler receives requests from the engine and enqueues them for feeding
|
||||
them later (also to the engine) when the engine requests them.
|
||||
The :ref:`scheduler <topics-scheduler>` receives requests from the engine and
|
||||
enqueues them for feeding them later (also to the engine) when the engine
|
||||
requests them.
|
||||
|
||||
.. _component-downloader:
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,34 @@
|
|||
.. _topics-scheduler:
|
||||
|
||||
=========
|
||||
Scheduler
|
||||
=========
|
||||
|
||||
.. module:: scrapy.core.scheduler
|
||||
|
||||
The scheduler component receives requests from the :ref:`engine <component-engine>`
|
||||
and stores them into persistent and/or non-persistent data structures.
|
||||
It also gets those requests and feeds them back to the engine when it
|
||||
asks for a next request to be downloaded.
|
||||
|
||||
|
||||
Overriding the default scheduler
|
||||
================================
|
||||
|
||||
You can use your own custom scheduler class by supplying its full
|
||||
Python path in the :setting:`SCHEDULER` setting.
|
||||
|
||||
|
||||
Minimal scheduler interface
|
||||
===========================
|
||||
|
||||
.. autoclass:: BaseScheduler
|
||||
:members:
|
||||
|
||||
|
||||
Default Scrapy scheduler
|
||||
========================
|
||||
|
||||
.. autoclass:: Scheduler
|
||||
:members:
|
||||
:special-members: __len__
|
||||
|
|
@ -1280,7 +1280,8 @@ SCHEDULER
|
|||
|
||||
Default: ``'scrapy.core.scheduler.Scheduler'``
|
||||
|
||||
The scheduler to use for crawling.
|
||||
The scheduler class to be used for crawling.
|
||||
See the :ref:`topics-scheduler` topic for details.
|
||||
|
||||
.. setting:: SCHEDULER_DEBUG
|
||||
|
||||
|
|
|
|||
|
|
@ -17,6 +17,7 @@ from scrapy import signals
|
|||
from scrapy.core.scraper import Scraper
|
||||
from scrapy.exceptions import DontCloseSpider, ScrapyDeprecationWarning
|
||||
from scrapy.http import Response, Request
|
||||
from scrapy.settings import BaseSettings
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.log import logformatter_adapter, failure_to_exc_info
|
||||
from scrapy.utils.misc import create_instance, load_object
|
||||
|
|
@ -73,12 +74,22 @@ class ExecutionEngine:
|
|||
self.spider: Optional[Spider] = None
|
||||
self.running = False
|
||||
self.paused = False
|
||||
self.scheduler_cls = load_object(crawler.settings["SCHEDULER"])
|
||||
self.scheduler_cls = self._get_scheduler_class(crawler.settings)
|
||||
downloader_cls = load_object(self.settings['DOWNLOADER'])
|
||||
self.downloader = downloader_cls(crawler)
|
||||
self.scraper = Scraper(crawler)
|
||||
self._spider_closed_callback = spider_closed_callback
|
||||
|
||||
def _get_scheduler_class(self, settings: BaseSettings) -> type:
|
||||
from scrapy.core.scheduler import BaseScheduler
|
||||
scheduler_cls = load_object(settings["SCHEDULER"])
|
||||
if not issubclass(scheduler_cls, BaseScheduler):
|
||||
raise TypeError(
|
||||
f"The provided scheduler class ({settings['SCHEDULER']})"
|
||||
" does not fully implement the scheduler interface"
|
||||
)
|
||||
return scheduler_cls
|
||||
|
||||
@inlineCallbacks
|
||||
def start(self) -> Deferred:
|
||||
if self.running:
|
||||
|
|
@ -301,7 +312,8 @@ class ExecutionEngine:
|
|||
start_requests = yield self.scraper.spidermw.process_start_requests(start_requests, spider)
|
||||
self.slot = Slot(start_requests, close_if_idle, nextcall, scheduler)
|
||||
self.spider = spider
|
||||
yield scheduler.open(spider)
|
||||
if hasattr(scheduler, "open"):
|
||||
yield scheduler.open(spider)
|
||||
yield self.scraper.open_spider(spider)
|
||||
self.crawler.stats.open_spider(spider)
|
||||
yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider)
|
||||
|
|
@ -345,8 +357,9 @@ class ExecutionEngine:
|
|||
dfd.addBoth(lambda _: self.scraper.close_spider(spider))
|
||||
dfd.addErrback(log_failure('Scraper close failure'))
|
||||
|
||||
dfd.addBoth(lambda _: self.slot.scheduler.close(reason))
|
||||
dfd.addErrback(log_failure('Scheduler close failure'))
|
||||
if hasattr(self.slot.scheduler, "close"):
|
||||
dfd.addBoth(lambda _: self.slot.scheduler.close(reason))
|
||||
dfd.addErrback(log_failure("Scheduler close failure"))
|
||||
|
||||
dfd.addBoth(lambda _: self.signals.send_catch_log_deferred(
|
||||
signal=signals.spider_closed, spider=spider, reason=reason,
|
||||
|
|
|
|||
|
|
@ -1,42 +1,179 @@
|
|||
import os
|
||||
import json
|
||||
import logging
|
||||
from os.path import join, exists
|
||||
import os
|
||||
from abc import abstractmethod
|
||||
from os.path import exists, join
|
||||
from typing import Optional, Type, TypeVar
|
||||
|
||||
from scrapy.utils.misc import load_object, create_instance
|
||||
from twisted.internet.defer import Deferred
|
||||
|
||||
from scrapy.crawler import Crawler
|
||||
from scrapy.http.request import Request
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.job import job_dir
|
||||
from scrapy.utils.misc import create_instance, load_object
|
||||
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Scheduler:
|
||||
class BaseSchedulerMeta(type):
|
||||
"""
|
||||
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.
|
||||
Metaclass to check scheduler classes against the necessary interface
|
||||
"""
|
||||
def __init__(self, dupefilter, jobdir=None, dqclass=None, mqclass=None,
|
||||
logunser=False, stats=None, pqclass=None, crawler=None):
|
||||
def __instancecheck__(cls, instance):
|
||||
return cls.__subclasscheck__(type(instance))
|
||||
|
||||
def __subclasscheck__(cls, subclass):
|
||||
return (
|
||||
hasattr(subclass, "has_pending_requests") and callable(subclass.has_pending_requests)
|
||||
and hasattr(subclass, "enqueue_request") and callable(subclass.enqueue_request)
|
||||
and hasattr(subclass, "next_request") and callable(subclass.next_request)
|
||||
)
|
||||
|
||||
|
||||
class BaseScheduler(metaclass=BaseSchedulerMeta):
|
||||
"""
|
||||
The scheduler component is responsible for storing requests received from
|
||||
the engine, and feeding them back upon request (also to the engine).
|
||||
|
||||
The original sources of said requests are:
|
||||
|
||||
* Spider: ``start_requests`` method, requests created for URLs in the ``start_urls`` attribute, request callbacks
|
||||
* Spider middleware: ``process_spider_output`` and ``process_spider_exception`` methods
|
||||
* Downloader middleware: ``process_request``, ``process_response`` and ``process_exception`` methods
|
||||
|
||||
The order in which the scheduler returns its stored requests (via the ``next_request`` method)
|
||||
plays a great part in determining the order in which those requests are downloaded.
|
||||
|
||||
The methods defined in this class constitute the minimal interface that the Scrapy engine will interact with.
|
||||
"""
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler: Crawler):
|
||||
"""
|
||||
Factory method which receives the current :class:`~scrapy.crawler.Crawler` object as argument.
|
||||
"""
|
||||
return cls()
|
||||
|
||||
def open(self, spider: Spider) -> Optional[Deferred]:
|
||||
"""
|
||||
Called when the spider is opened by the engine. It receives the spider
|
||||
instance as argument and it's useful to execute initialization code.
|
||||
|
||||
:param spider: the spider object for the current crawl
|
||||
:type spider: :class:`~scrapy.spiders.Spider`
|
||||
"""
|
||||
pass
|
||||
|
||||
def close(self, reason: str) -> Optional[Deferred]:
|
||||
"""
|
||||
Called when the spider is closed by the engine. It receives the reason why the crawl
|
||||
finished as argument and it's useful to execute cleaning code.
|
||||
|
||||
:param reason: a string which describes the reason why the spider was closed
|
||||
:type reason: :class:`str`
|
||||
"""
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
def has_pending_requests(self) -> bool:
|
||||
"""
|
||||
``True`` if the scheduler has enqueued requests, ``False`` otherwise
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
@abstractmethod
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
"""
|
||||
Process a request received by the engine.
|
||||
|
||||
Return ``True`` if the request is stored correctly, ``False`` otherwise.
|
||||
|
||||
If ``False``, the engine will fire a ``request_dropped`` signal, and
|
||||
will not make further attempts to schedule the request at a later time.
|
||||
For reference, the default Scrapy scheduler returns ``False`` when the
|
||||
request is rejected by the dupefilter.
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
@abstractmethod
|
||||
def next_request(self) -> Optional[Request]:
|
||||
"""
|
||||
Return the next :class:`~scrapy.http.Request` to be processed, or ``None``
|
||||
to indicate that there are no requests to be considered ready at the moment.
|
||||
|
||||
Returning ``None`` implies that no request from the scheduler will be sent
|
||||
to the downloader in the current reactor cycle. The engine will continue
|
||||
calling ``next_request`` until ``has_pending_requests`` is ``False``.
|
||||
"""
|
||||
raise NotImplementedError()
|
||||
|
||||
|
||||
SchedulerTV = TypeVar("SchedulerTV", bound="Scheduler")
|
||||
|
||||
|
||||
class Scheduler(BaseScheduler):
|
||||
"""
|
||||
Default Scrapy scheduler. This implementation also handles duplication
|
||||
filtering via the :setting:`dupefilter <DUPEFILTER_CLASS>`.
|
||||
|
||||
This scheduler stores requests into several priority queues (defined by the
|
||||
:setting:`SCHEDULER_PRIORITY_QUEUE` setting). In turn, said priority queues
|
||||
are backed by either memory or disk based queues (respectively defined by the
|
||||
:setting:`SCHEDULER_MEMORY_QUEUE` and :setting:`SCHEDULER_DISK_QUEUE` settings).
|
||||
|
||||
Request prioritization is almost entirely delegated to the priority queue. The only
|
||||
prioritization performed by this scheduler is using the disk-based queue if present
|
||||
(i.e. if the :setting:`JOBDIR` setting is defined) and falling back to the memory-based
|
||||
queue if a serialization error occurs. If the disk queue is not present, the memory one
|
||||
is used directly.
|
||||
|
||||
:param dupefilter: An object responsible for checking and filtering duplicate requests.
|
||||
The value for the :setting:`DUPEFILTER_CLASS` setting is used by default.
|
||||
:type dupefilter: :class:`scrapy.dupefilters.BaseDupeFilter` instance or similar:
|
||||
any class that implements the `BaseDupeFilter` interface
|
||||
|
||||
:param jobdir: The path of a directory to be used for persisting the crawl's state.
|
||||
The value for the :setting:`JOBDIR` setting is used by default.
|
||||
See :ref:`topics-jobs`.
|
||||
:type jobdir: :class:`str` or ``None``
|
||||
|
||||
:param dqclass: A class to be used as persistent request queue.
|
||||
The value for the :setting:`SCHEDULER_DISK_QUEUE` setting is used by default.
|
||||
:type dqclass: class
|
||||
|
||||
:param mqclass: A class to be used as non-persistent request queue.
|
||||
The value for the :setting:`SCHEDULER_MEMORY_QUEUE` setting is used by default.
|
||||
:type mqclass: class
|
||||
|
||||
:param logunser: A boolean that indicates whether or not unserializable requests should be logged.
|
||||
The value for the :setting:`SCHEDULER_DEBUG` setting is used by default.
|
||||
:type logunser: bool
|
||||
|
||||
:param stats: A stats collector object to record stats about the request scheduling process.
|
||||
The value for the :setting:`STATS_CLASS` setting is used by default.
|
||||
:type stats: :class:`scrapy.statscollectors.StatsCollector` instance or similar:
|
||||
any class that implements the `StatsCollector` interface
|
||||
|
||||
:param pqclass: A class to be used as priority queue for requests.
|
||||
The value for the :setting:`SCHEDULER_PRIORITY_QUEUE` setting is used by default.
|
||||
:type pqclass: class
|
||||
|
||||
:param crawler: The crawler object corresponding to the current crawl.
|
||||
:type crawler: :class:`scrapy.crawler.Crawler`
|
||||
"""
|
||||
def __init__(
|
||||
self,
|
||||
dupefilter,
|
||||
jobdir: Optional[str] = None,
|
||||
dqclass=None,
|
||||
mqclass=None,
|
||||
logunser: bool = False,
|
||||
stats=None,
|
||||
pqclass=None,
|
||||
crawler: Optional[Crawler] = None,
|
||||
):
|
||||
self.df = dupefilter
|
||||
self.dqdir = self._dqdir(jobdir)
|
||||
self.pqclass = pqclass
|
||||
|
|
@ -47,34 +184,57 @@ class Scheduler:
|
|||
self.crawler = crawler
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler):
|
||||
settings = crawler.settings
|
||||
dupefilter_cls = load_object(settings['DUPEFILTER_CLASS'])
|
||||
dupefilter = create_instance(dupefilter_cls, settings, crawler)
|
||||
pqclass = load_object(settings['SCHEDULER_PRIORITY_QUEUE'])
|
||||
dqclass = load_object(settings['SCHEDULER_DISK_QUEUE'])
|
||||
mqclass = load_object(settings['SCHEDULER_MEMORY_QUEUE'])
|
||||
logunser = settings.getbool('SCHEDULER_DEBUG')
|
||||
return cls(dupefilter, jobdir=job_dir(settings), logunser=logunser,
|
||||
stats=crawler.stats, pqclass=pqclass, dqclass=dqclass,
|
||||
mqclass=mqclass, crawler=crawler)
|
||||
def from_crawler(cls: Type[SchedulerTV], crawler) -> SchedulerTV:
|
||||
"""
|
||||
Factory method, initializes the scheduler with arguments taken from the crawl settings
|
||||
"""
|
||||
dupefilter_cls = load_object(crawler.settings['DUPEFILTER_CLASS'])
|
||||
return cls(
|
||||
dupefilter=create_instance(dupefilter_cls, crawler.settings, crawler),
|
||||
jobdir=job_dir(crawler.settings),
|
||||
dqclass=load_object(crawler.settings['SCHEDULER_DISK_QUEUE']),
|
||||
mqclass=load_object(crawler.settings['SCHEDULER_MEMORY_QUEUE']),
|
||||
logunser=crawler.settings.getbool('SCHEDULER_DEBUG'),
|
||||
stats=crawler.stats,
|
||||
pqclass=load_object(crawler.settings['SCHEDULER_PRIORITY_QUEUE']),
|
||||
crawler=crawler,
|
||||
)
|
||||
|
||||
def has_pending_requests(self):
|
||||
def has_pending_requests(self) -> bool:
|
||||
return len(self) > 0
|
||||
|
||||
def open(self, spider):
|
||||
def open(self, spider: Spider) -> Optional[Deferred]:
|
||||
"""
|
||||
(1) initialize the memory queue
|
||||
(2) initialize the disk queue if the ``jobdir`` attribute is a valid directory
|
||||
(3) return the result of the dupefilter's ``open`` method
|
||||
"""
|
||||
self.spider = spider
|
||||
self.mqs = self._mq()
|
||||
self.dqs = self._dq() if self.dqdir else None
|
||||
return self.df.open()
|
||||
|
||||
def close(self, reason):
|
||||
if self.dqs:
|
||||
def close(self, reason: str) -> Optional[Deferred]:
|
||||
"""
|
||||
(1) dump pending requests to disk if there is a disk queue
|
||||
(2) return the result of the dupefilter's ``close`` method
|
||||
"""
|
||||
if self.dqs is not None:
|
||||
state = self.dqs.close()
|
||||
assert isinstance(self.dqdir, str)
|
||||
self._write_dqs_state(self.dqdir, state)
|
||||
return self.df.close(reason)
|
||||
|
||||
def enqueue_request(self, request):
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
"""
|
||||
Unless the received request is filtered out by the Dupefilter, attempt to push
|
||||
it into the disk queue, falling back to pushing it into the memory queue.
|
||||
|
||||
Increment the appropriate stats, such as: ``scheduler/enqueued``,
|
||||
``scheduler/enqueued/disk``, ``scheduler/enqueued/memory``.
|
||||
|
||||
Return ``True`` if the request was stored successfully, ``False`` otherwise.
|
||||
"""
|
||||
if not request.dont_filter and self.df.request_seen(request):
|
||||
self.df.log(request, self.spider)
|
||||
return False
|
||||
|
|
@ -87,24 +247,35 @@ class Scheduler:
|
|||
self.stats.inc_value('scheduler/enqueued', spider=self.spider)
|
||||
return True
|
||||
|
||||
def next_request(self):
|
||||
def next_request(self) -> Optional[Request]:
|
||||
"""
|
||||
Return a :class:`~scrapy.http.Request` object from the memory queue,
|
||||
falling back to the disk queue if the memory queue is empty.
|
||||
Return ``None`` if there are no more enqueued requests.
|
||||
|
||||
Increment the appropriate stats, such as: ``scheduler/dequeued``,
|
||||
``scheduler/dequeued/disk``, ``scheduler/dequeued/memory``.
|
||||
"""
|
||||
request = self.mqs.pop()
|
||||
if request:
|
||||
if request is not None:
|
||||
self.stats.inc_value('scheduler/dequeued/memory', spider=self.spider)
|
||||
else:
|
||||
request = self._dqpop()
|
||||
if request:
|
||||
if request is not None:
|
||||
self.stats.inc_value('scheduler/dequeued/disk', spider=self.spider)
|
||||
if request:
|
||||
if request is not None:
|
||||
self.stats.inc_value('scheduler/dequeued', spider=self.spider)
|
||||
return request
|
||||
|
||||
def __len__(self):
|
||||
return len(self.dqs) + len(self.mqs) if self.dqs else len(self.mqs)
|
||||
def __len__(self) -> int:
|
||||
"""
|
||||
Return the total amount of enqueued requests
|
||||
"""
|
||||
return len(self.dqs) + len(self.mqs) if self.dqs is not None else len(self.mqs)
|
||||
|
||||
def _dqpush(self, request):
|
||||
def _dqpush(self, request: Request) -> bool:
|
||||
if self.dqs is None:
|
||||
return
|
||||
return False
|
||||
try:
|
||||
self.dqs.push(request)
|
||||
except ValueError as e: # non serializable request
|
||||
|
|
@ -115,18 +286,18 @@ class Scheduler:
|
|||
logger.warning(msg, {'request': request, 'reason': e},
|
||||
exc_info=True, extra={'spider': self.spider})
|
||||
self.logunser = False
|
||||
self.stats.inc_value('scheduler/unserializable',
|
||||
spider=self.spider)
|
||||
return
|
||||
self.stats.inc_value('scheduler/unserializable', spider=self.spider)
|
||||
return False
|
||||
else:
|
||||
return True
|
||||
|
||||
def _mqpush(self, request):
|
||||
def _mqpush(self, request: Request) -> None:
|
||||
self.mqs.push(request)
|
||||
|
||||
def _dqpop(self):
|
||||
if self.dqs:
|
||||
def _dqpop(self) -> Optional[Request]:
|
||||
if self.dqs is not None:
|
||||
return self.dqs.pop()
|
||||
return None
|
||||
|
||||
def _mq(self):
|
||||
""" Create a new priority queue instance, with in-memory storage """
|
||||
|
|
@ -150,21 +321,22 @@ class Scheduler:
|
|||
{'queuesize': len(q)}, extra={'spider': self.spider})
|
||||
return q
|
||||
|
||||
def _dqdir(self, jobdir):
|
||||
def _dqdir(self, jobdir: Optional[str]) -> Optional[str]:
|
||||
""" Return a folder name to keep disk queue state at """
|
||||
if jobdir:
|
||||
if jobdir is not None:
|
||||
dqdir = join(jobdir, 'requests.queue')
|
||||
if not exists(dqdir):
|
||||
os.makedirs(dqdir)
|
||||
return dqdir
|
||||
return None
|
||||
|
||||
def _read_dqs_state(self, dqdir):
|
||||
def _read_dqs_state(self, dqdir: str) -> list:
|
||||
path = join(dqdir, 'active.json')
|
||||
if not exists(path):
|
||||
return ()
|
||||
return []
|
||||
with open(path) as f:
|
||||
return json.load(f)
|
||||
|
||||
def _write_dqs_state(self, dqdir, state):
|
||||
def _write_dqs_state(self, dqdir: str, state: list) -> None:
|
||||
with open(join(dqdir, 'active.json'), 'w') as f:
|
||||
json.dump(state, f)
|
||||
|
|
|
|||
|
|
@ -1,7 +1,10 @@
|
|||
import os
|
||||
from typing import Optional
|
||||
|
||||
from scrapy.settings import BaseSettings
|
||||
|
||||
|
||||
def job_dir(settings):
|
||||
def job_dir(settings: BaseSettings) -> Optional[str]:
|
||||
path = settings['JOBDIR']
|
||||
if path and not os.path.exists(path):
|
||||
os.makedirs(path)
|
||||
|
|
|
|||
|
|
@ -0,0 +1,159 @@
|
|||
from typing import Dict, Optional
|
||||
from unittest import TestCase
|
||||
from urllib.parse import urljoin, urlparse
|
||||
|
||||
from testfixtures import LogCapture
|
||||
from twisted.internet import defer
|
||||
from twisted.trial.unittest import TestCase as TwistedTestCase
|
||||
|
||||
from scrapy.core.scheduler import BaseScheduler
|
||||
from scrapy.crawler import CrawlerRunner
|
||||
from scrapy.http import Request
|
||||
from scrapy.spiders import Spider
|
||||
from scrapy.utils.request import request_fingerprint
|
||||
|
||||
from tests.mockserver import MockServer
|
||||
|
||||
|
||||
PATHS = ["/a", "/b", "/c"]
|
||||
URLS = [urljoin("https://example.org", p) for p in PATHS]
|
||||
|
||||
|
||||
class MinimalScheduler:
|
||||
def __init__(self) -> None:
|
||||
self.requests: Dict[str, Request] = {}
|
||||
|
||||
def has_pending_requests(self) -> bool:
|
||||
return bool(self.requests)
|
||||
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
fp = request_fingerprint(request)
|
||||
if fp not in self.requests:
|
||||
self.requests[fp] = request
|
||||
return True
|
||||
return False
|
||||
|
||||
def next_request(self) -> Optional[Request]:
|
||||
if self.has_pending_requests():
|
||||
fp, request = self.requests.popitem()
|
||||
return request
|
||||
return None
|
||||
|
||||
|
||||
class SimpleScheduler(MinimalScheduler):
|
||||
def open(self, spider: Spider) -> defer.Deferred:
|
||||
return defer.succeed("open")
|
||||
|
||||
def close(self, reason: str) -> defer.Deferred:
|
||||
return defer.succeed("close")
|
||||
|
||||
def __len__(self) -> int:
|
||||
return len(self.requests)
|
||||
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
|
||||
def __init__(self, mockserver, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.start_urls = map(mockserver.url, PATHS)
|
||||
|
||||
def parse(self, response):
|
||||
return {"path": urlparse(response.url).path}
|
||||
|
||||
|
||||
class InterfaceCheckMixin:
|
||||
def test_scheduler_class(self):
|
||||
self.assertTrue(isinstance(self.scheduler, BaseScheduler))
|
||||
self.assertTrue(issubclass(self.scheduler.__class__, BaseScheduler))
|
||||
|
||||
|
||||
class BaseSchedulerTest(TestCase, InterfaceCheckMixin):
|
||||
def setUp(self):
|
||||
self.scheduler = BaseScheduler()
|
||||
|
||||
def test_methods(self):
|
||||
self.assertIsNone(self.scheduler.open(Spider("foo")))
|
||||
self.assertIsNone(self.scheduler.close("finished"))
|
||||
self.assertRaises(NotImplementedError, self.scheduler.has_pending_requests)
|
||||
self.assertRaises(NotImplementedError, self.scheduler.enqueue_request, Request("https://example.org"))
|
||||
self.assertRaises(NotImplementedError, self.scheduler.next_request)
|
||||
|
||||
|
||||
class MinimalSchedulerTest(TestCase, InterfaceCheckMixin):
|
||||
def setUp(self):
|
||||
self.scheduler = MinimalScheduler()
|
||||
|
||||
def test_open_close(self):
|
||||
with self.assertRaises(AttributeError):
|
||||
self.scheduler.open(Spider("foo"))
|
||||
with self.assertRaises(AttributeError):
|
||||
self.scheduler.close("finished")
|
||||
|
||||
def test_len(self):
|
||||
with self.assertRaises(AttributeError):
|
||||
self.scheduler.__len__()
|
||||
with self.assertRaises(TypeError):
|
||||
len(self.scheduler)
|
||||
|
||||
def test_enqueue_dequeue(self):
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
for url in URLS:
|
||||
self.assertTrue(self.scheduler.enqueue_request(Request(url)))
|
||||
self.assertFalse(self.scheduler.enqueue_request(Request(url)))
|
||||
self.assertTrue(self.scheduler.has_pending_requests)
|
||||
|
||||
dequeued = []
|
||||
while self.scheduler.has_pending_requests():
|
||||
request = self.scheduler.next_request()
|
||||
dequeued.append(request.url)
|
||||
self.assertEqual(set(dequeued), set(URLS))
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
|
||||
|
||||
class SimpleSchedulerTest(TwistedTestCase, InterfaceCheckMixin):
|
||||
def setUp(self):
|
||||
self.scheduler = SimpleScheduler()
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_enqueue_dequeue(self):
|
||||
open_result = yield self.scheduler.open(Spider("foo"))
|
||||
self.assertEqual(open_result, "open")
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
|
||||
for url in URLS:
|
||||
self.assertTrue(self.scheduler.enqueue_request(Request(url)))
|
||||
self.assertFalse(self.scheduler.enqueue_request(Request(url)))
|
||||
|
||||
self.assertTrue(self.scheduler.has_pending_requests())
|
||||
self.assertEqual(len(self.scheduler), len(URLS))
|
||||
|
||||
dequeued = []
|
||||
while self.scheduler.has_pending_requests():
|
||||
request = self.scheduler.next_request()
|
||||
dequeued.append(request.url)
|
||||
self.assertEqual(set(dequeued), set(URLS))
|
||||
|
||||
self.assertFalse(self.scheduler.has_pending_requests())
|
||||
self.assertEqual(len(self.scheduler), 0)
|
||||
|
||||
close_result = yield self.scheduler.close("")
|
||||
self.assertEqual(close_result, "close")
|
||||
|
||||
|
||||
class MinimalSchedulerCrawlTest(TwistedTestCase):
|
||||
scheduler_cls = MinimalScheduler
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_crawl(self):
|
||||
with MockServer() as mockserver:
|
||||
settings = {"SCHEDULER": self.scheduler_cls}
|
||||
with LogCapture() as log:
|
||||
yield CrawlerRunner(settings).crawl(TestSpider, mockserver)
|
||||
for path in PATHS:
|
||||
self.assertIn(f"{{'path': '{path}'}}", str(log))
|
||||
self.assertIn(f"'item_scraped_count': {len(PATHS)}", str(log))
|
||||
|
||||
|
||||
class SimpleSchedulerCrawlTest(MinimalSchedulerCrawlTest):
|
||||
scheduler_cls = SimpleScheduler
|
||||
Loading…
Reference in New Issue