From 180ca39b230590a2a2862c8d18c574661b1d16ad Mon Sep 17 00:00:00 2001 From: Andrey Rakhmatullin Date: Tue, 16 Dec 2025 13:35:07 +0500 Subject: [PATCH] Deprecate returning Deferreds from pipeline methods (#7179) * Add tests for exceptions in pipelines. * Deprecate returning Deferreds from pipeline process_item(). * Deprecate returning Deferreds from pipeline {open,close}_spider(). * Update the custom pipeline docs. --- docs/topics/item-pipeline.rst | 7 +- scrapy/middleware.py | 8 +- scrapy/pipelines/__init__.py | 59 ++++++++++--- scrapy/utils/defer.py | 7 ++ tests/test_pipelines.py | 161 +++++++++++++++++++++++++--------- 5 files changed, 182 insertions(+), 60 deletions(-) diff --git a/docs/topics/item-pipeline.rst b/docs/topics/item-pipeline.rst index e67cf06c8..8194fe043 100644 --- a/docs/topics/item-pipeline.rst +++ b/docs/topics/item-pipeline.rst @@ -33,9 +33,8 @@ implement the following method: `item` is an :ref:`item object `, see :ref:`supporting-item-types`. - :meth:`process_item` must either: return an :ref:`item object `, - return a :class:`~twisted.internet.defer.Deferred` or raise a - :exc:`~scrapy.exceptions.DropItem` exception. + :meth:`process_item` must either return an :ref:`item object ` + or raise a :exc:`~scrapy.exceptions.DropItem` exception. Dropped items are no longer processed by further pipeline components. @@ -52,6 +51,8 @@ Additionally, they may also implement the following methods: This method is called when the spider is closed. +Any of these methods may be defined as a coroutine function (``async def``). + Item pipeline example ===================== diff --git a/scrapy/middleware.py b/scrapy/middleware.py index 83362784b..c4e0ce16e 100644 --- a/scrapy/middleware.py +++ b/scrapy/middleware.py @@ -138,17 +138,21 @@ class MiddlewareManager(ABC): *args: Any, add_spider: bool = False, always_add_spider: bool = False, + warn_deferred: bool = False, ) -> _T: methods = cast( "Iterable[Callable[Concatenate[_T, _P], _T]]", self.methods[methodname] ) for method in methods: + warn = global_object_name(method) if warn_deferred else None if always_add_spider or ( add_spider and method in self._mw_methods_requiring_spider ): - obj = await ensure_awaitable(method(obj, *(*args, self._spider))) + obj = await ensure_awaitable( + method(obj, *(*args, self._spider)), _warn=warn + ) else: - obj = await ensure_awaitable(method(obj, *args)) + obj = await ensure_awaitable(method(obj, *args), _warn=warn) return obj def open_spider(self, spider: Spider) -> Deferred[list[None]]: # pragma: no cover diff --git a/scrapy/pipelines/__init__.py b/scrapy/pipelines/__init__.py index 398e895b8..21ad6fb1a 100644 --- a/scrapy/pipelines/__init__.py +++ b/scrapy/pipelines/__init__.py @@ -6,6 +6,7 @@ See documentation in docs/item-pipeline.rst from __future__ import annotations +import asyncio import warnings from typing import TYPE_CHECKING, Any, cast @@ -13,16 +14,13 @@ from twisted.internet.defer import Deferred, DeferredList from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.middleware import MiddlewareManager +from scrapy.utils.asyncio import is_asyncio_available from scrapy.utils.conf import build_component_list -from scrapy.utils.defer import ( - deferred_from_coro, - maybe_deferred_to_future, - maybeDeferred_coro, -) +from scrapy.utils.defer import deferred_from_coro, ensure_awaitable, maybeDeferred_coro from scrapy.utils.python import global_object_name if TYPE_CHECKING: - from collections.abc import Callable, Iterable + from collections.abc import Awaitable, Callable, Coroutine, Iterable from twisted.python.failure import Failure @@ -58,12 +56,19 @@ class ItemPipelineManager(MiddlewareManager): return deferred_from_coro(self.process_item_async(item)) async def process_item_async(self, item: Any) -> Any: - return await self._process_chain("process_item", item, add_spider=True) + return await self._process_chain( + "process_item", item, add_spider=True, warn_deferred=True + ) - def _process_parallel(self, methodname: str) -> Deferred[list[None]]: - methods = cast("Iterable[Callable[..., None]]", self.methods[methodname]) + def _process_parallel_dfd(self, methodname: str) -> Deferred[list[None]]: + methods = cast( + "Iterable[Callable[..., Coroutine[Any, Any, None] | Deferred[None] | None]]", + self.methods[methodname], + ) - def get_dfd(method: Callable[..., None]) -> Deferred[None]: + def get_dfd( + method: Callable[..., Coroutine[Any, Any, None] | Deferred[None] | None], + ) -> Deferred[None]: if method in self._mw_methods_requiring_spider: return maybeDeferred_coro(method, self._spider) return maybeDeferred_coro(method) @@ -80,6 +85,32 @@ class ItemPipelineManager(MiddlewareManager): d2.addErrback(eb) return d2 + async def _process_parallel_asyncio(self, methodname: str) -> list[None]: + methods = cast( + "Iterable[Callable[..., Coroutine[Any, Any, None] | Deferred[None] | None]]", + self.methods[methodname], + ) + if not methods: + return [] + + def get_awaitable( + method: Callable[..., Coroutine[Any, Any, None] | Deferred[None] | None], + ) -> Awaitable[None]: + if method in self._mw_methods_requiring_spider: + result = method(self._spider) + else: + result = method() + return ensure_awaitable(result, _warn=global_object_name(method)) + + awaitables = [get_awaitable(m) for m in methods] + await asyncio.gather(*awaitables) + return [None for _ in methods] + + async def _process_parallel(self, methodname: str) -> list[None]: + if is_asyncio_available(): + return await self._process_parallel_asyncio(methodname) + return await self._process_parallel_dfd(methodname) + def open_spider(self, spider: Spider) -> Deferred[list[None]]: warnings.warn( f"{global_object_name(type(self))}.open_spider() is deprecated, use open_spider_async() instead.", @@ -87,10 +118,10 @@ class ItemPipelineManager(MiddlewareManager): stacklevel=2, ) self._set_compat_spider(spider) - return self._process_parallel("open_spider") + return deferred_from_coro(self._process_parallel("open_spider")) async def open_spider_async(self) -> None: - await maybe_deferred_to_future(self._process_parallel("open_spider")) + await self._process_parallel("open_spider") def close_spider(self, spider: Spider) -> Deferred[list[None]]: warnings.warn( @@ -99,7 +130,7 @@ class ItemPipelineManager(MiddlewareManager): stacklevel=2, ) self._set_compat_spider(spider) - return self._process_parallel("close_spider") + return deferred_from_coro(self._process_parallel("close_spider")) async def close_spider_async(self) -> None: - await maybe_deferred_to_future(self._process_parallel("close_spider")) + await self._process_parallel("close_spider") diff --git a/scrapy/utils/defer.py b/scrapy/utils/defer.py index 2a09f99bb..afbe680ef 100644 --- a/scrapy/utils/defer.py +++ b/scrapy/utils/defer.py @@ -27,6 +27,7 @@ from twisted.python import failure from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.utils.asyncio import is_asyncio_available +from scrapy.utils.python import global_object_name if TYPE_CHECKING: from collections.abc import AsyncIterator, Callable @@ -426,6 +427,12 @@ def maybeDeferred_coro( return fail(failure.Failure(captureVars=Deferred.debug)) if isinstance(result, Deferred): + warnings.warn( + f"{global_object_name(f)} returned a Deferred, this is deprecated." + f" Please refactor this function to return a coroutine.", + ScrapyDeprecationWarning, + stacklevel=2, + ) return result if asyncio.isfuture(result) or inspect.isawaitable(result): return deferred_from_coro(result) diff --git a/tests/test_pipelines.py b/tests/test_pipelines.py index 8db9431ea..fc61d61ec 100644 --- a/tests/test_pipelines.py +++ b/tests/test_pipelines.py @@ -1,7 +1,8 @@ import asyncio +from typing import Any import pytest -from twisted.internet.defer import Deferred, inlineCallbacks, succeed +from twisted.internet.defer import Deferred, fail, succeed from scrapy import Request, Spider, signals from scrapy.crawler import Crawler @@ -42,6 +43,12 @@ class DeferredPipeline: item["pipeline_passed"] = True return item + def open_spider(self): + return succeed(None) + + def close_spider(self): + return succeed(None) + def process_item(self, item): d = Deferred() d.addCallback(self.cb) @@ -83,6 +90,36 @@ class AsyncDefNotAsyncioPipeline: return item +class ProcessItemExceptionPipeline: + def process_item(self, item): + raise ValueError("process_item error") + + +class ProcessItemExceptionDeferredPipeline: + def process_item(self, item): + return fail(ValueError("process_item error")) + + +class ProcessItemExceptionAsyncPipeline: + async def process_item(self, item): + raise ValueError("process_item error") + + +class OpenSpiderExceptionPipeline: + def open_spider(self): + raise ValueError("open_spider error") + + +class OpenSpiderExceptionDeferredPipeline: + def open_spider(self): + return fail(ValueError("open_spider error")) + + +class OpenSpiderExceptionAsyncPipeline: + async def open_spider(self): + raise ValueError("open_spider error") + + class ItemSpider(Spider): name = "itemspider" @@ -94,59 +131,55 @@ class ItemSpider(Spider): class TestPipeline: - @classmethod - def setup_class(cls): - cls.mockserver = MockServer() - cls.mockserver.__enter__() - - @classmethod - def teardown_class(cls): - cls.mockserver.__exit__(None, None, None) - def _on_item_scraped(self, item): assert isinstance(item, dict) assert item.get("pipeline_passed") self.items.append(item) - def _create_crawler(self, pipeline_class): + def _create_crawler(self, pipeline_class: type) -> Crawler: settings = { "ITEM_PIPELINES": {pipeline_class: 1}, } crawler = get_crawler(ItemSpider, settings) crawler.signals.connect(self._on_item_scraped, signals.item_scraped) - self.items = [] + self.items: list[Any] = [] return crawler - @inlineCallbacks - def test_simple_pipeline(self): - crawler = self._create_crawler(SimplePipeline) - yield crawler.crawl(mockserver=self.mockserver) + @pytest.mark.parametrize( + "pipeline_class", + [ + SimplePipeline, + AsyncDefPipeline, + pytest.param(AsyncDefAsyncioPipeline, marks=pytest.mark.only_asyncio), + pytest.param( + AsyncDefNotAsyncioPipeline, marks=pytest.mark.only_not_asyncio + ), + ], + ) + @deferred_f_from_coro_f + async def test_pipeline(self, mockserver: MockServer, pipeline_class: type) -> None: + crawler = self._create_crawler(pipeline_class) + await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver)) assert len(self.items) == 1 - @inlineCallbacks - def test_deferred_pipeline(self): + @deferred_f_from_coro_f + async def test_pipeline_deferred(self, mockserver: MockServer) -> None: crawler = self._create_crawler(DeferredPipeline) - yield crawler.crawl(mockserver=self.mockserver) - assert len(self.items) == 1 - - @inlineCallbacks - def test_asyncdef_pipeline(self): - crawler = self._create_crawler(AsyncDefPipeline) - yield crawler.crawl(mockserver=self.mockserver) - assert len(self.items) == 1 - - @pytest.mark.only_asyncio - @inlineCallbacks - def test_asyncdef_asyncio_pipeline(self): - crawler = self._create_crawler(AsyncDefAsyncioPipeline) - yield crawler.crawl(mockserver=self.mockserver) - assert len(self.items) == 1 - - @pytest.mark.only_not_asyncio - @inlineCallbacks - def test_asyncdef_not_asyncio_pipeline(self): - crawler = self._create_crawler(AsyncDefNotAsyncioPipeline) - yield crawler.crawl(mockserver=self.mockserver) + with ( + pytest.warns( + ScrapyDeprecationWarning, + match="DeferredPipeline.open_spider returned a Deferred", + ), + pytest.warns( + ScrapyDeprecationWarning, + match="DeferredPipeline.close_spider returned a Deferred", + ), + pytest.warns( + ScrapyDeprecationWarning, + match="DeferredPipeline.process_item returned a Deferred", + ), + ): + await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver)) assert len(self.items) == 1 @deferred_f_from_coro_f @@ -170,6 +203,52 @@ class TestPipeline: assert len(self.items) == 1 + @pytest.mark.parametrize( + "pipeline_class", + [ + ProcessItemExceptionPipeline, + pytest.param( + ProcessItemExceptionDeferredPipeline, + marks=pytest.mark.filterwarnings( + "ignore::scrapy.exceptions.ScrapyDeprecationWarning" + ), + ), + ProcessItemExceptionAsyncPipeline, + ], + ) + @deferred_f_from_coro_f + async def test_process_item_exception( + self, + caplog: pytest.LogCaptureFixture, + mockserver: MockServer, + pipeline_class: type, + ) -> None: + crawler = self._create_crawler(pipeline_class) + await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver)) + assert "Error processing {'field': 42}" in caplog.text + assert "process_item error" in caplog.text + + @pytest.mark.parametrize( + "pipeline_class", + [ + OpenSpiderExceptionPipeline, + pytest.param( + OpenSpiderExceptionDeferredPipeline, + marks=pytest.mark.filterwarnings( + "ignore::scrapy.exceptions.ScrapyDeprecationWarning" + ), + ), + OpenSpiderExceptionAsyncPipeline, + ], + ) + @deferred_f_from_coro_f + async def test_open_spider_exception( + self, mockserver: MockServer, pipeline_class: type + ) -> None: + crawler = self._create_crawler(pipeline_class) + with pytest.raises(ValueError, match="open_spider error"): + await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver)) + class TestCustomPipelineManager: def test_deprecated_process_item_spider_arg(self) -> None: @@ -374,7 +453,7 @@ class TestMiddlewareManagerSpider: match=r"ItemPipelineManager needs to access self\.crawler\.spider but it is None", ), ): - mwman.open_spider(DefaultSpider()) + await maybe_deferred_to_future(mwman.open_spider(DefaultSpider())) with pytest.raises( ValueError, match=r"ItemPipelineManager needs to access self\.crawler\.spider but it is None", @@ -390,7 +469,7 @@ class TestMiddlewareManagerSpider: match=r"ItemPipelineManager needs to access self\.crawler\.spider but it is None", ), ): - mwman.close_spider(DefaultSpider()) + await maybe_deferred_to_future(mwman.close_spider(DefaultSpider())) with pytest.raises( ValueError, match=r"ItemPipelineManager needs to access self\.crawler\.spider but it is None",