diff --git a/docs/topics/practices.rst b/docs/topics/practices.rst index c9f1c6b3f..f88fa02ac 100644 --- a/docs/topics/practices.rst +++ b/docs/topics/practices.rst @@ -513,7 +513,7 @@ run: Because :setting:`SPIDER_MODULES` is a list setting, you can include multiple modules by separating them with commas. -.. _crawl-optimization: +.. _crawl-bottlenecks: Identifying crawl bottlenecks ============================= diff --git a/docs/topics/request-response.rst b/docs/topics/request-response.rst index f673506e4..75158440b 100644 --- a/docs/topics/request-response.rst +++ b/docs/topics/request-response.rst @@ -901,7 +901,6 @@ Those are: * :reqmeta:`redirect_reasons` * :reqmeta:`redirect_urls` * :reqmeta:`referrer_policy` -* :reqmeta:`response_rough_size` * :reqmeta:`verbatim_url` .. reqmeta:: bindaddress @@ -1016,20 +1015,6 @@ The meta key is used set retry times per request. When set, the :reqmeta:`max_retry_times` meta key takes higher precedence over the :setting:`RETRY_TIMES` setting. -.. reqmeta:: response_rough_size - -response_rough_size -------------------- - -Overrides the :setting:`RESPONSE_ROUGH_SIZE` setting for this specific request. - -The value is the estimated size (in bytes) to count toward -:setting:`RESPONSE_MAX_ACTIVE_SIZE` while this request is being downloaded, -before its actual response size is known. Set to ``0`` to disable rough-size -counting for this request. - -.. versionadded:: VERSION - .. reqmeta:: verbatim_url verbatim_url diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index e73035244..84efb2c91 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -1828,34 +1828,23 @@ This counts both the size of response bodies that have passed through memory, and the :setting:`rough size ` of requests currently being downloaded. -When the total exceeds this value, Scrapy pauses sending new requests to the -downloader until it drops below the limit. +While the total is above this value, Scrapy pauses sending new requests to the +downloader. A higher value improves crawl speed at the cost of memory usage. +``0`` disables the limit. -If you set this to ``0``, the limit is disabled. - -Setting this to a lower value reduces memory usage at the cost of crawl speed. -Setting this to a higher value (or disabling it) improves crawl speed but may -cause memory issues when responses are large. - -When the limit is first reached, Scrapy logs an info-level message explaining -the situation. Check the ``request_backout_seconds/response_max_active_size`` -stat to see how long request processing has been paused due to this limit over -the course of a crawl. +Scrapy logs an info-level message the first time the limit is reached. To see +how long request processing was paused because of it over a whole crawl, check +the ``request_backout_seconds/response_max_active_size`` stat. .. caution:: - If your code stores strong references to :class:`~scrapy.http.Response` - objects (e.g. in a scheduled request's meta or in a component attribute), - the garbage collector cannot free them, and the total active size may not - drop below the limit. In that case your crawl might get stuck indefinitely. - Either avoid storing such references, or set this to ``0`` to disable the - limit. - - To check whether your crawl is stuck due to this, connect to the - :ref:`telnet console ` and run ``prefs()`` to see - the count of live :class:`~scrapy.http.Response` objects. If that count - is large and not decreasing, you likely have strong response references. - See :ref:`topics-leaks` for details. + Responses that your code keeps a strong reference to, e.g. in the + :attr:`.Request.meta` of a scheduled request or in a component + attribute, count toward this limit until that reference is gone, so + accumulating them can pause a crawl indefinitely. To find out, run + ``prefs()`` on the :ref:`telnet console ` and see + whether the count of live :class:`~scrapy.http.Response` objects keeps + growing; see :ref:`topics-leaks`. .. versionadded:: VERSION @@ -1864,18 +1853,14 @@ the course of a crawl. RESPONSE_ROUGH_SIZE ------------------- -Default: ``131072`` (128 kiB) +Default: ``None`` Estimated size (in bytes) to count toward :setting:`RESPONSE_MAX_ACTIVE_SIZE` -for each request that is currently being downloaded, before its actual response -size is known. +for each request being downloaded, whose actual response size is not known yet. -This allows :setting:`RESPONSE_MAX_ACTIVE_SIZE` to provide backpressure based -on the number of concurrent in-flight requests, not just already-received -responses. Once the response arrives, its actual body size is counted instead. - -You can override this value on a per-request basis via the -:reqmeta:`response_rough_size` request meta key. +``None`` means a quarter of :setting:`RESPONSE_MAX_ACTIVE_SIZE` split among +:setting:`CONCURRENT_REQUESTS` requests, so that requests being downloaded +cannot use up the whole limit on their own. ``0`` counts only responses. .. versionadded:: VERSION diff --git a/docs/topics/telnetconsole.rst b/docs/topics/telnetconsole.rst index 54379b4b9..0ff3f2df7 100644 --- a/docs/topics/telnetconsole.rst +++ b/docs/topics/telnetconsole.rst @@ -125,9 +125,8 @@ engine status:: len(engine._slot.scheduler.mqs) : 92 len(engine.scraper.slot.queue) : 0 len(engine.scraper.slot.active) : 0 - engine.scraper.slot.active_size : 0 + engine.downloader.middleware._total_active_size : 1310720 engine.scraper.slot.itemproc_size : 0 - engine.scraper.slot.needs_backout() : False Pause, resume and stop the Scrapy engine diff --git a/scrapy/_classutilities.py b/scrapy/_classutilities.py deleted file mode 100644 index b788898fa..000000000 --- a/scrapy/_classutilities.py +++ /dev/null @@ -1,102 +0,0 @@ -# https://github.com/david-salac/classutilities/issues/1 -# Modified copy of -# https://github.com/david-salac/classutilities/blob/a6e4a86331936d432afaa454ed4c963528165a61/src/classutilities/classproperty.py - -# Allows creating a class level property -from __future__ import annotations - -from typing import TYPE_CHECKING, Any - -if TYPE_CHECKING: - from collections.abc import Callable - - -class ClassPropertyContainer: - """ - Allows creating a class level property (functionality for - decorator). - """ - - def __init__(self, prop_get: Any, prop_set: Any = None): - """ - Container that allows having a class property decorator. - :param prop_get: Class property getter. - :param prop_set: Class property setter. - """ - self.prop_get: Any = prop_get - self.prop_set: Any = prop_set - - def __get__(self, obj: Any, cls: type | None = None) -> Any: - """ - Return the value of the class property. - :param obj: Instance of the class. - :param cls: Type of the class. - :return: Value of the class property. - """ - if cls is None: # pragma: no cover - cls = type(obj) - return self.prop_get.__get__(obj, cls)() - - def __set__(self, obj: Any, value: Any) -> None: - """ - Set the value of the class property. - :param obj: Instance of the class. - :param value: A value to be set. - """ - if not self.prop_set: # pragma: no cover - raise AttributeError("cannot set attribute") - _type: type = type(obj) - if _type == ClassPropertyMetaClass: - _type = obj - self.prop_set.__get__(obj, _type)(value) - - def setter( - self, - func: Callable[..., Any] | classmethod[Any, Any, Any] | staticmethod[Any, Any], - ) -> ClassPropertyContainer: - """ - Allows creating setter in a property like way. - :param func: Setter function. - :return: Setter object for the decorator. - """ - if not isinstance(func, (classmethod, staticmethod)): - func = classmethod(func) - self.prop_set = func - return self - - -def classproperty( - func: Callable[..., Any] | classmethod[Any, Any, Any] | staticmethod[Any, Any], -) -> ClassPropertyContainer: - """ - Create a decorator for a class level property. - :param func: This class method is decorated. - :return: Modified class method behaving like a class property. - """ - if not isinstance(func, (classmethod, staticmethod)): - # The method must be a classmethod (or staticmethod) - func = classmethod(func) - return ClassPropertyContainer(func) - - -class ClassPropertyMetaClass(type): - """ - Metaclass that allows creating a standard setter. - """ - - def __setattr__(cls, key: str, value: Any) -> None: - """Overloads setter for class""" - obj = None - if key in cls.__dict__: - obj = cls.__dict__.get(key) - if obj and isinstance(obj, ClassPropertyContainer): - return obj.__set__(cls, value) - - return super().__setattr__(key, value) - - -class ClassPropertiesMixin(metaclass=ClassPropertyMetaClass): - """ - This mixin allows using class properties setter (getter works - correctly even without this mixin) - """ diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index e93cf0e95..e59afdc48 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -6,9 +6,8 @@ from collections import deque from dataclasses import dataclass, field from datetime import datetime from logging import getLogger -from time import monotonic, time +from time import monotonic from typing import TYPE_CHECKING, Any -from warnings import warn from twisted.internet.defer import Deferred, inlineCallbacks from twisted.python.failure import Failure @@ -16,7 +15,6 @@ from twisted.python.failure import Failure from scrapy import Request, Spider, signals from scrapy.core.downloader.handlers import DownloadHandlers from scrapy.core.downloader.middleware import DownloaderMiddlewareManager -from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.resolver import dnscache from scrapy.utils.asyncio import ( AsyncioLoopingCall, @@ -47,7 +45,7 @@ if TYPE_CHECKING: logger = getLogger(__name__) -@dataclass(slots=True, eq=False, repr=False) +@dataclass(slots=True, eq=False) class Slot: """Downloader slot""" @@ -55,20 +53,13 @@ class Slot: delay: float randomize_delay: bool - active: set[Request] = field(default_factory=set, init=False) + active: set[Request] = field(default_factory=set, init=False, repr=False) queue: deque[tuple[Request, Deferred[Response]]] = field( - default_factory=deque, init=False + default_factory=deque, init=False, repr=False ) - transferring: set[Request] = field(default_factory=set, init=False) - lastseen: float = field(default=0, init=False) - latercall: CallLaterResult | None = field(default=None, init=False) - - def __repr__(self) -> str: - return ( - f"Slot(concurrency={self.concurrency!r}, " - f"delay={self.delay:.2f}, " - f"randomize_delay={self.randomize_delay!r})" - ) + transferring: set[Request] = field(default_factory=set, init=False, repr=False) + lastseen: float = field(default=0, init=False, repr=False) + latercall: CallLaterResult | None = field(default=None, init=False, repr=False) def free_transfer_slots(self) -> int: return self.concurrency - len(self.transferring) @@ -112,6 +103,7 @@ def _get_concurrency_delay( class Downloader: DOWNLOAD_SLOT = "download_slot" _SLOT_GC_INTERVAL: float = 60.0 # seconds + _GC_INTERVAL: float = 1.0 # seconds def __init__(self, crawler: Crawler): self.crawler: Crawler = crawler @@ -137,45 +129,8 @@ class Downloader: # (reason, start_time): current backout reason and when it began, or # (None, None) if not backing out. self._last_backout: tuple[str | None, float | None] = (None, None) - - deprecated_setting_priority = self.settings.getpriority( - "SCRAPER_SLOT_MAX_ACTIVE_SIZE" - ) - assert deprecated_setting_priority is not None - setting_priority = self.settings.getpriority("RESPONSE_MAX_ACTIVE_SIZE") - assert setting_priority is not None - if deprecated_setting_priority > 0: - if setting_priority >= deprecated_setting_priority: - warn( - ( - "The SCRAPER_SLOT_MAX_ACTIVE_SIZE setting is deprecated " - "and is being ignored because RESPONSE_MAX_ACTIVE_SIZE is " - "set with an equal or higher priority. Remove " - "SCRAPER_SLOT_MAX_ACTIVE_SIZE from your settings." - ), - ScrapyDeprecationWarning, - stacklevel=2, - ) - self._response_max_active_size = self.settings.getint( - "RESPONSE_MAX_ACTIVE_SIZE" - ) - else: - warn( - ( - "The SCRAPER_SLOT_MAX_ACTIVE_SIZE setting is deprecated, " - "use RESPONSE_MAX_ACTIVE_SIZE instead." - ), - ScrapyDeprecationWarning, - stacklevel=2, - ) - self._response_max_active_size = self.settings.getint( - "SCRAPER_SLOT_MAX_ACTIVE_SIZE" - ) - else: - self._response_max_active_size = self.settings.getint( - "RESPONSE_MAX_ACTIVE_SIZE" - ) - self._response_max_active_size_warned = False + self._max_active_size_warned = False + self._last_gc: float = 0 @inlineCallbacks @_warn_spider_arg @@ -183,7 +138,7 @@ class Downloader: self, request: Request, spider: Spider | None = None ) -> Generator[Deferred[Any], Any, Response | Request]: self.active.add(request) - rough_size = self.middleware._count_rough_size(request) + self.middleware._count_rough_size(request) try: result: Response | Request = yield ( deferred_from_coro( @@ -193,13 +148,13 @@ class Downloader: return result finally: self.active.remove(request) - self.middleware._discount_rough_size(rough_size) + self.middleware._discount_rough_size(request) def _record_backout(self, reason: str | None) -> None: last_reason, last_reason_start_time = self._last_backout if last_reason == reason: return - current_time = time() + current_time = monotonic() if last_reason is not None and self._stats is not None: assert last_reason_start_time is not None last_reason_seconds = current_time - last_reason_start_time @@ -214,38 +169,27 @@ class Downloader: if 0 < self.total_concurrency <= len(self.active): self._record_backout("concurrency") return True - if ( - self._response_max_active_size - and self.middleware.total_active_size >= self._response_max_active_size - ): - if not self._response_max_active_size_warned: - self._response_max_active_size_warned = True + max_active_size = self.middleware._max_active_size + if max_active_size and self.middleware._total_active_size >= max_active_size: + if not self._max_active_size_warned: + self._max_active_size_warned = True logger.info( - f"The active response size, i.e. the total size of all " - f"bodies from responses that have been processed by " - f"downloader middlewares and remain in memory, plus the " - f"rough sizes of in-flight requests, is " - f"{self.middleware.total_active_size} B. The " - f"RESPONSE_MAX_ACTIVE_SIZE setting sets its maximum value " - f"at {self._response_max_active_size} B. No more requests " - f"will be processed until active response size lowers. If " - f"your memory allows it, you may increase " - f"RESPONSE_MAX_ACTIVE_SIZE, which should increase your " - f"crawl speed. If your code keeps non-weak references to " - f"Response objects, e.g. in (scheduled) requests or in a " - f"container within a component, your crawl might get " - f"stuck indefinitely; you can set " - f"RESPONSE_MAX_ACTIVE_SIZE to 0 to disable this limit, " - f"but then your code might run out of memory. This " - f"message will only appear the first time this happens. " - f"To learn how often request processing has been paused " - f"during a crawl for this reason, see the " - f"request_backout_seconds/response_max_active_size stat." + f"Pausing request processing: the active response size " + f"({self.middleware._total_active_size} B) has reached " + f"RESPONSE_MAX_ACTIVE_SIZE ({max_active_size} B). See " + f"https://docs.scrapy.org/en/latest/topics/settings.html#response-max-active-size " + f"and the request_backout_seconds/response_max_active_size " + f"stat. This message is only logged once." ) self._record_backout("response_max_active_size") - # Force the garbage collection of response objects. Necessary for - # PyPy, which is lazier when it comes to garbage collection. - gc.collect() + # Responses are only freed once nothing references them, which for + # reference cycles, and for every response on PyPy, requires a + # garbage collection. A full collection is expensive, and this runs + # once per request while paused, hence the interval. + current_time = monotonic() + if current_time - self._last_gc >= self._GC_INTERVAL: + self._last_gc = current_time + gc.collect() return True self._record_backout(None) return False diff --git a/scrapy/core/downloader/middleware.py b/scrapy/core/downloader/middleware.py index 0a126191d..92e4504e2 100644 --- a/scrapy/core/downloader/middleware.py +++ b/scrapy/core/downloader/middleware.py @@ -32,18 +32,60 @@ if TYPE_CHECKING: from scrapy.settings import BaseSettings +def _get_max_active_size(settings: BaseSettings) -> int: + deprecated_priority = settings.getpriority("SCRAPER_SLOT_MAX_ACTIVE_SIZE") + priority = settings.getpriority("RESPONSE_MAX_ACTIVE_SIZE") + assert deprecated_priority is not None + assert priority is not None + if deprecated_priority <= 0: + return settings.getint("RESPONSE_MAX_ACTIVE_SIZE") + if priority >= deprecated_priority: + warnings.warn( + "The SCRAPER_SLOT_MAX_ACTIVE_SIZE setting is deprecated and is " + "being ignored because RESPONSE_MAX_ACTIVE_SIZE is set with an " + "equal or higher priority. Remove SCRAPER_SLOT_MAX_ACTIVE_SIZE " + "from your settings.", + ScrapyDeprecationWarning, + stacklevel=2, + ) + return settings.getint("RESPONSE_MAX_ACTIVE_SIZE") + warnings.warn( + "The SCRAPER_SLOT_MAX_ACTIVE_SIZE setting is deprecated, use " + "RESPONSE_MAX_ACTIVE_SIZE instead.", + ScrapyDeprecationWarning, + stacklevel=2, + ) + return settings.getint("SCRAPER_SLOT_MAX_ACTIVE_SIZE") + + +def _get_response_rough_size(settings: BaseSettings, max_active_size: int) -> int: + if settings.get("RESPONSE_ROUGH_SIZE") is not None: + return settings.getint("RESPONSE_ROUGH_SIZE") + concurrency = settings.getint("CONCURRENT_REQUESTS") + if not concurrency: + # Unlimited concurrency: there is no bound on the number of in-flight + # requests to spread a share of the limit over. + return 0 + # Requests being downloaded, whose response size is unknown, may take up to + # a quarter of the limit; the rest is for responses already in memory. + return max_active_size // (4 * concurrency) + + class DownloaderMiddlewareManager(MiddlewareManager): component_name = "downloader middleware" def __init__(self, *args: Any, **kwargs: Any) -> None: super().__init__(*args, **kwargs) - self.response_active_size = 0 + assert self.crawler is not None + settings = self.crawler.settings + self._max_active_size: int = _get_max_active_size(settings) + self._response_rough_size: int = _get_response_rough_size( + settings, self._max_active_size + ) + self._response_active_size = 0 self._tracked_responses: WeakSet[Response] = WeakSet() self._rough_active_size = 0 - assert self.crawler is not None - self._response_rough_size: int = self.crawler.settings.getint( - "RESPONSE_ROUGH_SIZE" - ) + self._rough_sizes: dict[Request, int] = {} @classmethod def _get_mwlist_from_settings(cls, settings: BaseSettings) -> list[Any]: @@ -63,28 +105,34 @@ class DownloaderMiddlewareManager(MiddlewareManager): self._check_mw_method_spider_arg(mw.process_exception) @property - def total_active_size(self) -> int: - """Sum of sizes of tracked responses and rough sizes of in-flight requests.""" - return self.response_active_size + self._rough_active_size + def _total_active_size(self) -> int: + return self._response_active_size + self._rough_active_size - def _count_rough_size(self, request: Request) -> int: - size: int = request.meta.get("response_rough_size", self._response_rough_size) - self._rough_active_size += size - return size + def _count_rough_size(self, request: Request) -> None: + # Counting the same request twice, which concurrent downloads of the + # same request object would do, would leak its rough size, since only + # one discount call can find it. + if request in self._rough_sizes: + return + self._rough_sizes[request] = self._response_rough_size + self._rough_active_size += self._response_rough_size - def _discount_rough_size(self, size: int) -> None: - self._rough_active_size -= size + def _discount_rough_size(self, request: Request) -> None: + self._rough_active_size -= self._rough_sizes.pop(request, 0) - def _count_response_size(self, response: Response) -> None: + def _count_response_size(self, response: Response, request: Request) -> None: + # The actual response size replaces the rough size estimate of its + # request. + self._discount_rough_size(request) if response in self._tracked_responses: return self._tracked_responses.add(response) size = len(response.body) - self.response_active_size += size + self._response_active_size += size finalize(response, partial(self._discount_response_size, size)) def _discount_response_size(self, size: int) -> None: - self.response_active_size -= size + self._response_active_size -= size def download( self, @@ -146,11 +194,11 @@ class DownloaderMiddlewareManager(MiddlewareManager): ) if response: if isinstance(response, Response): - self._count_response_size(response) + self._count_response_size(response, request) return response result = await download_func(request) if isinstance(result, Response): - self._count_response_size(result) + self._count_response_size(result, request) return result async def _process_response( @@ -175,11 +223,11 @@ class DownloaderMiddlewareManager(MiddlewareManager): ) if isinstance(response, Request): return response - self._count_response_size(response) + self._count_response_size(response, request) return response async def _process_exception( - self, exception: Exception, request: Request | Response + self, exception: Exception, request: Request ) -> Response | Request: for method in self.methods["process_exception"]: assert method is not None @@ -194,6 +242,6 @@ class DownloaderMiddlewareManager(MiddlewareManager): ) if response: if isinstance(response, Response): - self._count_response_size(response) + self._count_response_size(response, request) return response raise exception diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index 978ef9956..e1ccfe1fb 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -13,7 +13,6 @@ from twisted.internet.defer import Deferred, inlineCallbacks from twisted.python.failure import Failure from scrapy import Spider, signals -from scrapy._classutilities import ClassPropertiesMixin, classproperty from scrapy.core.spidermw import SpiderMiddlewareManager from scrapy.exceptions import ( CloseSpider, @@ -59,33 +58,45 @@ QueueTuple: TypeAlias = tuple[Response | Failure, Request, Deferred[None]] _UNSET = object() -class Slot(ClassPropertiesMixin): +def _warn_min_response_size() -> None: + warnings.warn( + "scrapy.core.scraper.Slot.MIN_RESPONSE_SIZE is deprecated.", + ScrapyDeprecationWarning, + stacklevel=3, + ) + + +class _MinResponseSize: + """Deprecated alias of ``_MIN_RESPONSE_SIZE``, readable and writable both on + the class and on its instances.""" + + def __get__(self, instance: Slot | None, owner: type[Slot] | None = None) -> int: + _warn_min_response_size() + target = instance if instance is not None else owner + assert target is not None + return target._MIN_RESPONSE_SIZE + + def __set__(self, instance: Slot, value: int) -> None: + _warn_min_response_size() + instance._MIN_RESPONSE_SIZE = value + + +class _SlotMeta(type): + # Class-level assignment bypasses the _MinResponseSize descriptor. + def __setattr__(cls, name: str, value: Any) -> None: + if name == "MIN_RESPONSE_SIZE": + _warn_min_response_size() + name = "_MIN_RESPONSE_SIZE" + super().__setattr__(name, value) + + +class Slot(metaclass=_SlotMeta): """Scraper slot (one per running spider)""" _MIN_RESPONSE_SIZE = 1024 - - @classmethod - def _get_min_response_size(cls) -> int: - warnings.warn( - "scrapy.core.scraper.Slot.MIN_RESPONSE_SIZE is deprecated.", - ScrapyDeprecationWarning, - stacklevel=2, - ) - return cls._MIN_RESPONSE_SIZE - - @classmethod - def _set_min_response_size(cls, value: int) -> None: - warnings.warn( - "scrapy.core.scraper.Slot.MIN_RESPONSE_SIZE is deprecated.", - ScrapyDeprecationWarning, - stacklevel=2, - ) - cls._MIN_RESPONSE_SIZE = value - - # Typed as Any because mypy does not understand class properties. - MIN_RESPONSE_SIZE: Any = classproperty(_get_min_response_size).setter( - _set_min_response_size - ) + # Any so that mypy allows class-level assignment, which _SlotMeta redirects + # to _MIN_RESPONSE_SIZE. + MIN_RESPONSE_SIZE: Any = _MinResponseSize() def __init__(self, max_active_size: Any = _UNSET): if max_active_size is _UNSET: @@ -111,9 +122,10 @@ class Slot(ClassPropertiesMixin): def active_size(self) -> int: warnings.warn( ( - "scrapy.core.scraper.Slot.active_size is deprecated. Read " - "scrapy.core.downloader.DownloaderMiddlewareManager.response_active_size " - "instead." + "scrapy.core.scraper.Slot.active_size is deprecated. The size " + "of responses in memory is now tracked by the downloader, and " + "no longer has a public API. If you have a use case for one, " + "please open a GitHub issue." ), ScrapyDeprecationWarning, stacklevel=2, @@ -124,14 +136,8 @@ class Slot(ClassPropertiesMixin): def active_size(self, value: int) -> None: warnings.warn( ( - "scrapy.core.scraper.Slot.active_size is deprecated. " - "scrapy.core.downloader.DownloaderMiddlewareManager.response_active_size " - "might work as a replacement, but modifying that attribute " - "might not be a good idea. If you have a use case for it, you " - "might want to bring it up in a GitHub issue, to discuss with " - "Scrapy developers if there is a better approach, or some " - "change we could implement in Scrapy to improve support for " - "your use case." + "scrapy.core.scraper.Slot.active_size is deprecated, and " + "setting it no longer has any effect on request processing." ), ScrapyDeprecationWarning, stacklevel=2, diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 7f874d1c6..7930802c4 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -538,7 +538,8 @@ SCHEDULER_START_MEMORY_QUEUE = "scrapy.squeues.FifoMemoryQueue" SCRAPER_SLOT_MAX_ACTIVE_SIZE = 5_000_000 RESPONSE_MAX_ACTIVE_SIZE = 5_000_000 -RESPONSE_ROUGH_SIZE = 131072 +# None means a share of RESPONSE_MAX_ACTIVE_SIZE based on CONCURRENT_REQUESTS. +RESPONSE_ROUGH_SIZE = None SPIDER_CONTRACTS = {} SPIDER_CONTRACTS_BASE = { diff --git a/scrapy/utils/defer.py b/scrapy/utils/defer.py index d0ca9385d..7c6235f29 100644 --- a/scrapy/utils/defer.py +++ b/scrapy/utils/defer.py @@ -530,14 +530,9 @@ def _schedule_coro(coro: Coroutine[Any, Any, Any]) -> None: alternative is calling :func:`scrapy.utils.defer.deferred_from_coro`, keeping the result, and adding proper exception handling (e.g. errbacks) to it. - - In both asyncio and non-asyncio modes the coroutine starts in the next - event-loop iteration, never in the current call stack. """ if not is_asyncio_available(): - from twisted.internet import reactor - - reactor.callLater(0, Deferred.fromCoroutine, coro) + Deferred.fromCoroutine(coro) return loop = asyncio.get_event_loop() loop.create_task(coro) # noqa: RUF006 diff --git a/scrapy/utils/engine.py b/scrapy/utils/engine.py index 085720f66..e8135651d 100644 --- a/scrapy/utils/engine.py +++ b/scrapy/utils/engine.py @@ -24,9 +24,8 @@ def get_engine_status(engine: ExecutionEngine) -> list[tuple[str, Any]]: "len(engine._slot.scheduler.mqs)", "len(engine.scraper.slot.queue)", "len(engine.scraper.slot.active)", - "engine.scraper.slot.active_size", + "engine.downloader.middleware._total_active_size", "engine.scraper.slot.itemproc_size", - "engine.scraper.slot.needs_backout()", ] checks: list[tuple[str, Any]] = [] diff --git a/tests/test_downloader.py b/tests/test_downloader.py index 990d0cd56..f0e56773c 100644 --- a/tests/test_downloader.py +++ b/tests/test_downloader.py @@ -5,18 +5,31 @@ import pytest from twisted.internet.defer import Deferred from scrapy import Request, Spider -from scrapy.core.downloader import Slot +from scrapy.crawler import Crawler from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.http import Response +from scrapy.utils.asyncio import sleep from scrapy.utils.defer import maybe_deferred_to_future from scrapy.utils.test import get_crawler from tests.utils.decorators import coroutine_test -class TestSlot: - def test_repr(self): - slot = Slot(concurrency=8, delay=0.1, randomize_delay=True) - assert repr(slot) == "Slot(concurrency=8, delay=0.10, randomize_delay=True)" +def _count_backout_logs(caplog: pytest.LogCaptureFixture) -> int: + return sum( + 1 + for record in caplog.records + if record.levelname == "INFO" + and str(record.msg).startswith("Pausing request processing") + ) + + +def _backout_stats(crawler: Crawler) -> dict[str, Any]: + assert crawler.stats + return { + k: v + for k, v in crawler.stats.get_stats().items() + if k.startswith("request_backout_seconds/") + } class OfflineSpider(Spider): @@ -83,7 +96,7 @@ class TestResponseMaxActiveSize: warnings.simplefilter("error", ScrapyDeprecationWarning) await maybe_deferred_to_future(crawler.crawl()) assert crawler.engine - assert crawler.engine.downloader._response_max_active_size == 5_000_000 + assert crawler.engine.downloader.middleware._max_active_size == 5_000_000 @coroutine_test async def test_custom(self): @@ -96,7 +109,7 @@ class TestResponseMaxActiveSize: warnings.simplefilter("error", ScrapyDeprecationWarning) await maybe_deferred_to_future(crawler.crawl()) assert crawler.engine - assert crawler.engine.downloader._response_max_active_size == 0 + assert crawler.engine.downloader.middleware._max_active_size == 0 @coroutine_test async def test_deprecated_default(self): @@ -108,7 +121,7 @@ class TestResponseMaxActiveSize: with pytest.warns(ScrapyDeprecationWarning) as warning_messages: await maybe_deferred_to_future(crawler.crawl()) assert crawler.engine - assert crawler.engine.downloader._response_max_active_size == 5_000_000 + assert crawler.engine.downloader.middleware._max_active_size == 5_000_000 _assert_scraper_slot_deprecation(warning_messages) @coroutine_test @@ -122,7 +135,7 @@ class TestResponseMaxActiveSize: with pytest.warns(ScrapyDeprecationWarning) as warning_messages: await maybe_deferred_to_future(crawler.crawl()) assert crawler.engine - assert crawler.engine.downloader._response_max_active_size == 0 + assert crawler.engine.downloader.middleware._max_active_size == 0 _assert_scraper_slot_deprecation(warning_messages) @coroutine_test @@ -142,7 +155,7 @@ class TestResponseMaxActiveSize: with pytest.warns(ScrapyDeprecationWarning) as warning_messages: await maybe_deferred_to_future(crawler.crawl()) assert crawler.engine - assert crawler.engine.downloader._response_max_active_size == 1 + assert crawler.engine.downloader.middleware._max_active_size == 1 _assert_scraper_slot_deprecation(warning_messages, ignored=True) @coroutine_test @@ -169,7 +182,7 @@ class TestResponseMaxActiveSize: with pytest.warns(ScrapyDeprecationWarning) as warning_messages: await maybe_deferred_to_future(crawler.crawl()) assert crawler.engine - assert crawler.engine.downloader._response_max_active_size == 2 + assert crawler.engine.downloader.middleware._max_active_size == 2 _assert_scraper_slot_deprecation(warning_messages) @@ -178,70 +191,86 @@ class TestResponseRoughSize: def use_caplog(self, caplog): self.caplog = caplog + @pytest.mark.parametrize( + ("settings_dict", "expected"), + [ + # A quarter of RESPONSE_MAX_ACTIVE_SIZE split among + # CONCURRENT_REQUESTS requests. + ({}, 78125), + ({"CONCURRENT_REQUESTS": 100}, 12500), + ({"RESPONSE_MAX_ACTIVE_SIZE": 8_000_000}, 125_000), + # Unlimited concurrency leaves no number of requests to split it + # among. + ({"CONCURRENT_REQUESTS": 0}, 0), + ({"RESPONSE_MAX_ACTIVE_SIZE": 0}, 0), + ({"RESPONSE_ROUGH_SIZE": 1}, 1), + ({"RESPONSE_ROUGH_SIZE": 0}, 0), + ], + ) @coroutine_test - async def test_default(self): - crawler = get_crawler(OfflineSpider) + async def test_value(self, settings_dict, expected): + crawler = get_crawler(OfflineSpider, settings_dict=settings_dict) with warnings.catch_warnings(): warnings.simplefilter("error", ScrapyDeprecationWarning) await maybe_deferred_to_future(crawler.crawl()) assert crawler.engine - assert crawler.engine.downloader.middleware._response_rough_size == 131072 + assert crawler.engine.downloader.middleware._response_rough_size == expected @coroutine_test - async def test_custom(self): - """Setting RESPONSE_ROUGH_SIZE to a custom value changes the rough size.""" - crawler = get_crawler(OfflineSpider, settings_dict={"RESPONSE_ROUGH_SIZE": 0}) - with warnings.catch_warnings(): - warnings.simplefilter("error", ScrapyDeprecationWarning) - await maybe_deferred_to_future(crawler.crawl()) - assert crawler.engine - assert crawler.engine.downloader.middleware._response_rough_size == 0 + async def test_response_replaces_rough_size(self): + """The rough size of a request stops counting as soon as the size of its + response is known.""" + sizes: list[tuple[int, int]] = [] - @coroutine_test - async def test_rough_size_per_request(self): - """response_rough_size meta key overrides RESPONSE_ROUGH_SIZE per request. + class DownloaderMiddleware: + def __init__(self, crawler): + self.crawler = crawler - A low RESPONSE_MAX_ACTIVE_SIZE is set so that only requests with the - custom rough size of 1 pass; without the override the default 131072 - would trigger backout and only one request would be downloaded.""" + @classmethod + def from_crawler(cls, crawler): + return cls(crawler) + + def process_response(self, request, response): + middleware = self.crawler.engine.downloader.middleware + sizes.append( + (middleware._rough_active_size, middleware._response_active_size) + ) + return response class TestSpider(Spider): name = "test" + start_urls = ["data:,a"] custom_settings = { - "RESPONSE_MAX_ACTIVE_SIZE": 512, + "DOWNLOADER_MIDDLEWARES": {DownloaderMiddleware: 0}, } - async def start(self): - yield Request("data:,a", meta={"response_rough_size": 1}) - yield Request("data:,b") - def parse(self, response): pass crawler = get_crawler(TestSpider) - self.caplog.clear() - with self.caplog.at_level("INFO"): - await maybe_deferred_to_future(crawler.crawl()) + await maybe_deferred_to_future(crawler.crawl()) - active_size_log_count = sum( - 1 - for r in self.caplog.records - if str(r.msg).startswith("The active response size") - and r.levelname == "INFO" - ) - assert active_size_log_count == 1 + assert sizes == [(0, 1)] + assert crawler.engine + assert crawler.engine.downloader.middleware._rough_sizes == {} @coroutine_test async def test_rough_size_triggers_backout(self): - """Rough sizes of in-flight requests count toward the backpressure limit. + """Rough sizes of requests being downloaded count toward the limit, even + if their responses turn out to be empty.""" - With RESPONSE_MAX_ACTIVE_SIZE=512 and RESPONSE_ROUGH_SIZE=1024, even a - response with an empty body should trigger the backout log.""" + class SlowDown: + """Keeps requests in flight long enough for the engine to check for + backout while their rough size is being counted.""" + + async def process_request(self, request): + await sleep(0.01) class TestSpider(Spider): name = "test" start_urls = ["data:,", "data:,"] custom_settings = { + "DOWNLOADER_MIDDLEWARES": {SlowDown: 0}, "RESPONSE_MAX_ACTIVE_SIZE": 512, "RESPONSE_ROUGH_SIZE": 1024, } @@ -254,26 +283,13 @@ class TestResponseRoughSize: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 1 + assert _count_backout_logs(self.caplog) == 1 expected_stats = { "request_backout_seconds/response_max_active_size": gt(0), "request_backout_seconds/total": gt(0), } - assert crawler.stats - actual_stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert expected_stats == actual_stats + assert _backout_stats(crawler) == expected_stats class TestRequestBackout: @@ -296,22 +312,9 @@ class TestRequestBackout: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 0 + assert _count_backout_logs(self.caplog) == 0 - assert crawler.stats - stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert stats == {} + assert _backout_stats(crawler) == {} @coroutine_test async def test_concurrency(self): @@ -351,26 +354,13 @@ class TestRequestBackout: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 0 + assert _count_backout_logs(self.caplog) == 0 expected_stats = { "request_backout_seconds/concurrency": gt(0), "request_backout_seconds/total": gt(0), } - assert crawler.stats - actual_stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert expected_stats == actual_stats + assert _backout_stats(crawler) == expected_stats @coroutine_test async def test_response_size(self): @@ -390,26 +380,13 @@ class TestRequestBackout: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 1 + assert _count_backout_logs(self.caplog) == 1 expected_stats = { "request_backout_seconds/response_max_active_size": gt(0), "request_backout_seconds/total": gt(0), } - assert crawler.stats - actual_stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert expected_stats == actual_stats + assert _backout_stats(crawler) == expected_stats @coroutine_test async def test_response_size_process_request(self): @@ -434,26 +411,13 @@ class TestRequestBackout: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 1 + assert _count_backout_logs(self.caplog) == 1 expected_stats = { "request_backout_seconds/response_max_active_size": gt(0), "request_backout_seconds/total": gt(0), } - assert crawler.stats - actual_stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert expected_stats == actual_stats + assert _backout_stats(crawler) == expected_stats @coroutine_test async def test_response_size_process_response(self): @@ -478,26 +442,13 @@ class TestRequestBackout: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 1 + assert _count_backout_logs(self.caplog) == 1 expected_stats = { "request_backout_seconds/response_max_active_size": gt(0), "request_backout_seconds/total": gt(0), } - assert crawler.stats - actual_stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert expected_stats == actual_stats + assert _backout_stats(crawler) == expected_stats @coroutine_test async def test_response_size_process_exception(self): @@ -529,26 +480,13 @@ class TestRequestBackout: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 1 + assert _count_backout_logs(self.caplog) == 1 expected_stats = { "request_backout_seconds/response_max_active_size": gt(0), "request_backout_seconds/total": gt(0), } - assert crawler.stats - actual_stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert expected_stats == actual_stats + assert _backout_stats(crawler) == expected_stats @coroutine_test async def test_response_size_download(self): @@ -584,23 +522,10 @@ class TestRequestBackout: with self.caplog.at_level("INFO"): await maybe_deferred_to_future(crawler.crawl()) - matching_log_count = 0 - for log_record in self.caplog.records: - if ( - str(log_record.msg).startswith("The active response size") - and log_record.levelname == "INFO" - ): - matching_log_count += 1 - assert matching_log_count == 1 + assert _count_backout_logs(self.caplog) == 1 expected_stats = { "request_backout_seconds/response_max_active_size": gt(0), "request_backout_seconds/total": gt(0), } - assert crawler.stats - actual_stats = { - k: v - for k, v in crawler.stats.get_stats().items() - if k.startswith("request_backout_seconds/") - } - assert expected_stats == actual_stats + assert _backout_stats(crawler) == expected_stats diff --git a/tests/test_scraper.py b/tests/test_scraper.py index 9f5599fd1..efad11109 100644 --- a/tests/test_scraper.py +++ b/tests/test_scraper.py @@ -145,9 +145,10 @@ class TestSlot: assert actual == 0 assert len(warning_messages) == 1 assert str(warning_messages[0].message) == ( - "scrapy.core.scraper.Slot.active_size is deprecated. Read " - "scrapy.core.downloader.DownloaderMiddlewareManager.response_active_size " - "instead." + "scrapy.core.scraper.Slot.active_size is deprecated. The size of " + "responses in memory is now tracked by the downloader, and no " + "longer has a public API. If you have a use case for one, please " + "open a GitHub issue." ) def test_active_size_write(self): @@ -158,14 +159,8 @@ class TestSlot: assert slot.active_size == 1 assert len(warning_messages) == 1 assert str(warning_messages[0].message) == ( - "scrapy.core.scraper.Slot.active_size is deprecated. " - "scrapy.core.downloader.DownloaderMiddlewareManager.response_active_size " - "might work as a replacement, but modifying that attribute " - "might not be a good idea. If you have a use case for it, you " - "might want to bring it up in a GitHub issue, to discuss with " - "Scrapy developers if there is a better approach, or some " - "change we could implement in Scrapy to improve support for " - "your use case." + "scrapy.core.scraper.Slot.active_size is deprecated, and setting " + "it no longer has any effect on request processing." ) def test_needs_backout_false(self): diff --git a/tox.ini b/tox.ini index 8b0004a6b..8331e2ab6 100644 --- a/tox.ini +++ b/tox.ini @@ -45,7 +45,6 @@ deps = pytest-cov >= 7.0.0 pytest-xdist sybil >= 1.3.0 # https://github.com/cjw296/sybil/issues/20#issuecomment-605433422 - pywin32; sys_platform == "win32" pytest-twisted >= 1.14.3 [testenv]