Move the implementation to the downloader

This commit is contained in:
Adrián Chaves 2024-07-08 09:17:24 +02:00
parent 2621b97c94
commit 8e534f708a
4 changed files with 119 additions and 97 deletions

View File

@ -4,6 +4,7 @@ import random
import warnings
from collections import deque
from datetime import datetime
from logging import getLogger
from time import time
from typing import (
TYPE_CHECKING,
@ -17,6 +18,7 @@ from typing import (
Union,
cast,
)
from warnings import warn
from twisted.internet import task
from twisted.internet.defer import Deferred
@ -35,6 +37,7 @@ if TYPE_CHECKING:
from scrapy.http import Response
from scrapy.settings import BaseSettings
logger = getLogger(__name__)
_T = TypeVar("_T")
@ -129,6 +132,25 @@ class Downloader:
self.per_slot_settings: Dict[str, Dict[str, Any]] = self.settings.getdict(
"DOWNLOAD_SLOTS", {}
)
self._stats = crawler.stats
default_response_max_active_size = 5000000
scraper_max_active_size = self.settings.getint(
"SCRAPER_MAX_ACTIVE_SIZE", 5000000
)
if scraper_max_active_size != default_response_max_active_size:
warn(
(
"The SCRAPER_MAX_ACTIVE_SIZE setting is deprecated, use "
"RESPONSE_MAX_ACTIVE_SIZE instead."
),
ScrapyDeprecationWarning,
)
default_response_max_active_size = scraper_max_active_size
self._response_max_active_size = self.settings.getint(
"RESPONSE_MAX_ACTIVE_SIZE", default_response_max_active_size
)
self._response_max_active_size_warned = False
def fetch(
self, request: Request, spider: Spider
@ -143,8 +165,42 @@ class Downloader:
)
return dfd.addBoth(_deactivate)
def _count_backout(self, reason):
self._stats.inc_value("request_backouts/total")
self._stats.inc_value(f"request_backouts/{reason}")
def needs_backout(self) -> bool:
return len(self.active) >= self.total_concurrency
if len(self.active) >= self.total_concurrency:
self._count_backout("concurrent_requests")
return True
if (
self._response_max_active_size
and self.middleware.response_active_size >= self._response_max_active_size
):
if not self._response_max_active_size_warned:
self._response_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, is "
f"{self.middleware.response_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, your crawl might get stuck "
f"indefinitely; you can set RESPONSE_MAX_ACTIVE_SIZE to 0 "
f"to disable this limit, but then your code might run out "
f"of memory. This message will only appear the first time "
f"this happens. To learn how often request processing has "
f"been paused during a crawl for this reason, see the "
f"request_backouts/response_max_active_size stat."
)
self._count_backout("response_max_active_size")
return True
return False
def _get_slot(self, request: Request, spider: Spider) -> Tuple[str, Slot]:
key = self.get_slot_key(request)

View File

