diff --git a/docs/topics/feed-exports.rst b/docs/topics/feed-exports.rst index 66768c97b..2f686fd0f 100644 --- a/docs/topics/feed-exports.rst +++ b/docs/topics/feed-exports.rst @@ -218,12 +218,13 @@ passed through the following settings: .. _temporary security credentials: https://docs.aws.amazon.com/IAM/latest/UserGuide/security-creds.html -You can also define a custom ACL, custom endpoint, and region name for exported -feeds using these settings: +You can also define a custom ACL, custom endpoint, region name and connection +pool size for exported feeds using these settings: - :setting:`FEED_STORAGE_S3_ACL` - :setting:`AWS_ENDPOINT_URL` - :setting:`AWS_REGION_NAME` +- :setting:`AWS_MAX_POOL_CONNECTIONS` The default value for the ``overwrite`` key in the :setting:`FEEDS` for this storage backend is: ``True``. diff --git a/docs/topics/media-pipeline.rst b/docs/topics/media-pipeline.rst index 4ceb4732a..b16066d0c 100644 --- a/docs/topics/media-pipeline.rst +++ b/docs/topics/media-pipeline.rst @@ -268,6 +268,9 @@ For self-hosting you also might feel the need not to use SSL and not to verify S AWS_USE_SSL = False # or True (None by default) AWS_VERIFY = False # or True (None by default) +To reuse connections for as many files as you check or upload in parallel, set +:setting:`AWS_MAX_POOL_CONNECTIONS` accordingly. + .. _canned ACLs: https://docs.aws.amazon.com/AmazonS3/latest/userguide/acl-overview.html#canned-acl .. _Minio: https://github.com/minio/minio .. _Zenko CloudServer: https://www.zenko.io/cloudserver/ diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 8068824e3..1b6851d04 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -458,6 +458,26 @@ Default: ``None`` Endpoint URL used for S3-like storage, for example Minio or s3.scality. +.. setting:: AWS_MAX_POOL_CONNECTIONS + +AWS_MAX_POOL_CONNECTIONS +------------------------ + +.. versionadded:: VERSION + +Default: ``None`` + +Maximum number of connections that AWS clients, such as those of the +:ref:`S3 feed storage backend ` and of the +:ref:`S3 media pipeline storage backend `, keep in their +connection pool. + +If ``None``, the value of :setting:`REACTOR_THREADPOOL_MAXSIZE` is used. + +Values lower than the number of parallel AWS calls do not limit those calls, but +their connections are closed instead of reused, which hurts performance, and +``Connection pool is full, discarding connection`` warnings are logged. + .. setting:: AWS_REGION_NAME AWS_REGION_NAME diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index c2997921d..448279546 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -28,6 +28,7 @@ from scrapy import Spider, signals from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning from scrapy.extensions.postprocessing import PostProcessingManager from scrapy.utils.asyncio import is_asyncio_available, run_in_thread +from scrapy.utils.boto import _get_max_pool_connections from scrapy.utils.conf import feed_complete_default_values_from_settings from scrapy.utils.defer import deferred_from_coro, ensure_awaitable from scrapy.utils.ftp import ftp_store_file @@ -213,11 +214,14 @@ class S3FeedStorage(BlockingFeedStorage): feed_options: dict[str, Any] | None = None, session_token: str | None = None, region_name: str | None = None, + max_pool_connections: int | None = None, ): try: import boto3.session # noqa: PLC0415 except ImportError: raise NotConfigured("missing boto3 library") from None + from botocore.config import Config # noqa: PLC0415 + u = urlparse(uri) assert u.hostname self.bucketname: str = u.hostname @@ -228,6 +232,7 @@ class S3FeedStorage(BlockingFeedStorage): self.acl: str | None = acl self.endpoint_url: str | None = endpoint_url self.region_name: str | None = region_name + self.max_pool_connections: int | None = max_pool_connections boto3_session = boto3.session.Session() self.s3_client = boto3_session.client( @@ -237,6 +242,11 @@ class S3FeedStorage(BlockingFeedStorage): aws_session_token=self.session_token, endpoint_url=self.endpoint_url, region_name=self.region_name, + config=( + Config(max_pool_connections=self.max_pool_connections) + if self.max_pool_connections is not None + else None + ), ) if feed_options and feed_options.get("overwrite", True) is False: @@ -262,6 +272,7 @@ class S3FeedStorage(BlockingFeedStorage): acl=crawler.settings["FEED_STORAGE_S3_ACL"] or None, endpoint_url=crawler.settings["AWS_ENDPOINT_URL"] or None, region_name=crawler.settings["AWS_REGION_NAME"] or None, + max_pool_connections=_get_max_pool_connections(crawler.settings), feed_options=feed_options, ) diff --git a/scrapy/pipelines/files.py b/scrapy/pipelines/files.py index 8e8082332..55a3676e5 100644 --- a/scrapy/pipelines/files.py +++ b/scrapy/pipelines/files.py @@ -37,7 +37,7 @@ from scrapy.pipelines.media import ( _MediaRequestFiltered, ) from scrapy.utils.asyncio import run_in_thread -from scrapy.utils.boto import is_botocore_available +from scrapy.utils.boto import _get_max_pool_connections, is_botocore_available from scrapy.utils.datatypes import CaseInsensitiveDict from scrapy.utils.defer import deferred_from_coro, ensure_awaitable from scrapy.utils.ftp import ftp_store_file @@ -164,6 +164,9 @@ class S3FilesStore: AWS_REGION_NAME = None AWS_USE_SSL = None AWS_VERIFY = None + # Overridden from settings.AWS_MAX_POOL_CONNECTIONS in + # FilesPipeline.from_crawler(); None means the botocore default + AWS_MAX_POOL_CONNECTIONS: int | None = None POLICY = "private" # Overridden from settings.FILES_STORE_S3_ACL in FilesPipeline.from_crawler() HEADERS: ClassVar[dict[str, str]] = { @@ -174,7 +177,13 @@ class S3FilesStore: if not is_botocore_available(): raise NotConfigured("missing botocore library") import botocore.session # noqa: PLC0415 + from botocore.config import Config # noqa: PLC0415 + config = ( + Config(max_pool_connections=self.AWS_MAX_POOL_CONNECTIONS) + if self.AWS_MAX_POOL_CONNECTIONS is not None + else None + ) session = botocore.session.get_session() self.s3_client = session.create_client( "s3", @@ -185,6 +194,7 @@ class S3FilesStore: region_name=self.AWS_REGION_NAME, use_ssl=self.AWS_USE_SSL, verify=self.AWS_VERIFY, + config=config, ) if not uri.startswith("s3://"): raise ValueError(f"Incorrect URI scheme in {uri}, expected 's3'") @@ -522,6 +532,7 @@ class FilesPipeline(MediaPipeline): s3store.AWS_REGION_NAME = settings["AWS_REGION_NAME"] s3store.AWS_USE_SSL = settings["AWS_USE_SSL"] s3store.AWS_VERIFY = settings["AWS_VERIFY"] + s3store.AWS_MAX_POOL_CONNECTIONS = _get_max_pool_connections(settings) s3store.POLICY = settings["FILES_STORE_S3_ACL"] gcs_store: type[GCSFilesStore] = cast( diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 993f2436d..a44b36c8a 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -28,6 +28,7 @@ __all__ = [ "AUTOTHROTTLE_TARGET_CONCURRENCY", "AWS_ACCESS_KEY_ID", "AWS_ENDPOINT_URL", + "AWS_MAX_POOL_CONNECTIONS", "AWS_REGION_NAME", "AWS_SECRET_ACCESS_KEY", "AWS_SESSION_TOKEN", @@ -229,6 +230,7 @@ AUTOTHROTTLE_TARGET_CONCURRENCY = 1.0 AWS_ACCESS_KEY_ID = None AWS_SECRET_ACCESS_KEY = None AWS_ENDPOINT_URL = None +AWS_MAX_POOL_CONNECTIONS = None AWS_REGION_NAME = None AWS_SESSION_TOKEN = None AWS_USE_SSL = None diff --git a/scrapy/utils/boto.py b/scrapy/utils/boto.py index 2a77ee2ac..76ee0e7ec 100644 --- a/scrapy/utils/boto.py +++ b/scrapy/utils/boto.py @@ -1,7 +1,22 @@ """Boto/botocore helpers""" +from __future__ import annotations + from importlib.util import find_spec +from typing import TYPE_CHECKING + +if TYPE_CHECKING: + from scrapy.settings import BaseSettings def is_botocore_available() -> bool: return find_spec("botocore") is not None + + +def _get_max_pool_connections(settings: BaseSettings) -> int: + """Return the maximum number of connections that AWS clients may keep in + their connection pool. + """ + return settings.getint("AWS_MAX_POOL_CONNECTIONS") or settings.getint( + "REACTOR_THREADPOOL_MAXSIZE" + ) diff --git a/tests/test_feedexport_storages.py b/tests/test_feedexport_storages.py index 66488540f..b1e3787fd 100644 --- a/tests/test_feedexport_storages.py +++ b/tests/test_feedexport_storages.py @@ -381,6 +381,41 @@ class TestS3FeedStorage: assert storage.region_name == region_name assert storage.s3_client._client_config.region_name == region_name + def test_init_without_max_pool_connections(self) -> None: + storage = S3FeedStorage("s3://mybucket/export.csv", "access_key", "secret_key") + assert storage.max_pool_connections is None + config: Any = storage.s3_client.meta.config + assert config.max_pool_connections == 10 + + def test_init_with_max_pool_connections(self) -> None: + storage = S3FeedStorage( + "s3://mybucket/export.csv", + "access_key", + "secret_key", + max_pool_connections=30, + ) + assert storage.max_pool_connections == 30 + config: Any = storage.s3_client.meta.config + assert config.max_pool_connections == 30 + + @pytest.mark.parametrize( + ("settings", "expected"), + [ + ({}, 10), + ({"REACTOR_THREADPOOL_MAXSIZE": 20}, 20), + ({"AWS_MAX_POOL_CONNECTIONS": 30}, 30), + ({"AWS_MAX_POOL_CONNECTIONS": 30, "REACTOR_THREADPOOL_MAXSIZE": 20}, 30), + ], + ) + def test_from_crawler_max_pool_connections( + self, settings: dict[str, Any], expected: int + ) -> None: + crawler = get_crawler(settings_dict=settings) + storage = S3FeedStorage.from_crawler(crawler, "s3://mybucket/export.csv") + assert storage.max_pool_connections == expected + config: Any = storage.s3_client.meta.config + assert config.max_pool_connections == expected + @coroutine_test async def test_store_without_acl(self): storage = S3FeedStorage( diff --git a/tests/test_pipeline_files.py b/tests/test_pipeline_files.py index 97daa514e..4e7fb118b 100644 --- a/tests/test_pipeline_files.py +++ b/tests/test_pipeline_files.py @@ -895,6 +895,33 @@ class TestS3FilesStore: stub.assert_no_pending_responses() + def test_default_max_pool_connections(self) -> None: + store = S3FilesStore("s3://mybucket/prefix/") + config: Any = store.s3_client.meta.config + assert config.max_pool_connections == 10 + + @pytest.mark.parametrize( + ("settings", "expected"), + [ + ({}, 10), + ({"REACTOR_THREADPOOL_MAXSIZE": 20}, 20), + ({"AWS_MAX_POOL_CONNECTIONS": 30}, 30), + ({"AWS_MAX_POOL_CONNECTIONS": 30, "REACTOR_THREADPOOL_MAXSIZE": 20}, 30), + ], + ) + def test_max_pool_connections( + self, monkeypatch: pytest.MonkeyPatch, settings: dict[str, Any], expected: int + ) -> None: + # restores the value that FilesPipeline.from_crawler() sets on the class + monkeypatch.setattr(S3FilesStore, "AWS_MAX_POOL_CONNECTIONS", None) + crawler = get_crawler( + settings_dict={"FILES_STORE": "s3://mybucket/prefix/", **settings} + ) + store = FilesPipeline.from_crawler(crawler).store + assert isinstance(store, S3FilesStore) + config: Any = store.s3_client.meta.config + assert config.max_pool_connections == expected + class TestGCSFilesStore: @staticmethod