diff --git a/docs/topics/throttling.rst b/docs/topics/throttling.rst index 565bbd09f..07b0fa0a5 100644 --- a/docs/topics/throttling.rst +++ b/docs/topics/throttling.rst @@ -38,48 +38,15 @@ The main throttling :ref:`settings ` are: - .. setting:: DOWNLOAD_DELAY - :setting:`DOWNLOAD_DELAY` (default: ``1`` (:ref:`fallback `: ``0``)) + :setting:`DOWNLOAD_DELAY` (default: ``1`` + (:ref:`fallback `: ``0``)) Minimum seconds between any two requests to the same domain. - Even if you have multiple slots, requests to the same domain cannot be sent - more frequently than this delay. - To target a specific number of requests per minute (RPM) *per domain*, set this to ``60 / RPM``. For example, ``DOWNLOAD_DELAY = 1.0`` for 60 RPM, or ``DOWNLOAD_DELAY = 2.0`` for 30 RPM. -- .. setting:: DOWNLOAD_DELAY_PER_SLOT - - :setting:`DOWNLOAD_DELAY_PER_SLOT` (default: ``None``) - - Minimum seconds to wait between two consecutive requests sent to the same - download slot. Unlike :setting:`DOWNLOAD_DELAY`, which applies per domain - (:ref:`throttling scope `), this delay is per slot. - - When ``None`` (default), the per-slot delay falls back to - :setting:`DOWNLOAD_DELAY`, preserving the historical behavior where - :setting:`DOWNLOAD_DELAY` was enforced per slot. - - The wait time is measured from when the previous request was sent. - -For example, with ``DOWNLOAD_DELAY = 1.0`` (and, by default, a single download -slot per domain), requests to the same domain are sent at most once per second: - -.. code-block:: text - - T=0.0s: Request 1 sent - T=1.0s: Request 2 sent - T=2.0s: Request 3 sent - -:setting:`DOWNLOAD_DELAY` (per :ref:`throttling scope `) and -:setting:`DOWNLOAD_DELAY_PER_SLOT` (per download slot) are enforced -independently. By default each domain is both its own scope and its own -download slot, so both apply to the same requests and the effective minimum -spacing is the larger of the two; they only differ when requests are grouped -into custom :ref:`scopes ` or download slots (via the -``download_slot`` request meta key). - When configuring these settings, note that: - :setting:`CONCURRENT_REQUESTS` caps ``CONCURRENT_REQUESTS_PER_DOMAIN``. diff --git a/scrapy/core/downloader/__init__.py b/scrapy/core/downloader/__init__.py index a9778d1d7..e22b1192d 100644 --- a/scrapy/core/downloader/__init__.py +++ b/scrapy/core/downloader/__init__.py @@ -1,127 +1,208 @@ from __future__ import annotations import random -from collections import deque +import warnings +from collections.abc import Iterator, Mapping from dataclasses import dataclass, field -from datetime import datetime -from time import monotonic from typing import TYPE_CHECKING, Any -from twisted.internet.defer import Deferred, inlineCallbacks -from twisted.python.failure import Failure +from twisted.internet.defer import inlineCallbacks 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, - CallLaterResult, - call_later, - create_looping_call, -) from scrapy.utils.decorators import _warn_spider_arg -from scrapy.utils.defer import ( - _defer_sleep_async, - _schedule_coro, - deferred_from_coro, - maybe_deferred_to_future, -) -from scrapy.utils.deprecate import warn_on_deprecated_spider_attribute +from scrapy.utils.defer import _defer_sleep_async, deferred_from_coro +from scrapy.utils.deprecate import create_deprecated_class from scrapy.utils.httpobj import urlparse_cached if TYPE_CHECKING: from collections.abc import Generator - from twisted.internet.task import LoopingCall + from twisted.internet.defer import Deferred from scrapy.crawler import Crawler from scrapy.http import Response from scrapy.settings import BaseSettings from scrapy.signalmanager import SignalManager + from scrapy.throttling import ThrottlingScopeManagerProtocol @dataclass(slots=True, eq=False) -class Slot: +class _Slot: """Downloader slot""" - concurrency: int - delay: float - randomize_delay: bool - active: set[Request] = field(default_factory=set, init=False, repr=False) - queue: deque[tuple[Request, Deferred[Response]]] = field( - default_factory=deque, init=False, repr=False - ) 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) + + +Slot = create_deprecated_class( + "Slot", + _Slot, + old_class_path="scrapy.core.downloader.Slot", + subclass_warn_message=("{cls} inherits from the deprecated Slot class."), + instance_warn_message=("The Slot class is deprecated."), +) + + +class _DeprecatedSlotView: + """Deprecated per-domain slot view backed by the downloader and throttler.""" + + __slots__ = ("_downloader", "_key", "_scope") + + def __init__( + self, + downloader: Downloader, + key: str, + scope: ThrottlingScopeManagerProtocol | None, + ) -> None: + self._downloader = downloader + self._key = key + self._scope = scope + + @property + def active(self) -> set[Request]: + return { + r + for r in self._downloader.active + if r.meta.get(Downloader.DOWNLOAD_SLOT) == self._key + } + + @property + def transferring(self) -> set[Request]: + return { + r + for r in self._downloader._transferring + if r.meta.get(Downloader.DOWNLOAD_SLOT) == self._key + } + + @property + def lastseen(self) -> float: + return 0.0 + + @property + def delay(self) -> float: + if self._scope is not None: + return self._scope._delay # type: ignore[union-attr] + return 0.0 + + @delay.setter + def delay(self, value: float) -> None: + if self._scope is not None: + self._scope._delay = value # type: ignore[union-attr] + + @property + def randomize_delay(self) -> bool: + if self._scope is not None: + return self._scope._randomize # type: ignore[union-attr] + return False + + @property + def concurrency(self) -> int: + warnings.warn( + "Slot.concurrency is deprecated. Per-slot concurrency limits are " + "now managed by the throttling system.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + if self._scope is not None: + return self._scope._concurrency or 0 # type: ignore[union-attr] + return 0 def free_transfer_slots(self) -> int: - return self.concurrency - len(self.transferring) + concurrency = self._scope._concurrency if self._scope is not None else 0 # type: ignore[union-attr] + return concurrency - len(self.transferring) def download_delay(self) -> float: + delay = self.delay if self.randomize_delay: - return random.uniform(0.5 * self.delay, 1.5 * self.delay) # noqa: S311 - return self.delay + return random.uniform(0.5 * delay, 1.5 * delay) # noqa: S311 + return delay def close(self) -> None: - if self.latercall: - self.latercall.cancel() - self.latercall = None + pass + + def __repr__(self) -> str: + return f"_DeprecatedSlotView({self._key!r})" def __str__(self) -> str: - return ( - f"" + return f"_DeprecatedSlotView({self._key!r})" + + +class _DeprecatedSlotsView(Mapping): + """Deprecated mapping view of active downloads, keyed by slot name.""" + + __slots__ = ("_downloader", "_throttler") + + def __init__(self, downloader: Downloader, throttler: Any) -> None: + self._downloader = downloader + self._throttler = throttler + + def _active_keys(self) -> set[str]: + return { + r.meta[Downloader.DOWNLOAD_SLOT] + for r in self._downloader.active + if Downloader.DOWNLOAD_SLOT in r.meta + } + + def __getitem__(self, key: str) -> _DeprecatedSlotView: + if key not in self._active_keys(): + raise KeyError(key) + scope = ( + self._throttler._get_scope_manager(key) + if self._throttler is not None + else None ) + return _DeprecatedSlotView(self._downloader, key, scope) + def __iter__(self) -> Iterator[str]: + return iter(self._active_keys()) -def _get_concurrency_delay( - concurrency: int, spider: Spider, settings: BaseSettings -) -> tuple[int, float]: - delay: float = settings.getfloat("DOWNLOAD_DELAY") - if hasattr(spider, "download_delay"): - delay = spider.download_delay - if settings.get("DOWNLOAD_DELAY_PER_SLOT") is not None: - delay = settings.getfloat("DOWNLOAD_DELAY_PER_SLOT") + def __len__(self) -> int: + return len(self._active_keys()) - if hasattr(spider, "max_concurrent_requests"): # pragma: no cover - warn_on_deprecated_spider_attribute( - "max_concurrent_requests", "CONCURRENT_REQUESTS" - ) - concurrency = spider.max_concurrent_requests - - return concurrency, delay + def __contains__(self, key: object) -> bool: + return key in self._active_keys() class Downloader: DOWNLOAD_SLOT = "download_slot" - _SLOT_GC_INTERVAL: float = 60.0 # seconds def __init__(self, crawler: Crawler): self.crawler: Crawler = crawler self.settings: BaseSettings = crawler.settings self.signals: SignalManager = crawler.signals - self.slots: dict[str, Slot] = {} self.active: set[Request] = set() + self._transferring: set[Request] = set() self.handlers: DownloadHandlers = DownloadHandlers(crawler) self.total_concurrency: int = self.settings.getint("CONCURRENT_REQUESTS") - self.domain_concurrency: int = self.settings.getint( - "CONCURRENT_REQUESTS_PER_DOMAIN" - ) self.ip_concurrency: int = self.settings.getint("CONCURRENT_REQUESTS_PER_IP") - self.randomize_delay: bool = self.settings.getbool("RANDOMIZE_DOWNLOAD_DELAY") self.middleware: DownloaderMiddlewareManager = ( DownloaderMiddlewareManager.from_crawler(crawler) ) - self._slot_gc_loop: AsyncioLoopingCall | LoopingCall | None = None self.per_slot_settings: dict[str, dict[str, Any]] = self.settings.getdict( "DOWNLOAD_SLOTS" ) + if self.per_slot_settings: + warnings.warn( + "The DOWNLOAD_SLOTS setting is deprecated. Use THROTTLING_SCOPES for " + "per-domain configuration instead.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + for slot_settings in self.per_slot_settings.values(): + for deprecated_key in ("concurrency", "delay", "randomize_delay"): + if deprecated_key in slot_settings: + warnings.warn( + f"The '{deprecated_key}' key in DOWNLOAD_SLOTS is deprecated." + " Use THROTTLING_SCOPES to configure per-domain settings" + " instead.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) @inlineCallbacks @_warn_spider_arg @@ -142,94 +223,53 @@ class Downloader: def needs_backout(self) -> bool: return len(self.active) >= self.total_concurrency + @property + def slots(self) -> _DeprecatedSlotsView: + warnings.warn( + "Downloader.slots is deprecated. Use the throttling manager API instead.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + return _DeprecatedSlotsView(self, self.crawler.throttler) + @_warn_spider_arg def _get_slot( self, request: Request, spider: Spider | None = None - ) -> tuple[str, Slot]: - key = self.get_slot_key(request) - if key not in self.slots: - assert self.crawler.spider - slot_settings = self.per_slot_settings.get(key, {}) - conc = self.ip_concurrency or self.domain_concurrency - conc, delay = _get_concurrency_delay( - conc, self.crawler.spider, self.settings - ) - conc, delay = ( - slot_settings.get("concurrency", conc), - slot_settings.get("delay", delay), - ) - randomize_delay = slot_settings.get("randomize_delay", self.randomize_delay) - new_slot = Slot(conc, delay, randomize_delay) - self.slots[key] = new_slot - self._start_slot_gc() + ) -> tuple[str, _DeprecatedSlotView]: + key = self._get_slot_key(request) + scope = ( + self.crawler.throttler._get_scope_manager(key) + if self.crawler.throttler is not None + else None + ) + return key, _DeprecatedSlotView(self, key, scope) - return key, self.slots[key] + def _get_slot_key(self, request: Request) -> str: + throttler = self.crawler.throttler + if throttler is not None: + return throttler.get_slot_key(request) + return self.get_slot_key(request) def get_slot_key(self, request: Request) -> str: - meta_slot: str | None = request.meta.get(self.DOWNLOAD_SLOT) - if meta_slot is not None: - return meta_slot - - key = urlparse_cached(request).hostname or "" + key = urlparse_cached(request).netloc or "" if self.ip_concurrency: key = dnscache.get(key, key) - return key - # passed as download_func into self.middleware.download() in self.fetch() async def _enqueue_request(self, request: Request) -> Response: - key, slot = self._get_slot(request) + key = self._get_slot_key(request) request.meta[self.DOWNLOAD_SLOT] = key - slot.active.add(request) self.signals.send_catch_log( signal=signals.request_reached_downloader, request=request, spider=self.crawler.spider, ) - d: Deferred[Response] = Deferred() - slot.queue.append((request, d)) - self._process_queue(slot) + return await self._download(request) + + async def _download(self, request: Request) -> Response: + self._transferring.add(request) try: - return await maybe_deferred_to_future(d) # fired in _wait_for_download() - finally: - slot.active.remove(request) - - def _process_queue(self, slot: Slot) -> None: - if slot.latercall: - # block processing until slot.latercall is called - return - - # Delay queue processing if a download_delay is configured - now = monotonic() - delay = slot.download_delay() - if delay: - penalty = delay - now + slot.lastseen - if penalty > 0: - slot.latercall = call_later(penalty, self._latercall, slot) - return - - # Process enqueued requests if there are free slots to transfer for this slot - while slot.queue and slot.free_transfer_slots() > 0: - slot.lastseen = now - request, queue_dfd = slot.queue.popleft() - _schedule_coro(self._wait_for_download(slot, request, queue_dfd)) - # prevent burst if inter-request delays were configured - if delay: - self._process_queue(slot) - break - - def _latercall(self, slot: Slot) -> None: - slot.latercall = None - self._process_queue(slot) - - async def _download(self, slot: Slot, request: Request) -> Response: - # The order is very important for the following logic. Do not change! - slot.transferring.add(request) - try: - # 1. Download the response response: Response = await self.handlers.download_request_async(request) - # 2. Notify response_downloaded listeners about the recent download - # before querying queue for next request self.signals.send_catch_log( signal=signals.response_downloaded, response=response, @@ -241,46 +281,12 @@ class Downloader: await _defer_sleep_async() raise finally: - # 3. After response arrives, remove the request from transferring - # state to free up the transferring slot so it can be used by the - # following requests (perhaps those which came from the downloader - # middleware itself) - slot.transferring.remove(request) - self._process_queue(slot) + self._transferring.discard(request) self.signals.send_catch_log( signal=signals.request_left_downloader, request=request, spider=self.crawler.spider, ) - async def _wait_for_download( - self, slot: Slot, request: Request, queue_dfd: Deferred[Response] - ) -> None: - try: - response = await self._download(slot, request) - except Exception: - queue_dfd.errback(Failure()) - else: - queue_dfd.callback(response) # awaited in _enqueue_request() - def close(self) -> None: - self._stop_slot_gc() - for slot in self.slots.values(): - slot.close() - - def _slot_gc(self, age: float = 60) -> None: - mintime = monotonic() - age - for key, slot in list(self.slots.items()): - if not slot.active and slot.lastseen + slot.delay < mintime: - self.slots.pop(key).close() - - def _start_slot_gc(self) -> None: - if self._slot_gc_loop: - return - self._slot_gc_loop = create_looping_call(self._slot_gc) - self._slot_gc_loop.start(self._SLOT_GC_INTERVAL, now=False) - - def _stop_slot_gc(self) -> None: - if self._slot_gc_loop: - self._slot_gc_loop.stop() - self._slot_gc_loop = None + pass diff --git a/scrapy/core/downloader/handlers/_base_streaming.py b/scrapy/core/downloader/handlers/_base_streaming.py index 5a669c732..b7aba8fd0 100644 --- a/scrapy/core/downloader/handlers/_base_streaming.py +++ b/scrapy/core/downloader/handlers/_base_streaming.py @@ -86,8 +86,16 @@ class BaseStreamingDownloadHandler(BaseHttpDownloadHandler, ABC, Generic[_Respon self._proxy_auth_encoding: str = crawler.settings.get("HTTPPROXY_AUTH_ENCODING") # these are useful for many handlers but used in different ways by them self._pool_size_total: int = crawler.settings.getint("CONCURRENT_REQUESTS") - self._pool_size_per_host: int = crawler.settings.getint( - "CONCURRENT_REQUESTS_PER_DOMAIN" + scope_concurrencies = [ + scope["concurrency"] + for scope in crawler.settings.getdict("THROTTLING_SCOPES").values() + if "concurrency" in scope + ] + self._pool_size_per_host: int = max( + [ + crawler.settings.getint("THROTTLING_SCOPE_CONCURRENCY"), + *scope_concurrencies, + ] ) @staticmethod diff --git a/scrapy/core/downloader/handlers/http11.py b/scrapy/core/downloader/handlers/http11.py index 1fc504c59..c399551e7 100644 --- a/scrapy/core/downloader/handlers/http11.py +++ b/scrapy/core/downloader/handlers/http11.py @@ -92,8 +92,16 @@ class HTTP11DownloadHandler(BaseHttpDownloadHandler): from twisted.internet import reactor self._pool: HTTPConnectionPool = HTTPConnectionPool(reactor, persistent=True) - self._pool.maxPersistentPerHost = crawler.settings.getint( - "CONCURRENT_REQUESTS_PER_DOMAIN" + scope_concurrencies = [ + scope["concurrency"] + for scope in crawler.settings.getdict("THROTTLING_SCOPES").values() + if "concurrency" in scope + ] + self._pool.maxPersistentPerHost = max( + [ + crawler.settings.getint("THROTTLING_SCOPE_CONCURRENCY"), + *scope_concurrencies, + ] ) self._pool._factory.noisy = False diff --git a/scrapy/crawler.py b/scrapy/crawler.py index c18e061f8..492877e7a 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -98,6 +98,8 @@ class Crawler: return self.addons.load_settings(self.settings) + self._warn_on_deprecated_default_settings() + self._apply_spider_download_delay() self.stats = load_object(self.settings["STATS_CLASS"])(self) lf_cls: type[LogFormatter] = load_object(self.settings["LOG_FORMATTER"]) @@ -157,6 +159,46 @@ class Crawler: "Overridden settings:\n%(settings)s", {"settings": pprint.pformat(d)} ) + def _apply_spider_download_delay(self) -> None: + spider = self.spider if self.spider is not None else self.spidercls + if not hasattr(spider, "download_delay"): + return + delay_prio = self.settings.getpriority("DOWNLOAD_DELAY") or 0 + if delay_prio >= SETTINGS_PRIORITIES["spider"]: + warnings.warn( + "The 'download_delay' spider attribute is deprecated. " + "It is also being ignored because DOWNLOAD_DELAY is already set " + "at spider or higher priority. Remove the 'download_delay' " + "attribute from your spider.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + else: + warnings.warn( + "The 'download_delay' spider attribute is deprecated. Use the " + "DOWNLOAD_DELAY setting or per-domain THROTTLING_SCOPES instead.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + self.settings.set( + "DOWNLOAD_DELAY", spider.download_delay, priority="spider" + ) + + def _warn_on_deprecated_default_settings(self) -> None: + default_priority = SETTINGS_PRIORITIES["default"] + for setting_name, current_default, future_default in ( + ("THROTTLING_SCOPE_CONCURRENCY", 8, 1), + ): + if self.settings.getpriority(setting_name) == default_priority: + warnings.warn( + f"The default value of {setting_name} will change from " + f"{current_default!r} to {future_default!r} in a future " + f"Scrapy version. Explicitly set {setting_name} in your " + f"settings to silence this warning.", + category=ScrapyDeprecationWarning, + stacklevel=3, + ) + def _apply_reactorless_default_settings(self) -> None: """Change some setting defaults when not using a Twisted reactor. diff --git a/scrapy/pqueues.py b/scrapy/pqueues.py index dfee687b4..4394f7407 100644 --- a/scrapy/pqueues.py +++ b/scrapy/pqueues.py @@ -268,18 +268,20 @@ class DownloaderInterface: def __init__(self, crawler: Crawler): assert crawler.engine self.downloader: Downloader = crawler.engine.downloader + self._throttler: ThrottlingManagerProtocol | None = crawler.throttler - def stats(self, possible_slots: Iterable[str]) -> list[tuple[int, str]]: - return [(self._active_downloads(slot), slot) for slot in possible_slots] + def stats(self, possible_slots: Iterable[str]) -> list[tuple[float, str]]: + return [(self._slot_load(slot), slot) for slot in possible_slots] def get_slot_key(self, request: Request) -> str: + if self._throttler is not None: + return self._throttler.get_slot_key(request) return self.downloader.get_slot_key(request) - def _active_downloads(self, slot: str) -> int: - """Return a number of requests in a Downloader for a given slot""" - if slot not in self.downloader.slots: - return 0 - return len(self.downloader.slots[slot].active) + def _slot_load(self, slot: str) -> float: + if self._throttler is not None: + return self._throttler.get_scope_load(slot) + return 0.0 class DownloaderAwarePriorityQueue: @@ -363,19 +365,19 @@ class DownloaderAwarePriorityQueue: for slot, startprios in slot_startprios.items(): self.pqueues[slot] = self.pqfactory(slot, startprios) - def _next_slot(self, stats: list[tuple[int, str]], *, update_state: bool) -> str: + def _next_slot(self, stats: list[tuple[float, str]], *, update_state: bool) -> str: last = self._last_selected_slot - min_active: int | None = None + min_load: float | None = None best_slot: str | None = None best_slot_after_last: str | None = None - for active, slot in stats: - if min_active is None or active < min_active: - min_active = active + for load, slot in stats: + if min_load is None or load < min_load: + min_load = load best_slot = slot best_slot_after_last = None if last is not None and slot > last: best_slot_after_last = slot - elif active == min_active: + elif load == min_load: if best_slot is None or slot < best_slot: best_slot = slot if ( diff --git a/scrapy/settings/__init__.py b/scrapy/settings/__init__.py index be298e9b1..833c71f44 100644 --- a/scrapy/settings/__init__.py +++ b/scrapy/settings/__init__.py @@ -166,6 +166,22 @@ class BaseSettings(MutableMapping[str, Any]): stacklevel=2, ) + if name == "THROTTLING_SCOPE_CONCURRENCY": + per_domain_prio = self.getpriority("CONCURRENT_REQUESTS_PER_DOMAIN") or 0 + new_prio = self.getpriority(name) or 0 + if ( + per_domain_prio > SETTINGS_PRIORITIES["default"] + and new_prio <= SETTINGS_PRIORITIES["default"] + ): + warnings.warn( + "The CONCURRENT_REQUESTS_PER_DOMAIN setting is deprecated, use " + "THROTTLING_SCOPE_CONCURRENCY instead.", + ScrapyDeprecationWarning, + stacklevel=2, + ) + per_domain_val = self["CONCURRENT_REQUESTS_PER_DOMAIN"] + return per_domain_val if per_domain_val is not None else default + return self[name] if self[name] is not None else default def getbool(self, name: str, default: bool = False) -> bool: diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 586c2ecfd..6bfd73777 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -294,7 +294,6 @@ DNS_TIMEOUT = 60 DOWNLOAD_BIND_ADDRESS = None DOWNLOAD_DELAY = 0 -DOWNLOAD_DELAY_PER_SLOT = None DOWNLOAD_FAIL_ON_DATALOSS = True @@ -598,6 +597,7 @@ THROTTLING_SCOPES = {} THROTTLING_WINDOW = 60.0 THROTTLING_ROBOTSTXT_OBEY = True THROTTLING_ROBOTSTXT_MAX_DELAY = 60.0 +THROTTLING_SCOPE_CONCURRENCY = 8 THROTTLING_SCOPE_LIMIT = 100000 THROTTLING_SCOPE_MAX_IDLE = 3600.0 THROTTLING_DEBUG = False diff --git a/scrapy/throttling.py b/scrapy/throttling.py index ebb1a903e..040e15a28 100644 --- a/scrapy/throttling.py +++ b/scrapy/throttling.py @@ -5,6 +5,7 @@ import datetime as dt import logging import random import time +import warnings from collections import OrderedDict from collections.abc import Awaitable, Callable, Iterable from email.utils import parsedate_to_datetime @@ -16,6 +17,7 @@ from twisted.internet.defer import Deferred from typing_extensions import NotRequired, Self from scrapy import signals +from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.utils.asyncio import sleep, wait_for_first from scrapy.utils.httpobj import urlparse_cached from scrapy.utils.misc import build_from_crawler, load_object @@ -352,6 +354,15 @@ class ThrottlingManagerProtocol(Protocol): requests are time-blocked. """ + def get_slot_key(self, request: Request) -> str: + """Return a single string key for *request*, derived from its scopes. + + For a single scope this is the scope ID itself; for multiple scopes + the sorted scope IDs are joined with ``"+"``. This is the synchronous + counterpart of :meth:`get_scopes`, used wherever a plain string key is + needed (e.g. scheduler priority queues). + """ + def get_scope_load(self, scope_id: str) -> float: """Return the current load of the scope identified by *scope_id*: its active sends divided by its concurrency limit (or by the global @@ -519,8 +530,22 @@ class ThrottlingManager: scopes = request.meta.get("throttling_scopes") if scopes is not None: return cast("RequestScopes", scopes) + download_slot = request.meta.get("download_slot") + if download_slot is not None: + warnings.warn( + "The 'download_slot' request meta key is deprecated. Use " + "'throttling_scopes' instead.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + return download_slot return urlparse_cached(request).netloc + def get_slot_key(self, request: Request) -> str: + scopes = self._resolve_scopes_sync(request) + scope_ids = sorted(iter_scopes(scopes)) + return "+".join(scope_ids) if scope_ids else "" + def _cached_scope_values( self, request: Request ) -> list[tuple[ScopeID, float | None]]: @@ -1121,7 +1146,7 @@ class ThrottlingScopeManager: elif self._rampup_enabled: self._concurrency = self._min_concurrency else: - self._concurrency = None + self._concurrency = settings.getint("THROTTLING_SCOPE_CONCURRENCY") or None # Used as the load denominator when the scope enforces no explicit # concurrency limit (see get_load()). self._global_concurrency: int = settings.getint("CONCURRENT_REQUESTS") diff --git a/scrapy/utils/test.py b/scrapy/utils/test.py index 90ed70262..fc86cbe7b 100644 --- a/scrapy/utils/test.py +++ b/scrapy/utils/test.py @@ -66,13 +66,13 @@ def get_crawler( will be used to populate the crawler settings with a project level priority. """ - # When needed, useful settings can be added here, e.g. ones that prevent - # deprecation warnings. settings: dict[str, Any] = { "TELNETCONSOLE_ENABLED": False, **get_reactor_settings(), **(settings_dict or {}), } + if prevent_warnings: + settings.setdefault("THROTTLING_SCOPE_CONCURRENCY", 8) runner: CrawlerRunnerBase if is_reactor_installed(): runner = CrawlerRunner(settings) diff --git a/tests/test_core_downloader.py b/tests/test_core_downloader.py index cb14b588d..19272a0a9 100644 --- a/tests/test_core_downloader.py +++ b/tests/test_core_downloader.py @@ -14,7 +14,7 @@ from twisted.web import server, static from twisted.web.client import Agent, BrowserLikePolicyForHTTPS, readBody from twisted.web.client import Response as TxResponse -from scrapy.core.downloader import Downloader, Slot, _get_concurrency_delay, tls +from scrapy.core.downloader import Downloader, Slot, _Slot, tls from scrapy.core.downloader.contextfactory import ( _load_context_factory_from_settings, _ScrapyClientContextFactory, @@ -41,33 +41,20 @@ if TYPE_CHECKING: class TestSlot: def test_repr(self): - slot = Slot(concurrency=8, delay=0.1, randomize_delay=True) - assert repr(slot) == "Slot(concurrency=8, delay=0.1, randomize_delay=True)" + slot = _Slot() + assert repr(slot) == "_Slot()" + def test_deprecated(self): + with pytest.warns(ScrapyDeprecationWarning, match="Slot class is deprecated"): + Slot() -class TestGetConcurrencyDelay: - def test_default(self): - crawler = get_crawler() - concurrency, delay = _get_concurrency_delay( - 8, DefaultSpider(), crawler.settings - ) - assert (concurrency, delay) == (8, 0.0) + def test_deprecated_subclass(self): + with pytest.warns( + ScrapyDeprecationWarning, match="inherits from the deprecated Slot" + ): - def test_spider_download_delay(self): - crawler = get_crawler() - spider = DefaultSpider() - spider.download_delay = 2.5 - _concurrency, delay = _get_concurrency_delay(8, spider, crawler.settings) - assert delay == 2.5 - - def test_download_delay_per_slot(self): - crawler = get_crawler(settings_dict={"DOWNLOAD_DELAY_PER_SLOT": 3.0}) - # DOWNLOAD_DELAY_PER_SLOT takes precedence over both DOWNLOAD_DELAY and - # the spider's download_delay attribute. - spider = DefaultSpider() - spider.download_delay = 2.5 - _concurrency, delay = _get_concurrency_delay(8, spider, crawler.settings) - assert delay == 3.0 + class MySlot(Slot): + pass @pytest.mark.requires_reactor # this test is related to the Twisted HTTP code diff --git a/tests/test_crawler.py b/tests/test_crawler.py index 0cddfd0ed..5588775cf 100644 --- a/tests/test_crawler.py +++ b/tests/test_crawler.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio import logging import re +import warnings from pathlib import Path from typing import Any, ClassVar @@ -81,6 +82,60 @@ class TestCrawler(TestBaseCrawler): assert not settings.frozen assert crawler.settings.frozen + @pytest.mark.parametrize( + ("setting_name", "current_default", "future_default"), + [ + ("THROTTLING_SCOPE_CONCURRENCY", 8, 1), + ], + ) + def test_deprecated_default_settings_warn( + self, setting_name: str, current_default: Any, future_default: Any + ) -> None: + crawler = Crawler(DefaultSpider) + with pytest.warns( + ScrapyDeprecationWarning, + match=rf"The default value of {setting_name} will change from {current_default!r} to {future_default!r}", + ): + crawler._apply_settings() + + @pytest.mark.parametrize( + ("setting_name", "current_default"), + [ + ("THROTTLING_SCOPE_CONCURRENCY", 8), + ], + ) + def test_deprecated_default_settings_no_warn_when_set( + self, setting_name: str, current_default: int + ) -> None: + crawler = Crawler(DefaultSpider, {setting_name: current_default}) + with warnings.catch_warnings(): + warnings.simplefilter("error", ScrapyDeprecationWarning) + crawler._apply_settings() + + def test_spider_download_delay_deprecated(self) -> None: + class DelaySpider(DefaultSpider): + download_delay = 2.5 + + crawler = Crawler(DelaySpider, {"THROTTLING_SCOPE_CONCURRENCY": 8}) + with pytest.warns( + ScrapyDeprecationWarning, match="'download_delay' spider attribute" + ): + crawler._apply_settings() + assert crawler.settings.getfloat("DOWNLOAD_DELAY") == 2.5 + + def test_spider_download_delay_overridden_by_setting(self) -> None: + class DelaySpider(DefaultSpider): + download_delay = 2.5 + + crawler = Crawler(DelaySpider, {"THROTTLING_SCOPE_CONCURRENCY": 8}) + crawler.settings.set("DOWNLOAD_DELAY", 5.0, priority="spider") + with pytest.warns( + ScrapyDeprecationWarning, + match="'download_delay' spider attribute.*being ignored", + ): + crawler._apply_settings() + assert crawler.settings.getfloat("DOWNLOAD_DELAY") == 5.0 + def test_crawler_accepts_dict(self) -> None: crawler = get_crawler(DefaultSpider, {"foo": "bar"}) assert crawler.settings["foo"] == "bar" diff --git a/tests/test_downloaderslotssettings.py b/tests/test_downloaderslotssettings.py index 9d58a6e09..5540c2e98 100644 --- a/tests/test_downloaderslotssettings.py +++ b/tests/test_downloaderslotssettings.py @@ -1,39 +1,32 @@ -import time from typing import Any +from urllib.parse import urlparse import pytest from scrapy import Request -from scrapy.core.downloader import Downloader, Slot +from scrapy.core.downloader import Downloader from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.utils.spider import DefaultSpider from scrapy.utils.test import get_crawler from tests.mockserver.http import MockServer from tests.spiders import MetaSpider -from tests.utils.decorators import coroutine_test, inline_callbacks_test +from tests.utils.decorators import coroutine_test class DownloaderSlotsSettingsTestSpider(MetaSpider): name = "downloader_slots" custom_settings = { - "DOWNLOAD_DELAY": 1, - "RANDOMIZE_DOWNLOAD_DELAY": False, "DOWNLOAD_SLOTS": { - "quotes.toscrape.com": { - "concurrency": 1, - "delay": 2, - "randomize_delay": False, - "throttle": False, - }, - "books.toscrape.com": {"delay": 3, "randomize_delay": False}, + "quotes.toscrape.com": {"concurrency": 1}, + "books.toscrape.com": {"concurrency": 2}, }, } def __init__(self, *args: Any, **kwargs: Any): super().__init__(*args, **kwargs) assert self.mockserver - self.default_slot = self.mockserver.host + self.default_slot = urlparse(self.mockserver.url("/")).netloc self.times: dict[str, list[float]] = {} async def start(self): @@ -45,65 +38,86 @@ class DownloaderSlotsSettingsTestSpider(MetaSpider): def parse(self, response): slot = response.meta.get("download_slot", self.default_slot) - self.times[slot].append(time.time()) + self.times[slot].append(response.meta.get("download_latency")) url = self.mockserver.url(f"/?downloader_slot={slot}&req=2") yield Request(url, callback=self.not_parse, meta={"download_slot": slot}) def not_parse(self, response): slot = response.meta.get("download_slot", self.default_slot) - self.times[slot].append(time.time()) - - -class TestCrawl: - @classmethod - def setup_class(cls): - cls.mockserver = MockServer() - cls.mockserver.__enter__() - - @classmethod - def teardown_class(cls): - cls.mockserver.__exit__(None, None, None) - - @inline_callbacks_test - def test_delay(self): - crawler = get_crawler(DownloaderSlotsSettingsTestSpider) - yield crawler.crawl(mockserver=self.mockserver) - slots = crawler.engine.downloader.slots - times = crawler.spider.times - tolerance = 0.3 - - delays_real = {k: v[1] - v[0] for k, v in times.items()} - error_delta = { - k: 1 - min(delays_real[k], v.delay) / max(delays_real[k], v.delay) - for k, v in slots.items() - } - - assert max(list(error_delta.values())) < tolerance + self.times[slot].append(response.meta.get("download_latency")) @coroutine_test -async def test_params(): - params = { - "concurrency": 1, - "delay": 2, - "randomize_delay": False, - } - settings = { - "DOWNLOAD_SLOTS": { - "example.com": params, - }, - } +async def test_concurrency_key_deprecated(): + settings = {"DOWNLOAD_SLOTS": {"example.com": {"concurrency": 3}}} crawler = get_crawler(DefaultSpider, settings_dict=settings) crawler.spider = crawler._create_spider() + with pytest.warns(ScrapyDeprecationWarning) as warns: + downloader = Downloader(crawler) + messages = [str(w.message) for w in warns] + assert any("DOWNLOAD_SLOTS setting is deprecated" in m for m in messages) + assert any("'concurrency' key in DOWNLOAD_SLOTS" in m for m in messages) + downloader._get_slot(Request("https://example.com")) + downloader.close() + + +@coroutine_test +async def test_download_slots_deprecated(): + settings = {"DOWNLOAD_SLOTS": {"example.com": {"concurrency": 2}}} + crawler = get_crawler(DefaultSpider, settings_dict=settings) + crawler.spider = crawler._create_spider() + with pytest.warns( + ScrapyDeprecationWarning, match="DOWNLOAD_SLOTS setting is deprecated" + ): + Downloader(crawler).close() + + +@coroutine_test +async def test_slots_deprecated(): + crawler = get_crawler(DefaultSpider) + crawler.spider = crawler._create_spider() downloader = Downloader(crawler) request = Request("https://example.com") - _, actual = downloader._get_slot(request) + request.meta[Downloader.DOWNLOAD_SLOT] = "example.com" + downloader.active.add(request) + with pytest.warns(ScrapyDeprecationWarning, match="Downloader.slots is deprecated"): + slot = downloader.slots.get("example.com") + assert slot is not None + assert isinstance(slot.active, set) + assert request in slot.active + downloader.active.discard(request) + downloader.close() + + +@coroutine_test +async def test_download_slot_meta_deprecated(): + crawler = get_crawler(DefaultSpider) + crawler.spider = crawler._create_spider() + downloader = Downloader(crawler) + request = Request("https://example.com") + request.meta["download_slot"] = "custom" + with pytest.warns( + ScrapyDeprecationWarning, match="'download_slot' request meta key is deprecated" + ): + key, _ = downloader._get_slot(request) + downloader.close() + assert key == "custom" + + +@coroutine_test +async def test_delay_deprecated(): + settings = { + "DOWNLOAD_SLOTS": {"example.com": {"delay": 2, "randomize_delay": False}} + } + crawler = get_crawler(DefaultSpider, settings_dict=settings) + crawler.spider = crawler._create_spider() + with pytest.warns(ScrapyDeprecationWarning) as warns: + downloader = Downloader(crawler) + messages = [str(w.message) for w in warns] + assert any("DOWNLOAD_SLOTS setting is deprecated" in m for m in messages) + assert any("'delay' key in DOWNLOAD_SLOTS" in m for m in messages) + downloader._get_slot(Request("https://example.com")) downloader.close() - expected = Slot(**params) - for param in params: - assert getattr(expected, param) == getattr(actual, param), ( - f"Slot.{param}: {getattr(expected, param)!r} != {getattr(actual, param)!r}" - ) @coroutine_test @@ -122,7 +136,7 @@ async def test_get_slot_deprecated_spider_arg(): downloader.close() assert key1 == key2 - assert slot1 == slot2 + assert slot1._key == slot2._key @pytest.mark.parametrize( @@ -132,6 +146,7 @@ async def test_get_slot_deprecated_spider_arg(): "scrapy.pqueues.DownloaderAwarePriorityQueue", ], ) +@pytest.mark.filterwarnings("ignore::scrapy.exceptions.ScrapyDeprecationWarning") @coroutine_test async def test_none_slot_with_priority_queue( mockserver: MockServer, priority_queue_class: str diff --git a/tests/test_engine.py b/tests/test_engine.py index 2b15c40d1..a7f81e19d 100644 --- a/tests/test_engine.py +++ b/tests/test_engine.py @@ -856,7 +856,7 @@ class TestEngineCloseSpider: engine = ExecutionEngine(crawler, lambda _: None) crawler.engine = engine await engine.open_spider_async() - del engine.downloader.slots + engine.downloader.close = Mock(side_effect=Exception("close failed")) await engine.close_spider_async() assert "Downloader close failure" in caplog.text diff --git a/tests/test_pqueues.py b/tests/test_pqueues.py index c45687a3f..cf7f94f04 100644 --- a/tests/test_pqueues.py +++ b/tests/test_pqueues.py @@ -4,7 +4,6 @@ from unittest.mock import Mock import pytest import queuelib -from scrapy.core.downloader import Downloader from scrapy.http.request import Request from scrapy.pqueues import ( DownloaderAwarePriorityQueue, @@ -192,22 +191,22 @@ class TestDownloaderAwarePriorityQueue: # No active downloads are tracked in the downloader, so every slot has # the same score and tie-breaking must not starve a slot. req_a1 = Request("https://example.org/a1") - req_a1.meta[Downloader.DOWNLOAD_SLOT] = "slot-a" + req_a1.meta["throttling_scopes"] = "slot-a" req_b1 = Request("https://example.org/b1") - req_b1.meta[Downloader.DOWNLOAD_SLOT] = "slot-b" + req_b1.meta["throttling_scopes"] = "slot-b" req_a2 = Request("https://example.org/a2") - req_a2.meta[Downloader.DOWNLOAD_SLOT] = "slot-a" + req_a2.meta["throttling_scopes"] = "slot-a" req_b2 = Request("https://example.org/b2") - req_b2.meta[Downloader.DOWNLOAD_SLOT] = "slot-b" + req_b2.meta["throttling_scopes"] = "slot-b" for request in (req_a1, req_b1, req_a2, req_b2): self.queue.push(request) slots = [ - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], + self.queue.pop().meta["throttling_scopes"], + self.queue.pop().meta["throttling_scopes"], + self.queue.pop().meta["throttling_scopes"], + self.queue.pop().meta["throttling_scopes"], ] assert slots == ["slot-a", "slot-b", "slot-a", "slot-b"] @@ -216,48 +215,49 @@ class TestDownloaderAwarePriorityQueue: # If the selected slot becomes empty, rotation should continue from # that slot marker to avoid restarting from the smallest slot. req_a1 = Request("https://example.org/a1") - req_a1.meta[Downloader.DOWNLOAD_SLOT] = "slot-a" + req_a1.meta["throttling_scopes"] = "slot-a" req_a2 = Request("https://example.org/a2") - req_a2.meta[Downloader.DOWNLOAD_SLOT] = "slot-a" + req_a2.meta["throttling_scopes"] = "slot-a" req_b1 = Request("https://example.org/b1") - req_b1.meta[Downloader.DOWNLOAD_SLOT] = "slot-b" + req_b1.meta["throttling_scopes"] = "slot-b" req_c1 = Request("https://example.org/c1") - req_c1.meta[Downloader.DOWNLOAD_SLOT] = "slot-c" + req_c1.meta["throttling_scopes"] = "slot-c" for request in (req_a1, req_a2, req_b1, req_c1): self.queue.push(request) slots = [ - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], - self.queue.pop().meta[Downloader.DOWNLOAD_SLOT], + self.queue.pop().meta["throttling_scopes"], + self.queue.pop().meta["throttling_scopes"], + self.queue.pop().meta["throttling_scopes"], + self.queue.pop().meta["throttling_scopes"], ] assert slots == ["slot-a", "slot-b", "slot-c", "slot-a"] def test_pop_prefers_slot_with_fewer_active_downloads(self): - downloader = self.queue._downloader_interface.downloader + throttler = self.queue._downloader_interface._throttler + assert throttler is not None req_a = Request("https://example.org/a") - req_a.meta[Downloader.DOWNLOAD_SLOT] = "slot-a" + req_a.meta["throttling_scopes"] = "slot-a" req_b = Request("https://example.org/b") - req_b.meta[Downloader.DOWNLOAD_SLOT] = "slot-b" + req_b.meta["throttling_scopes"] = "slot-b" req_c = Request("https://example.org/c") - req_c.meta[Downloader.DOWNLOAD_SLOT] = "slot-c" + req_c.meta["throttling_scopes"] = "slot-c" for req in (req_a, req_b, req_c): self.queue.push(req) - downloader.increment("slot-a") - downloader.increment("slot-c") + throttler._get_scope_manager("slot-a")._active = 1 + throttler._get_scope_manager("slot-c")._active = 1 popped = self.queue.pop() assert popped.url == req_b.url def test_contains(self): req = Request("https://example.org/") - req.meta[Downloader.DOWNLOAD_SLOT] = "example-slot" + req.meta["throttling_scopes"] = "example-slot" assert "example-slot" not in self.queue self.queue.push(req) assert "example-slot" in self.queue diff --git a/tests/test_throttling.py b/tests/test_throttling.py index abcb29604..ef07bc381 100644 --- a/tests/test_throttling.py +++ b/tests/test_throttling.py @@ -570,8 +570,12 @@ class TestThrottlingScopeManager: scope.record_backoff(now=0.0) assert scope._delay == pytest.approx(3.0) - def test_no_scope_concurrency_limit_by_default(self): + def test_default_scope_concurrency(self): scope = _scope_manager() + assert scope._concurrency == 8 + + def test_no_scope_concurrency_limit_when_zero(self): + scope = _scope_manager(settings={"THROTTLING_SCOPE_CONCURRENCY": 0}) assert scope._concurrency is None for _ in range(100): scope.record_sent(now=0.0)