@ -6,7 +6,9 @@ See documentation in docs/topics/downloader-middleware.rst
from __future__ import annotations
from functools import partial
from typing import TYPE_CHECKING, Any, Callable, Generator, List, Union, cast
from weakref import WeakSet, finalize
from twisted.internet.defer import Deferred, inlineCallbacks
@ -26,6 +28,11 @@ if TYPE_CHECKING:
class DownloaderMiddlewareManager(MiddlewareManager):
component_name = "downloader middleware"
def __init__(self, *args, **kwargs):
super().__init__(*args, **kwargs)
self.response_active_size = 0
self._tracked_responses = WeakSet()
@classmethod
def _get_mwlist_from_settings(cls, settings: BaseSettings) -> List[Any]:
return build_component_list(settings.getwithbase("DOWNLOADER_MIDDLEWARES"))
@ -38,6 +45,29 @@ class DownloaderMiddlewareManager(MiddlewareManager):
if hasattr(mw, "process_exception"):
self.methods["process_exception"].appendleft(mw.process_exception)
def _count_response_size(self, response: Response) -> None:
if response in self._tracked_responses:
return
self._tracked_responses.add(response)
size = len(response.body)
from logging import getLogger
logger = getLogger(__name__)
logger.debug(
f"{self.response_active_size=} += {size=}{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:
from logging import getLogger
logger = getLogger(__name__)
logger.debug(
f"{self.response_active_size=} -= {size=}{self.response_active_size - size}"
)
self.response_active_size -= size
def download(
self,
download_func: Callable[[Request, Spider], Deferred[Response]],
@ -50,42 +80,44 @@ class DownloaderMiddlewareManager(MiddlewareManager):
) -> Generator[Deferred[Any], Any, Union[Response, Request]]:
for method in self.methods["process_request"]:
method = cast(Callable, method)
response = yield deferred_from_coro(
result = yield deferred_from_coro(
method(request=request, spider=spider)
)
if response is not None and not isinstance(
response, (Response, Request)
):
if result is not None and not isinstance(result, (Response, Request)):
raise _InvalidOutput(
f"Middleware {method.__qualname__} must return None, Response or "
f"Request, got {response.__class__.__name__}"
f"Request, got {result.__class__.__name__}"
)
if response:
return response
if isinstance(result, Response):
self._count_response_size(result)
if result:
return result
return (yield download_func(request, spider))
@inlineCallbacks
def process_response(
response: Union[Response, Request]
) -> Generator[Deferred[Any], Any, Union[Response, Request]]:
if response is None:
result = response
if result is None:
raise TypeError("Received None in process_response")
elif isinstance(response, Request):
return response
elif isinstance(result, Request):
return result
for method in self.methods["process_response"]:
method = cast(Callable, method)
response = yield deferred_from_coro(
method(request=request, response=response, spider=spider)
result = yield deferred_from_coro(
method(request=request, response=result, spider=spider)
)
if not isinstance(response, (Response, Request)):
if not isinstance(result, (Response, Request)):
raise _InvalidOutput(
f"Middleware {method.__qualname__} must return Response or Request, "
f"got {type(response)}"
f"got {type(result)}"
)
if isinstance(response, Request):
return response
return response
if isinstance(result, Request):
return result
self._count_response_size(result)
return result
@inlineCallbacks
def process_exception(
@ -94,18 +126,18 @@ class DownloaderMiddlewareManager(MiddlewareManager):
exception = failure.value
for method in self.methods["process_exception"]:
method = cast(Callable, method)
response = yield deferred_from_coro(
result = yield deferred_from_coro(
method(request=request, exception=exception, spider=spider)
)
if response is not None and not isinstance(
response, (Response, Request)
):
if result is not None and not isinstance(result, (Response, Request)):
raise _InvalidOutput(
f"Middleware {method.__qualname__} must return None, Response or "
f"Request, got {type(response)}"
f"Request, got {type(result)}"
)
if response:
return response
if isinstance(result, Response):
self._count_response_size(result)
if result:
return result
return failure
deferred: Deferred[Union[Response, Request]] = mustbe_deferred(

View File

@ -23,7 +23,6 @@ from typing import (
Union,
cast,
)
from warnings import warn
from twisted.internet.defer import Deferred, inlineCallbacks, succeed
from twisted.internet.task import LoopingCall
@ -32,12 +31,7 @@ from twisted.python.failure import Failure
from scrapy import signals
from scrapy.core.downloader import Downloader
from scrapy.core.scraper import Scraper
from scrapy.exceptions import (
CloseSpider,
DontCloseSpider,
IgnoreRequest,
ScrapyDeprecationWarning,
)
from scrapy.exceptions import CloseSpider, DontCloseSpider, IgnoreRequest
from scrapy.http import Request, Response
from scrapy.logformatter import LogFormatter
from scrapy.settings import Settings
@ -123,24 +117,6 @@ class ExecutionEngine:
)
self.start_time: Optional[float] = None
default_response_max_active_size = 5000000
scraper_max_active_size = self.settings.getint(
"SCRAPER_SLOT_MAX_ACTIVE_SIZE", 5000000
)
if scraper_max_active_size != default_response_max_active_size:
warn(
(
"The SCRAPER_SLOT_MAX_ACTIVE_SIZE setting is deprecated, "
"use RESPONSE_MAX_ACTIVE_SIZE instead."
),
ScrapyDeprecationWarning,
)
default_response_max_active_size = scraper_max_active_size
self._response_max_active_size = self.settings.getint(
"RESPONSE_MAX_ACTIVE_SIZE", default_response_max_active_size
)
self._response_max_active_size_warned = False
def _get_scheduler_class(self, settings: BaseSettings) -> Type[BaseScheduler]:
from scrapy.core.scheduler import BaseScheduler
@ -234,43 +210,13 @@ class ExecutionEngine:
if self.spider_is_idle() and self.slot.close_if_idle:
self._spider_idle()
def _count_backout(self, reason):
self.crawler.stats.inc_value("request_backouts/total")
self.crawler.stats.inc_value(f"request_backouts/{reason}")
def _needs_backout(self) -> bool:
assert self.slot is not None # typing
if not self.running or bool(self.slot.closing):
return True
if self.downloader.needs_backout():
self._count_backout("concurrent_requests")
return True
if (
self._response_max_active_size
and Response._ACTIVE_SIZE >= self._response_max_active_size
):
if not self._response_max_active_size_warned:
self._response_max_active_size_warned = True
logger.info(
f"The total size of all bodies from active responses "
f"({Response._ACTIVE_SIZE} B) has surpassed its maximum "
f"({self._response_max_active_size} B) configured through "
f"the RESPONSE_MAX_ACTIVE_SIZE setting; no more requests "
f"will be processed until that changes. If your memory "
f"allows, increase RESPONSE_MAX_ACTIVE_SIZE to increase "
f"your crawl speed. If your code keeps non-weak "
f"references to Response objects, your crawl might get "
f"stuck indefinitely; you can set "
f"RESPONSE_MAX_ACTIVE_SIZE to 0 to disable this pause of "
f"request processing, but then your code might run out of "
f"memory. This message will only appear the first time "
f"this happens. To learn how often request processing has "
f"been paused during a crawl for this reason, see the "
f"request_backouts/response_max_active_size stat."
)
self._count_backout("response_max_active_size")
return True
return False
return (
not self.running
or bool(self.slot.closing)
or self.downloader.needs_backout()
)
def _next_request_from_scheduler(self) -> Optional[Deferred[None]]:
assert self.slot is not None # typing

View File

@ -7,7 +7,6 @@ See documentation in docs/topics/request-response.rst
from __future__ import annotations
from functools import partial
from typing import (
TYPE_CHECKING,
Any,
@ -25,7 +24,6 @@ from typing import (
overload,
)
from urllib.parse import urljoin
from weakref import finalize
from scrapy.exceptions import NotSupported
from scrapy.http.headers import Headers
@ -72,12 +70,6 @@ class Response(object_ref):
Currently used by :meth:`Response.replace`.
"""
_ACTIVE_SIZE = 0
@classmethod
def _finalize(cls, size):
Response._ACTIVE_SIZE -= size
def __init__(
self,
url: str,
@ -100,10 +92,6 @@ class Response(object_ref):
self.ip_address: Union[IPv4Address, IPv6Address, None] = ip_address
self.protocol: Optional[str] = protocol
size = len(self._body)
Response._ACTIVE_SIZE += size
finalize(self, partial(Response._finalize, size))
@property
def cb_kwargs(self) -> Dict[str, Any]:
try: