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.
This commit is contained in:
Andrey Rakhmatullin 2025-12-16 13:35:07 +05:00 committed by GitHub
parent 5a7e132486
commit 180ca39b23
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
5 changed files with 182 additions and 60 deletions

View File

@ -33,9 +33,8 @@ implement the following method:
`item` is an :ref:`item object <item-types>`, see
:ref:`supporting-item-types`.
:meth:`process_item` must either: return an :ref:`item object <item-types>`,
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 <item-types>`
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
=====================

View File

@ -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

View File

@ -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")

View File

@ -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)

View File

@ -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",