diff --git a/docs/news.rst b/docs/news.rst index 9e826f71a..634b1d934 100644 --- a/docs/news.rst +++ b/docs/news.rst @@ -37,7 +37,7 @@ Backward-incompatible changes :ref:`spider middleware ` no longer stops the crawl. -- In ``scrapy.core.engine.ExecutionEngine``: +- In :class:`~scrapy.core.engine.ExecutionEngine`: - The second parameter of ``open_spider()``, ``start_requests``, has been removed. The start requests are determined by the ``spider`` parameter @@ -46,6 +46,10 @@ Backward-incompatible changes - The ``slot`` attribute has been renamed to ``_slot`` and should not be used. + To access the running scheduler, previously at ``_slot.scheduler``, use + the :attr:`~scrapy.core.engine.ExecutionEngine.scheduler` attribute of + the running engine instead. + - In ``scrapy.core.engine``, the ``Slot`` class has been renamed to ``_Slot`` and should not be used. diff --git a/docs/topics/api.rst b/docs/topics/api.rst index 8e8f3a0c9..54b4b4830 100644 --- a/docs/topics/api.rst +++ b/docs/topics/api.rst @@ -84,14 +84,7 @@ how you :ref:`configure the downloader middlewares For an introduction on extensions and a list of available extensions on Scrapy see :ref:`topics-extensions`. - .. attribute:: engine - - The execution engine, which coordinates the core crawling logic - between the scheduler, downloader and spiders. - - Some extension may want to access the Scrapy engine, to inspect or - modify the downloader and scheduler behaviour, although this is an - advanced use and this API is not yet stable. + .. autoattribute:: engine .. attribute:: spider @@ -285,4 +278,4 @@ Engine API ========== .. autoclass:: scrapy.core.engine.ExecutionEngine() - :members: needs_backout + :members: needs_backout, scheduler diff --git a/docs/topics/scheduler.rst b/docs/topics/scheduler.rst index 5d0c40df6..411ebe0f5 100644 --- a/docs/topics/scheduler.rst +++ b/docs/topics/scheduler.rst @@ -54,6 +54,8 @@ Built-in scheduler ------------------ .. autoclass:: Scheduler() + :members: __len__, pause, unpause + :member-order: bysource .. _priority-queues: diff --git a/docs/topics/spiders.rst b/docs/topics/spiders.rst index 867331a13..c24dd7dd9 100644 --- a/docs/topics/spiders.rst +++ b/docs/topics/spiders.rst @@ -374,9 +374,8 @@ method of a spider. By default, Scrapy does not try to send :meth:`~scrapy.Spider.start` requests in order. Instead, it prioritizes reaching :setting:`CONCURRENT_REQUESTS` and -:ref:`scheduling ` start requests. But you can :ref:`force a -specific order ` or :ref:`delay scheduling -` as needed. +:ref:`scheduling ` start requests. However, you can change +that. .. The request send order when all start requests and callback requests have @@ -411,8 +410,9 @@ Start request order ------------------- To force a specific request order, override the :meth:`~scrapy.Spider.start` -method to set :attr:`Request.priority `. For -example: +method to set :attr:`Request.priority `. + +For example: - To send start requests before other requests: @@ -424,7 +424,7 @@ example: item_or_request = item_or_request.replace(priority=1) yield item_or_request -- To send start requests in order: +- To send start requests in yield order: .. code-block:: python @@ -439,6 +439,36 @@ example: You can also :ref:`customize the scheduler ` if you need more control over request prioritization. +.. note:: By default, :ref:`the first few requests are sent in yield order + ` regardless, for performance reasons. + +.. _start-requests-front-load: + +Start request front loading +--------------------------- + +By default, the first few requests yielded by :meth:`~scrapy.Spider.start` are +sent in yield order, even if you :ref:`try to sort them differently +`. This is because, until :setting:`CONCURRENT_REQUESTS` +is reached, start requests are removed from the scheduler immediately after +they are scheduled. + +If you do not yield the highest-priority start requests first, and you want +request priorities to be respected since the first request, pause the scheduler +while yielding your start requests: + +.. code-block:: python + + async def start(self): + self.crawler.engine.scheduler.pause() + async for item_or_request in super().start(): + yield item_or_request + self.crawler.engine.scheduler.unpause() + +.. note:: If you use a :ref:`custom scheduler `, make sure it + supports pausing and unpausing, like + :class:`~scrapy.core.scheduler.Scheduler` does. + .. _start-requests-lazy: Lazy start request scheduling diff --git a/docs/topics/telnetconsole.rst b/docs/topics/telnetconsole.rst index 3e9bbe56e..b7526273e 100644 --- a/docs/topics/telnetconsole.rst +++ b/docs/topics/telnetconsole.rst @@ -116,8 +116,8 @@ using the telnet console:: engine.spider_is_idle() : False engine._slot.closing : False len(engine._slot.inprogress) : 16 - len(engine._slot.scheduler.dqs or []) : 0 - len(engine._slot.scheduler.mqs) : 92 + len(engine.scheduler.dqs or []) : 0 + len(engine.scheduler.mqs) : 92 len(engine.scraper.slot.queue) : 0 len(engine.scraper.slot.active) : 0 engine.scraper.slot.active_size : 0 diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 565f61919..ba7e4e049 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -16,6 +16,7 @@ from twisted.internet.defer import Deferred, succeed from twisted.python.failure import Failure from scrapy import signals +from scrapy.core.scheduler import BaseScheduler from scrapy.core.scraper import Scraper, _HandleOutputDeferred from scrapy.exceptions import CloseSpider, DontCloseSpider, IgnoreRequest from scrapy.http import Request, Response @@ -32,7 +33,6 @@ if TYPE_CHECKING: from collections.abc import AsyncIterator, Callable from scrapy.core.downloader import Downloader - from scrapy.core.scheduler import BaseScheduler from scrapy.crawler import Crawler from scrapy.logformatter import LogFormatter from scrapy.settings import BaseSettings, Settings @@ -50,13 +50,11 @@ class _Slot: self, close_if_idle: bool, nextcall: CallLaterOnce[None], - scheduler: BaseScheduler, ) -> None: self.closing: Deferred[None] | None = None self.inprogress: set[Request] = set() self.close_if_idle: bool = close_if_idle self.nextcall: CallLaterOnce[None] = nextcall - self.scheduler: BaseScheduler = scheduler def add_request(self, request: Request) -> None: self.inprogress.add(request) @@ -78,6 +76,12 @@ class _Slot: class ExecutionEngine: + """The engine handles some core :ref:`components `. + + You can access the running engine at :attr:`Crawler.engine + `. + """ + _MIN_BACK_IN_SECONDS = 0.001 _MAX_BACK_IN_SECONDS = 5.0 @@ -86,6 +90,9 @@ class ExecutionEngine: crawler: Crawler, spider_closed_callback: Callable[[Spider], Deferred[None] | None], ) -> None: + #: :ref:`Scheduler ` in use. + self.scheduler: BaseScheduler | None = None + self.crawler: Crawler = crawler self.settings: Settings = crawler.settings self.signals: SignalManager = crawler.signals @@ -113,8 +120,6 @@ class ExecutionEngine: raise def _get_scheduler_class(self, settings: BaseSettings) -> type[BaseScheduler]: - from scrapy.core.scheduler import BaseScheduler - scheduler_cls: type[BaseScheduler] = load_object(settings["SCHEDULER"]) if not issubclass(scheduler_cls, BaseScheduler): raise TypeError( @@ -214,12 +219,13 @@ class ExecutionEngine: def _scheduler_has_pending_requests(self) -> bool: assert self._slot is not None # typing + assert self.scheduler is not None # typing try: - return self._slot.scheduler.has_pending_requests() + return self.scheduler.has_pending_requests() except Exception as exception: exception_traceback = format_exc() logger.error( - f"{global_object_name(self._slot.scheduler.has_pending_requests)} raised an exception: {exception}.\n{exception_traceback}", + f"{global_object_name(self.scheduler.has_pending_requests)} raised an exception: {exception}.\n{exception_traceback}", exc_info=True, ) return False @@ -290,13 +296,14 @@ class ExecutionEngine: def _start_scheduled_request(self) -> Deferred[None] | None: assert self._slot is not None # typing assert self.spider is not None # typing + assert self.scheduler is not None # typing try: - request = self._slot.scheduler.next_request() + request = self.scheduler.next_request() except Exception as exception: exception_traceback = format_exc() logger.exception( - f"{global_object_name(self._slot.scheduler.next_request)} raised an exception: {exception}\n{exception_traceback}" + f"{global_object_name(self.scheduler.next_request)} raised an exception: {exception}\n{exception_traceback}" ) return None if request is None: @@ -380,7 +387,7 @@ class ExecutionEngine: self._slot.nextcall.schedule() # type: ignore[union-attr] def _schedule_request(self, request: Request, spider: Spider) -> None: - assert self._slot is not None # typing + assert self.scheduler is not None # typing request_scheduled_result = self.signals.send_catch_log( signals.request_scheduled, request=request, @@ -391,11 +398,11 @@ class ExecutionEngine: if isinstance(result, Failure) and isinstance(result.value, IgnoreRequest): return try: - request_was_enqueued = self._slot.scheduler.enqueue_request(request) + request_was_enqueued = self.scheduler.enqueue_request(request) except Exception as exception: exception_traceback = format_exc() logger.error( - f"{global_object_name(self._slot.scheduler.enqueue_request)} raised an exception: {exception}\n{exception_traceback}" + f"{global_object_name(self.scheduler.enqueue_request)} raised an exception: {exception}\n{exception_traceback}" ) request_was_enqueued = False if not request_was_enqueued: @@ -468,12 +475,12 @@ class ExecutionEngine: logger.info("Spider opened", extra={"spider": spider}) self.spider = spider nextcall = CallLaterOnce(self._start_scheduled_requests) - scheduler = build_from_crawler(self.scheduler_cls, self.crawler) - self._slot = _Slot(close_if_idle, nextcall, scheduler) + self.scheduler = build_from_crawler(self.scheduler_cls, self.crawler) + self._slot = _Slot(close_if_idle, nextcall) self._start = await maybe_deferred_to_future( self.scraper.spidermw.process_start(spider) ) - if hasattr(scheduler, "open") and (d := scheduler.open(spider)): + if hasattr(self.scheduler, "open") and (d := self.scheduler.open(spider)): await maybe_deferred_to_future(d) await maybe_deferred_to_future(self.scraper.open_spider(spider)) assert self.crawler.stats @@ -535,8 +542,8 @@ class ExecutionEngine: dfd.addBoth(lambda _: self.scraper.close_spider(spider)) dfd.addErrback(log_failure("Scraper close failure")) - if hasattr(self._slot.scheduler, "close"): - dfd.addBoth(lambda _: cast(_Slot, self._slot).scheduler.close(reason)) + if hasattr(self.scheduler, "close"): + dfd.addBoth(lambda _: cast(BaseScheduler, self.scheduler).close(reason)) dfd.addErrback(log_failure("Scheduler close failure")) dfd.addBoth( diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 098b7d6e4..2dec0d732 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -164,6 +164,7 @@ class Scheduler(BaseScheduler): self.logunser: bool = logunser self.stats: StatsCollector | None = stats self.crawler: Crawler | None = crawler + self._paused = False @classmethod def from_crawler(cls, crawler: Crawler) -> Self: @@ -210,6 +211,8 @@ class Scheduler(BaseScheduler): return True def next_request(self) -> Request | None: + if self._paused: + return None request: Request | None = self.mqs.pop() assert self.stats is not None if request is not None: @@ -223,8 +226,22 @@ class Scheduler(BaseScheduler): return request def __len__(self) -> int: + """Return the number of pending requests.""" return len(self.dqs) + len(self.mqs) if self.dqs is not None else len(self.mqs) + def pause(self) -> None: + """Stop sending pending requests. + + It does not affect enqueing. + + See :ref:`start-requests-front-load` for an example. + """ + self._paused = True + + def unpause(self) -> None: + """Resume sending pending requests.""" + self._paused = False + def _dqpush(self, request: Request) -> bool: if self.dqs is None: return False diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 749096db5..0cc6cac3d 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -65,6 +65,9 @@ class Crawler: if isinstance(settings, dict) or settings is None: settings = Settings(settings) + #: Running engine. + self.engine: ExecutionEngine | None = None + self.spidercls: type[Spider] = spidercls self.settings: Settings = settings.copy() self.spidercls.update_settings(self.settings) @@ -82,7 +85,6 @@ class Crawler: self.logformatter: LogFormatter | None = None self.request_fingerprinter: RequestFingerprinterProtocol | None = None self.spider: Spider | None = None - self.engine: ExecutionEngine | None = None def _update_root_log_handler(self) -> None: if get_scrapy_root_handler() is not None: diff --git a/scrapy/utils/engine.py b/scrapy/utils/engine.py index 1e0c53212..a3dbdb04a 100644 --- a/scrapy/utils/engine.py +++ b/scrapy/utils/engine.py @@ -20,8 +20,8 @@ def get_engine_status(engine: ExecutionEngine) -> list[tuple[str, Any]]: "engine.spider_is_idle()", "engine._slot.closing", "len(engine._slot.inprogress)", - "len(engine._slot.scheduler.dqs or [])", - "len(engine._slot.scheduler.mqs)", + "len(engine.scheduler.dqs or [])", + "len(engine.scheduler.mqs)", "len(engine.scraper.slot.queue)", "len(engine.scraper.slot.active)", "engine.scraper.slot.active_size", diff --git a/tests/test_engine.py b/tests/test_engine.py index de6fb1175..c8e0d8b25 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -493,7 +493,8 @@ def test_request_scheduled_signal(caplog): yield engine._start = start() - engine._slot = _Slot(False, Mock(), scheduler) + engine.scheduler = scheduler + engine._slot = _Slot(False, Mock()) crawler.signals.connect(signal_handler, request_scheduled) keep_request = Request("https://keep.example") engine._schedule_request(keep_request, spider) diff --git a/tests/test_engine_loop.py b/tests/test_engine_loop.py index 507db93db..ef947f441 100644 --- a/tests/test_engine_loop.py +++ b/tests/test_engine_loop.py @@ -14,7 +14,7 @@ from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future from scrapy.utils.test import get_crawler from .mockserver import MockServer -from .test_scheduler import MemoryScheduler, PriorityScheduler +from .test_scheduler import MemoryScheduler async def sleep(seconds: float = ExecutionEngine._MIN_BACK_IN_SECONDS) -> None: @@ -36,7 +36,7 @@ class MainTestCase(TestCase): async def start(self): yield Request("data:,a") - self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,b")) + self.crawler.engine.scheduler.enqueue_request(Request("data:,b")) raise RuntimeError def parse(self, response): @@ -305,16 +305,21 @@ class RequestSendOrderTestCase(TestCase): ) ) - @pytest.mark.skip(reason="not implemented yet") @deferred_f_from_coro_f async def test_front_load(self): class TestSpider(Spider): name = "test" async def start(self): + assert self.crawler.engine is not None # typing + assert isinstance( + self.crawler.engine.scheduler, MemoryScheduler + ) # typing self.crawler.engine.scheduler.pause() - yield Request("data:,b", priority=0) - yield Request("data:,a", priority=1) + # By pausing the scheduler, a is scheduled before b is sent, + # and since the scheduler uses a LIFO queue, a is sent first. + yield Request("data:,b") + yield Request("data:,a") self.crawler.engine.scheduler.unpause() def parse(self, response): @@ -325,7 +330,7 @@ class RequestSendOrderTestCase(TestCase): def track_url(request, spider): actual_urls.append(request.url) - settings = {"SCHEDULER": PriorityScheduler} + settings = {"SCHEDULER": MemoryScheduler} crawler = get_crawler(TestSpider, settings_dict=settings) crawler.signals.connect(track_url, signals.request_reached_downloader) await maybe_deferred_to_future(crawler.crawl()) @@ -418,28 +423,28 @@ class RequestSendOrderTestCase(TestCase): # Let request 1 be processed. await spider.crawler.signals.wait_for(signals.scheduler_empty) - spider.crawler.engine._slot.scheduler.pause() - spider.crawler.engine._slot.scheduler.enqueue_request(_request(2)) + spider.crawler.engine.scheduler.pause() + spider.crawler.engine.scheduler.enqueue_request(_request(2)) # During this time, the scheduler reports having requests but # returns None. await spider.crawler.signals.wait_for(signals.scheduler_empty) - spider.crawler.engine._slot.scheduler.unpause() + spider.crawler.engine.scheduler.unpause() # The scheduler request is processed. await spider.crawler.signals.wait_for(signals.scheduler_empty) yield _request(3) - spider.crawler.engine._slot.scheduler.pause() - spider.crawler.engine._slot.scheduler.enqueue_request(_request(4)) + spider.crawler.engine.scheduler.pause() + spider.crawler.engine.scheduler.enqueue_request(_request(4)) # The last start request is processed during the time until the # delayed call below, proving that the start iteration can # finish before a scheduler “sleep” without causing the # scheduler to finish. - reactor.callLater(0, spider.crawler.engine._slot.scheduler.unpause) + reactor.callLater(0, spider.crawler.engine.scheduler.unpause) await maybe_deferred_to_future( self._test_request_order(