Add ItemPipelineManager.process_item_async() and ensure_awaitable(). (#7005)

* Add ItemPipelineManager.process_item_async() and ensure_awaitable().

* Wording.
This commit is contained in:
Andrey Rakhmatullin 2025-08-12 19:57:37 +05:00 committed by GitHub
parent 57f539d7a5
commit d3e15a10cf
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
11 changed files with 245 additions and 72 deletions

View File

@ -50,7 +50,7 @@ class DownloaderMiddlewareManager(MiddlewareManager):
if argument_is_required(download_func, "spider"):
warnings.warn(
"The spider argument of download_func is deprecated"
" and will not be passed in the future Scrapy versions.",
" and will not be passed in future Scrapy versions.",
ScrapyDeprecationWarning,
stacklevel=2,
)

View File

@ -10,7 +10,6 @@ from __future__ import annotations
import asyncio
import logging
import warnings
from collections.abc import AsyncIterator, Callable, Coroutine, Generator
from time import time
from traceback import format_exc
from typing import TYPE_CHECKING, Any
@ -36,6 +35,7 @@ from scrapy.utils.asyncio import (
from scrapy.utils.defer import (
_schedule_coro,
deferred_from_coro,
ensure_awaitable,
maybe_deferred_to_future,
)
from scrapy.utils.deprecate import argument_is_required
@ -45,6 +45,8 @@ from scrapy.utils.python import global_object_name
from scrapy.utils.reactor import CallLaterOnce
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Callable, Coroutine, Generator
from twisted.internet.task import LoopingCall
from scrapy.core.downloader import Downloader
@ -135,7 +137,7 @@ class ExecutionEngine:
if self._downloader_fetch_needs_spider:
warnings.warn(
f"The fetch() method of {global_object_name(downloader_cls)} requires a spider argument,"
f" this is deprecated and the argument will not be passed in the future Scrapy versions.",
f" this is deprecated and the argument will not be passed in future Scrapy versions.",
ScrapyDeprecationWarning,
stacklevel=2,
)
@ -626,11 +628,6 @@ class ExecutionEngine:
self.spider = None
try:
aw = self._spider_closed_callback(spider)
# TODO: replace with ensure_awaitable() when we add it
if isinstance(aw, Coroutine):
await aw
elif isinstance(aw, Deferred):
await maybe_deferred_to_future(aw)
await ensure_awaitable(self._spider_closed_callback(spider))
except Exception:
log_failure("Error running spider_closed_callback")

View File

@ -9,7 +9,7 @@ from collections import deque
from collections.abc import AsyncIterator
from typing import TYPE_CHECKING, Any, TypeVar, Union
from twisted.internet.defer import Deferred, inlineCallbacks, maybeDeferred
from twisted.internet.defer import Deferred, inlineCallbacks
from twisted.python.failure import Failure
from scrapy import Spider, signals
@ -21,6 +21,7 @@ from scrapy.exceptions import (
ScrapyDeprecationWarning,
)
from scrapy.http import Request, Response
from scrapy.pipelines import ItemPipelineManager
from scrapy.utils.asyncio import _parallel_asyncio, is_asyncio_available
from scrapy.utils.defer import (
_defer_sleep_async,
@ -28,12 +29,13 @@ from scrapy.utils.defer import (
aiter_errback,
deferred_f_from_coro_f,
deferred_from_coro,
ensure_awaitable,
iter_errback,
maybe_deferred_to_future,
parallel,
parallel_async,
)
from scrapy.utils.deprecate import argument_is_required
from scrapy.utils.deprecate import argument_is_required, method_is_overridden
from scrapy.utils.log import failure_to_exc_info, logformatter_adapter
from scrapy.utils.misc import load_object, warn_on_generator_with_return_value
from scrapy.utils.python import global_object_name
@ -44,7 +46,6 @@ if TYPE_CHECKING:
from scrapy.crawler import Crawler
from scrapy.logformatter import LogFormatter
from scrapy.pipelines import ItemPipelineManager
from scrapy.signalmanager import SignalManager
@ -109,22 +110,49 @@ class Scraper:
crawler.settings["ITEM_PROCESSOR"]
)
self.itemproc: ItemPipelineManager = itemproc_cls.from_crawler(crawler)
self._itemproc_needs_spider: dict[str, bool] = {}
for method in (
itemproc_methods = [
"open_spider",
"close_spider",
"process_item",
]
if not hasattr(self.itemproc, "process_item_async"):
warnings.warn(
f"{global_object_name(itemproc_cls)} doesn't define a process_item_async() method,"
f" this is deprecated and the method will be required in future Scrapy versions.",
ScrapyDeprecationWarning,
stacklevel=2,
)
itemproc_methods.append("process_item")
self._itemproc_has_process_async = False
elif (
issubclass(itemproc_cls, ItemPipelineManager)
and method_is_overridden(itemproc_cls, ItemPipelineManager, "process_item")
and not method_is_overridden(
itemproc_cls, ItemPipelineManager, "process_item_async"
)
):
warnings.warn(
f"{global_object_name(itemproc_cls)} overrides process_item() but doesn't override process_item_async()."
f" This is deprecated. process_item() will be used, but in future Scrapy versions process_item_async() will be used instead.",
ScrapyDeprecationWarning,
stacklevel=2,
)
itemproc_methods.append("process_item")
self._itemproc_has_process_async = False
else:
self._itemproc_has_process_async = True
self._itemproc_needs_spider: dict[str, bool] = {}
for method in itemproc_methods:
self._itemproc_needs_spider[method] = argument_is_required(
getattr(self.itemproc, method), "spider"
)
if self._itemproc_needs_spider[method]:
warnings.warn(
f"The {method}() method of {global_object_name(itemproc_cls)} requires a spider argument,"
f" this is deprecated and the argument will not be passed in the future Scrapy versions.",
f" this is deprecated and the argument will not be passed in future Scrapy versions.",
ScrapyDeprecationWarning,
stacklevel=2,
)
self.concurrent_items: int = crawler.settings.getint("CONCURRENT_ITEMS")
self.crawler: Crawler = crawler
self.signals: SignalManager = crawler.signals
@ -305,9 +333,7 @@ class Scraper:
output.raiseException()
# else the errback returned actual output (like a callback),
# which needs to be passed to iterate_spider_output()
return await maybe_deferred_to_future(
maybeDeferred(iterate_spider_output, output)
)
return await ensure_awaitable(iterate_spider_output(output))
def handle_spider_error(
self,
@ -467,11 +493,14 @@ class Scraper:
assert self.crawler.spider is not None # typing
self.slot.itemproc_size += 1
try:
if self._itemproc_needs_spider["process_item"]:
d = self.itemproc.process_item(item, self.crawler.spider)
if self._itemproc_has_process_async:
output = await self.itemproc.process_item_async(item)
else:
d = self.itemproc.process_item(item)
output = await maybe_deferred_to_future(d)
if self._itemproc_needs_spider["process_item"]:
d = self.itemproc.process_item(item, self.crawler.spider)
else:
d = self.itemproc.process_item(item)
output = await maybe_deferred_to_future(d)
except DropItem as ex:
logkws = self.logformatter.dropped(item, ex, response, self.crawler.spider)
if logkws is not None:

View File

@ -420,15 +420,13 @@ class SpiderMiddlewareManager(MiddlewareManager):
self._check_deprecated_start_requests_use()
if self._use_start_requests:
sync_start = iter(self._spider.start_requests())
sync_start = await maybe_deferred_to_future(
self._process_chain("process_start_requests", sync_start, self._spider)
sync_start = await self._process_chain(
"process_start_requests", sync_start, self._spider
)
start: AsyncIterator[Any] = as_async_generator(sync_start)
else:
start = self._spider.start()
start = await maybe_deferred_to_future(
self._process_chain("process_start", start)
)
start = await self._process_chain("process_start", start)
return start
def _check_deprecated_start_requests_use(self):

View File

@ -8,7 +8,7 @@ from collections import defaultdict, deque
from typing import TYPE_CHECKING, Any, TypeVar, cast
from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning
from scrapy.utils.defer import process_chain, process_parallel
from scrapy.utils.defer import _process_chain, process_parallel
from scrapy.utils.misc import build_from_crawler, load_object
from scrapy.utils.python import global_object_name
@ -172,11 +172,11 @@ class MiddlewareManager(ABC):
)
return process_parallel(methods, obj, *args)
def _process_chain(self, methodname: str, obj: _T, *args: Any) -> Deferred[_T]:
async def _process_chain(self, methodname: str, obj: _T, *args: Any) -> _T:
methods = cast(
"Iterable[Callable[Concatenate[_T, _P], _T]]", self.methods[methodname]
)
return process_chain(methods, obj, *args)
return await _process_chain(methods, obj, *args)
def open_spider(self, spider: Spider | None = None) -> Deferred[list[None]]:
if spider:

View File

@ -6,11 +6,14 @@ See documentation in docs/item-pipeline.rst
from __future__ import annotations
import warnings
from typing import TYPE_CHECKING, Any
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.middleware import MiddlewareManager
from scrapy.utils.conf import build_component_list
from scrapy.utils.defer import deferred_f_from_coro_f
from scrapy.utils.defer import deferred_from_coro
from scrapy.utils.python import global_object_name
if TYPE_CHECKING:
from twisted.internet.defer import Deferred
@ -29,12 +32,17 @@ class ItemPipelineManager(MiddlewareManager):
def _add_middleware(self, pipe: Any) -> None:
super()._add_middleware(pipe)
if hasattr(pipe, "process_item"):
self.methods["process_item"].append(
deferred_f_from_coro_f(pipe.process_item)
)
self.methods["process_item"].append(pipe.process_item)
def process_item(self, item: Any, spider: Spider | None = None) -> Deferred[Any]:
if spider:
self._warn_spider_arg("process_item")
self._set_compat_spider(spider)
return self._process_chain("process_item", item, self._spider)
warnings.warn(
f"{global_object_name(type(self))}.process_item() is deprecated, use process_item_async() instead.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
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, self._spider)

View File

@ -300,6 +300,11 @@ def process_chain(
**kw: _P.kwargs,
) -> Deferred[_T]:
"""Return a Deferred built by chaining the given callbacks"""
warnings.warn(
"process_chain() is deprecated.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
d: Deferred[_T] = Deferred()
for x in callbacks:
d.addCallback(x, *a, **kw)
@ -307,6 +312,19 @@ def process_chain(
return d
async def _process_chain(
callables: Iterable[Callable[Concatenate[_T, _P], _T | Awaitable[_T]]],
input_: _T,
*a: _P.args,
**kw: _P.kwargs,
) -> _T:
"""Chain the given (potentialy asynchronous) callables."""
result = input_
for callable_ in callables:
result = await ensure_awaitable(callable_(result, *a, **kw))
return result
def process_chain_both(
callbacks: Iterable[Callable[Concatenate[_T, _P], Any]],
errbacks: Iterable[Callable[Concatenate[Failure, _P], Any]],
@ -360,7 +378,7 @@ def iter_errback(
*a: _P.args,
**kw: _P.kwargs,
) -> Iterable[_T]:
"""Wraps an iterable calling an errback if an error is caught while
"""Wrap an iterable calling an errback if an error is caught while
iterating it.
"""
it = iter(iterable)
@ -379,7 +397,7 @@ async def aiter_errback(
*a: _P.args,
**kw: _P.kwargs,
) -> AsyncIterator[_T]:
"""Wraps an async iterable calling an errback if an error is caught while
"""Wrap an async iterable calling an errback if an error is caught while
iterating it. Similar to :func:`scrapy.utils.defer.iter_errback`.
"""
it = aiterable.__aiter__()
@ -401,8 +419,8 @@ def deferred_from_coro(o: _T2) -> _T2: ...
def deferred_from_coro(o: Awaitable[_T] | _T2) -> Deferred[_T] | _T2:
"""Converts a coroutine or other awaitable object into a Deferred,
or returns the object as is if it isn't a coroutine."""
"""Convert a coroutine or other awaitable object into a Deferred,
or return the object as is if it isn't a coroutine."""
if isinstance(o, Deferred):
return o
if inspect.isawaitable(o):
@ -418,7 +436,7 @@ def deferred_from_coro(o: Awaitable[_T] | _T2) -> Deferred[_T] | _T2:
def deferred_f_from_coro_f(
coro_f: Callable[_P, Awaitable[_T]],
) -> Callable[_P, Deferred[_T]]:
"""Converts a coroutine function into a function that returns a Deferred.
"""Convert a coroutine function into a function that returns a Deferred.
The coroutine function will be called at the time when the wrapper is called. Wrapper args will be passed to it.
This is useful for callback chains, as callback functions are called with the previous callback result.
@ -450,10 +468,7 @@ def maybeDeferred_coro(
def deferred_to_future(d: Deferred[_T]) -> Future[_T]:
"""
.. versionadded:: 2.6.0
Return an :class:`asyncio.Future` object that wraps *d*.
"""Return an :class:`asyncio.Future` object that wraps *d*.
This function requires
:class:`~twisted.internet.asyncioreactor.AsyncioSelectorReactor` to be
@ -472,6 +487,8 @@ def deferred_to_future(d: Deferred[_T]) -> Future[_T]:
deferred = self.crawler.engine.download(additional_request)
additional_response = await deferred_to_future(deferred)
.. versionadded:: 2.6.0
.. versionchanged:: VERSION
This function no longer installs an asyncio loop if called before the
Twisted asyncio reactor is installed. A :exc:`RuntimeError` is raised
@ -483,10 +500,7 @@ def deferred_to_future(d: Deferred[_T]) -> Future[_T]:
def maybe_deferred_to_future(d: Deferred[_T]) -> Deferred[_T] | Future[_T]:
"""
.. versionadded:: 2.6.0
Return *d* as an object that can be awaited from a :ref:`Scrapy callable
"""Return *d* as an object that can be awaited from a :ref:`Scrapy callable
defined as a coroutine <coroutine-support>`.
What you can await in Scrapy callables defined as coroutines depends on the
@ -507,6 +521,8 @@ def maybe_deferred_to_future(d: Deferred[_T]) -> Deferred[_T] | Future[_T]:
additional_request = scrapy.Request('https://example.org/price')
deferred = self.crawler.engine.download(additional_request)
additional_response = await maybe_deferred_to_future(deferred)
.. versionadded:: 2.6.0
"""
if not is_asyncio_available():
return d
@ -526,3 +542,32 @@ def _schedule_coro(coro: Coroutine[Any, Any, Any]) -> None:
return
loop = asyncio.get_event_loop()
loop.create_task(coro) # noqa: RUF006
@overload
def ensure_awaitable(o: Awaitable[_T]) -> Awaitable[_T]: ...
@overload
def ensure_awaitable(o: _T) -> Awaitable[_T]: ...
def ensure_awaitable(o: _T | Awaitable[_T]) -> Awaitable[_T]:
"""Convert any value to an awaitable object.
For a :class:`~twisted.internet.defer.Deferred` object, use
:func:`maybe_deferred_to_future` to wrap it into a suitable object. For an
awaitable object of a different type, return it as is. For any other
value, return a coroutine that completes with that value.
.. versionadded:: VERSION
"""
if isinstance(o, Deferred):
return maybe_deferred_to_future(o)
if inspect.isawaitable(o):
return o
async def coro() -> _T:
return o
return coro()

View File

@ -183,6 +183,8 @@ def method_is_overridden(subclass: type, base_class: type, method_name: str) ->
... pass
>>> class Sub4(Sub2):
... pass
>>> method_is_overridden(Base, Base, 'foo')
False
>>> method_is_overridden(Sub1, Base, 'foo')
False
>>> method_is_overridden(Sub2, Base, 'foo')

View File

@ -5,7 +5,7 @@ from unittest import mock
import pytest
from twisted.internet import error
from twisted.internet.defer import Deferred, maybeDeferred
from twisted.internet.defer import Deferred
from twisted.python import failure
from scrapy.downloadermiddlewares.robotstxt import RobotsTxtMiddleware
@ -15,7 +15,11 @@ from scrapy.http import Request, Response, TextResponse
from scrapy.http.request import NO_CALLBACK
from scrapy.settings import Settings
from scrapy.utils.asyncio import call_later
from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future
from scrapy.utils.defer import (
deferred_f_from_coro_f,
ensure_awaitable,
maybe_deferred_to_future,
)
from tests.test_robotstxt_interface import rerp_available
if TYPE_CHECKING:
@ -214,9 +218,7 @@ Disallow: /some/randome/page.html
self, request: Request, middleware: RobotsTxtMiddleware
) -> None:
spider = None # not actually used
result = await maybe_deferred_to_future(
maybeDeferred(middleware.process_request, request, spider) # type: ignore[call-overload]
)
result = await ensure_awaitable(middleware.process_request(request, spider)) # type: ignore[arg-type]
assert result is None
async def assertIgnored(
@ -224,9 +226,7 @@ Disallow: /some/randome/page.html
) -> None:
spider = None # not actually used
with pytest.raises(IgnoreRequest):
await maybe_deferred_to_future(
maybeDeferred(middleware.process_request, request, spider) # type: ignore[call-overload]
)
await ensure_awaitable(middleware.process_request(request, spider)) # type: ignore[arg-type]
def assertRobotsTxtRequested(self, base_url: str) -> None:
calls = self.crawler.engine.download_async.call_args_list

View File

@ -1,12 +1,13 @@
import asyncio
import pytest
from twisted.internet.defer import Deferred, inlineCallbacks
from twisted.internet.defer import Deferred, inlineCallbacks, succeed
from scrapy import Request, Spider, signals
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.pipelines import ItemPipelineManager
from scrapy.utils.asyncio import call_later
from scrapy.utils.conf import build_component_list
from scrapy.utils.defer import (
deferred_f_from_coro_f,
deferred_to_future,
@ -147,20 +148,48 @@ class TestCustomPipelineManager:
itemproc = CustomPipelineManager.from_crawler(crawler)
with pytest.warns(
ScrapyDeprecationWarning,
match=r"Passing a spider argument to CustomPipelineManager.process_item\(\) is deprecated",
match=r"CustomPipelineManager.process_item\(\) is deprecated, use process_item_async\(\)",
):
itemproc.process_item({}, crawler.spider)
@deferred_f_from_coro_f
async def test_deprecated_spider_arg_integration(
self, mockserver: MockServer
) -> None:
async def test_integration_recommended(self, mockserver: MockServer) -> None:
class CustomPipelineManager(ItemPipelineManager):
async def process_item_async(self, item):
return await super().process_item_async(item)
items = []
def _on_item_scraped(item):
assert isinstance(item, dict)
assert item.get("pipeline_passed")
items.append(item)
crawler = get_crawler(
ItemSpider,
{
"ITEM_PROCESSOR": CustomPipelineManager,
"ITEM_PIPELINES": {SimplePipeline: 1},
},
)
crawler.spider = crawler._create_spider()
crawler.signals.connect(_on_item_scraped, signals.item_scraped)
await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver))
assert len(items) == 1
@deferred_f_from_coro_f
async def test_integration_no_async_subclass(self, mockserver: MockServer) -> None:
class CustomPipelineManager(ItemPipelineManager):
def open_spider(self, spider): # pylint: disable=signature-differs
return super().open_spider(spider)
def process_item(self, item, spider): # pylint: disable=signature-differs
return super().process_item(item, spider)
with pytest.warns(
ScrapyDeprecationWarning,
match=r"CustomPipelineManager.process_item\(\) is deprecated, use process_item_async\(\)",
):
return super().process_item(item, spider)
items = []
@ -193,7 +222,73 @@ class TestCustomPipelineManager:
),
pytest.warns(
ScrapyDeprecationWarning,
match=r"Passing a spider argument to CustomPipelineManager.process_item\(\) is deprecated",
match=r"CustomPipelineManager overrides process_item\(\) but doesn't override process_item_async\(\)",
),
):
await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver))
assert len(items) == 1
@deferred_f_from_coro_f
async def test_integration_no_async_not_subclass(
self, mockserver: MockServer
) -> None:
class CustomPipelineManager:
def __init__(self, crawler):
self.pipelines = [
p()
for p in build_component_list(
crawler.settings.getwithbase("ITEM_PIPELINES")
)
]
@classmethod
def from_crawler(cls, crawler):
return cls(crawler)
def open_spider(self, spider):
return succeed(None)
def close_spider(self, spider):
return succeed(None)
def process_item(self, item, spider):
for pipeline in self.pipelines:
item = pipeline.process_item(item, spider)
return succeed(item)
items = []
def _on_item_scraped(item):
assert isinstance(item, dict)
assert item.get("pipeline_passed")
items.append(item)
crawler = get_crawler(
ItemSpider,
{
"ITEM_PROCESSOR": CustomPipelineManager,
"ITEM_PIPELINES": {SimplePipeline: 1},
},
)
crawler.spider = crawler._create_spider()
crawler.signals.connect(_on_item_scraped, signals.item_scraped)
with (
pytest.warns(
ScrapyDeprecationWarning,
match=r"CustomPipelineManager doesn't define a process_item_async\(\) method",
),
pytest.warns(
ScrapyDeprecationWarning,
match=r"The open_spider\(\) method of .+\.CustomPipelineManager requires a spider argument",
),
pytest.warns(
ScrapyDeprecationWarning,
match=r"The close_spider\(\) method of .+\.CustomPipelineManager requires a spider argument",
),
pytest.warns(
ScrapyDeprecationWarning,
match=r"The process_item\(\) method of .+\.CustomPipelineManager requires a spider argument",
),
):
await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver))

View File

@ -7,10 +7,10 @@ from typing import TYPE_CHECKING, Any
import pytest
from twisted.internet.defer import Deferred, inlineCallbacks, succeed
from twisted.python.failure import Failure
from scrapy.utils.asyncgen import as_async_generator, collect_asyncgen
from scrapy.utils.defer import (
_process_chain,
aiter_errback,
deferred_f_from_coro_f,
deferred_from_coro,
@ -19,7 +19,6 @@ from scrapy.utils.defer import (
maybe_deferred_to_future,
mustbe_deferred,
parallel_async,
process_chain,
process_parallel,
)
@ -79,7 +78,7 @@ def cb3(value, arg1, arg2):
def cb_fail(value, arg1, arg2):
return Failure(TypeError())
raise TypeError
def eb1(failure, arg1, arg2):
@ -87,13 +86,13 @@ def eb1(failure, arg1, arg2):
class TestDeferUtils:
@inlineCallbacks
def test_process_chain(self):
x = yield process_chain([cb1, cb2, cb3], "res", "v1", "v2")
@deferred_f_from_coro_f
async def test_process_chain(self):
x = await _process_chain([cb1, cb2, cb3], "res", "v1", "v2")
assert x == "(cb3 (cb2 (cb1 res v1 v2) v1 v2) v1 v2)"
with pytest.raises(TypeError):
yield process_chain([cb1, cb_fail, cb3], "res", "v1", "v2")
await _process_chain([cb1, cb_fail, cb3], "res", "v1", "v2")
@inlineCallbacks
def test_process_parallel(self):