Re-implement front-load support

This commit is contained in:
Adrián Chaves 2025-03-31 11:49:28 +02:00
parent 27ae12dc59
commit a04e85e08b
11 changed files with 112 additions and 51 deletions

View File

@ -37,7 +37,7 @@ Backward-incompatible changes
:ref:`spider middleware <topics-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.

View File

@ -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

View File

@ -54,6 +54,8 @@ Built-in scheduler
------------------
.. autoclass:: Scheduler()
:members: __len__, pause, unpause
:member-order: bysource
.. _priority-queues:

View File

@ -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 <topics-scheduler>` start requests. But you can :ref:`force a
specific order <start-requests-order>` or :ref:`delay scheduling
<start-requests-lazy>` as needed.
:ref:`scheduling <topics-scheduler>` 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 <scrapy.http.Request.priority>`. For
example:
method to set :attr:`Request.priority <scrapy.http.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 <topics-scheduler>` if you need
more control over request prioritization.
.. note:: By default, :ref:`the first few requests are sent in yield order
<start-requests-front-load>` 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
<start-requests-order>`. 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 <custom-scheduler>`, make sure it
supports pausing and unpausing, like
:class:`~scrapy.core.scheduler.Scheduler` does.
.. _start-requests-lazy:
Lazy start request scheduling

View File

@ -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

View File

@ -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 <topics-components>`.
You can access the running engine at :attr:`Crawler.engine
<scrapy.crawler.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 <topics-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(

View File

@ -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

View File

@ -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:

View File

@ -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",

View File

@ -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)

View File

@ -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(