mirror of https://github.com/scrapy/scrapy.git
Add the call_later() wrapper. (#6858)
This commit is contained in:
parent
c6698b9fe8
commit
5902aab25c
|
|
@ -14,7 +14,12 @@ from scrapy.core.downloader.handlers import DownloadHandlers
|
|||
from scrapy.core.downloader.middleware import DownloaderMiddlewareManager
|
||||
from scrapy.exceptions import ScrapyDeprecationWarning
|
||||
from scrapy.resolver import dnscache
|
||||
from scrapy.utils.asyncio import AsyncioLoopingCall, create_looping_call
|
||||
from scrapy.utils.asyncio import (
|
||||
AsyncioLoopingCall,
|
||||
CallLaterResult,
|
||||
call_later,
|
||||
create_looping_call,
|
||||
)
|
||||
from scrapy.utils.defer import (
|
||||
deferred_from_coro,
|
||||
maybe_deferred_to_future,
|
||||
|
|
@ -50,7 +55,7 @@ class Slot:
|
|||
self.queue: deque[tuple[Request, Deferred[Response]]] = deque()
|
||||
self.transferring: set[Request] = set()
|
||||
self.lastseen: float = 0
|
||||
self.latercall = None
|
||||
self.latercall: CallLaterResult | None = None
|
||||
|
||||
def free_transfer_slots(self) -> int:
|
||||
return self.concurrency - len(self.transferring)
|
||||
|
|
@ -61,8 +66,9 @@ class Slot:
|
|||
return self.delay
|
||||
|
||||
def close(self) -> None:
|
||||
if self.latercall and self.latercall.active():
|
||||
if self.latercall:
|
||||
self.latercall.cancel()
|
||||
self.latercall = None
|
||||
|
||||
def __repr__(self) -> str:
|
||||
cls_name = self.__class__.__name__
|
||||
|
|
@ -191,9 +197,8 @@ class Downloader:
|
|||
slot.active.remove(request)
|
||||
|
||||
def _process_queue(self, spider: Spider, slot: Slot) -> None:
|
||||
from twisted.internet import reactor
|
||||
|
||||
if slot.latercall and slot.latercall.active():
|
||||
if slot.latercall:
|
||||
# block processing until slot.latercall is called
|
||||
return
|
||||
|
||||
# Delay queue processing if a download_delay is configured
|
||||
|
|
@ -202,9 +207,7 @@ class Downloader:
|
|||
if delay:
|
||||
penalty = delay - now + slot.lastseen
|
||||
if penalty > 0:
|
||||
slot.latercall = reactor.callLater(
|
||||
penalty, self._process_queue, spider, slot
|
||||
)
|
||||
slot.latercall = call_later(penalty, self._latercall, spider, slot)
|
||||
return
|
||||
|
||||
# Process enqueued requests if there are free slots to transfer for this slot
|
||||
|
|
@ -218,6 +221,10 @@ class Downloader:
|
|||
self._process_queue(spider, slot)
|
||||
break
|
||||
|
||||
def _latercall(self, spider: Spider, slot: Slot) -> None:
|
||||
slot.latercall = None
|
||||
self._process_queue(spider, slot)
|
||||
|
||||
async def _download(self, slot: Slot, request: Request, spider: Spider) -> Response:
|
||||
# The order is very important for the following logic. Do not change!
|
||||
slot.transferring.add(request)
|
||||
|
|
|
|||
|
|
@ -12,9 +12,15 @@ from typing import TYPE_CHECKING, Any
|
|||
|
||||
from scrapy import Request, Spider, signals
|
||||
from scrapy.exceptions import NotConfigured
|
||||
from scrapy.utils.asyncio import create_looping_call
|
||||
from scrapy.utils.asyncio import (
|
||||
AsyncioLoopingCall,
|
||||
CallLaterResult,
|
||||
call_later,
|
||||
create_looping_call,
|
||||
)
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from twisted.internet.task import LoopingCall
|
||||
from twisted.python.failure import Failure
|
||||
|
||||
# typing.Self requires Python 3.11
|
||||
|
|
@ -31,6 +37,12 @@ class CloseSpider:
|
|||
def __init__(self, crawler: Crawler):
|
||||
self.crawler: Crawler = crawler
|
||||
|
||||
# for CLOSESPIDER_TIMEOUT
|
||||
self.task: CallLaterResult | None = None
|
||||
|
||||
# for CLOSESPIDER_TIMEOUT_NO_ITEM
|
||||
self.task_no_item: AsyncioLoopingCall | LoopingCall | None = None
|
||||
|
||||
self.close_on: dict[str, Any] = {
|
||||
"timeout": crawler.settings.getfloat("CLOSESPIDER_TIMEOUT"),
|
||||
"itemcount": crawler.settings.getint("CLOSESPIDER_ITEMCOUNT"),
|
||||
|
|
@ -92,14 +104,12 @@ class CloseSpider:
|
|||
self.crawler.engine.close_spider(spider, "closespider_pagecount_no_item")
|
||||
|
||||
def spider_opened(self, spider: Spider) -> None:
|
||||
from twisted.internet import reactor
|
||||
|
||||
assert self.crawler.engine
|
||||
self.task = reactor.callLater(
|
||||
self.task = call_later(
|
||||
self.close_on["timeout"],
|
||||
self.crawler.engine.close_spider,
|
||||
spider,
|
||||
reason="closespider_timeout",
|
||||
"closespider_timeout",
|
||||
)
|
||||
|
||||
def item_scraped(self, item: Any, spider: Spider) -> None:
|
||||
|
|
@ -110,13 +120,14 @@ class CloseSpider:
|
|||
self.crawler.engine.close_spider(spider, "closespider_itemcount")
|
||||
|
||||
def spider_closed(self, spider: Spider) -> None:
|
||||
task = getattr(self, "task", None)
|
||||
if task and task.active():
|
||||
task.cancel()
|
||||
if self.task:
|
||||
self.task.cancel()
|
||||
self.task = None
|
||||
|
||||
task_no_item = getattr(self, "task_no_item", None)
|
||||
if task_no_item and task_no_item.running:
|
||||
task_no_item.stop()
|
||||
if self.task_no_item:
|
||||
if self.task_no_item.running:
|
||||
self.task_no_item.stop()
|
||||
self.task_no_item = None
|
||||
|
||||
def spider_opened_no_item(self, spider: Spider) -> None:
|
||||
self.task_no_item = create_looping_call(self._count_items_produced, spider)
|
||||
|
|
|
|||
|
|
@ -15,10 +15,14 @@ from scrapy.utils.asyncgen import as_async_generator
|
|||
from scrapy.utils.reactor import is_asyncio_reactor_installed, is_reactor_installed
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from twisted.internet.base import DelayedCall
|
||||
|
||||
# typing.Concatenate and typing.ParamSpec require Python 3.10
|
||||
from typing_extensions import Concatenate, ParamSpec
|
||||
# typing.Self, typing.TypeVarTuple and typing.Unpack require Python 3.11
|
||||
from typing_extensions import Concatenate, ParamSpec, Self, TypeVarTuple, Unpack
|
||||
|
||||
_P = ParamSpec("_P")
|
||||
_Ts = TypeVarTuple("_Ts")
|
||||
|
||||
|
||||
_T = TypeVar("_T")
|
||||
|
|
@ -192,3 +196,60 @@ def create_looping_call(
|
|||
if is_asyncio_available():
|
||||
return AsyncioLoopingCall(func, *args, **kwargs)
|
||||
return LoopingCall(func, *args, **kwargs)
|
||||
|
||||
|
||||
def call_later(
|
||||
delay: float, func: Callable[[Unpack[_Ts]], object], *args: Unpack[_Ts]
|
||||
) -> CallLaterResult:
|
||||
"""Schedule a function to be called after a delay.
|
||||
|
||||
This uses either ``loop.call_later()`` or ``reactor.callLater()``, depending
|
||||
on whether asyncio support is available.
|
||||
"""
|
||||
if is_asyncio_available():
|
||||
loop = asyncio.get_event_loop()
|
||||
return CallLaterResult.from_asyncio(loop.call_later(delay, func, *args))
|
||||
|
||||
from twisted.internet import reactor
|
||||
|
||||
return CallLaterResult.from_twisted(reactor.callLater(delay, func, *args))
|
||||
|
||||
|
||||
class CallLaterResult:
|
||||
"""An universal result for :func:`call_later`, wrapping either
|
||||
:class:`asyncio.TimerHandle` or :class:`twisted.internet.base.DelayedCall`.
|
||||
|
||||
The provided API is close to the :class:`asyncio.TimerHandle` one: there is
|
||||
no ``active()`` (as there is no such public API in
|
||||
:class:`asyncio.TimerHandle`) but ``cancel()`` can be called on already
|
||||
called or cancelled instances.
|
||||
"""
|
||||
|
||||
_timer_handle: asyncio.TimerHandle | None = None
|
||||
_delayed_call: DelayedCall | None = None
|
||||
|
||||
@classmethod
|
||||
def from_asyncio(cls, timer_handle: asyncio.TimerHandle) -> Self:
|
||||
"""Create a CallLaterResult from an asyncio TimerHandle."""
|
||||
o = cls()
|
||||
o._timer_handle = timer_handle
|
||||
return o
|
||||
|
||||
@classmethod
|
||||
def from_twisted(cls, delayed_call: DelayedCall) -> Self:
|
||||
"""Create a CallLaterResult from a Twisted DelayedCall."""
|
||||
o = cls()
|
||||
o._delayed_call = delayed_call
|
||||
return o
|
||||
|
||||
def cancel(self) -> None:
|
||||
"""Cancel the underlying delayed call.
|
||||
|
||||
Does nothing if the delayed call was already called or cancelled.
|
||||
"""
|
||||
if self._timer_handle:
|
||||
self._timer_handle.cancel()
|
||||
self._timer_handle = None
|
||||
elif self._delayed_call and self._delayed_call.active():
|
||||
self._delayed_call.cancel()
|
||||
self._delayed_call = None
|
||||
|
|
|
|||
|
|
@ -16,13 +16,14 @@ if TYPE_CHECKING:
|
|||
from asyncio import AbstractEventLoop, AbstractEventLoopPolicy
|
||||
from collections.abc import Callable
|
||||
|
||||
from twisted.internet.base import DelayedCall
|
||||
from twisted.internet.protocol import ServerFactory
|
||||
from twisted.internet.tcp import Port
|
||||
|
||||
# typing.ParamSpec requires Python 3.10
|
||||
from typing_extensions import ParamSpec
|
||||
|
||||
from scrapy.utils.asyncio import CallLaterResult
|
||||
|
||||
_P = ParamSpec("_P")
|
||||
|
||||
_T = TypeVar("_T")
|
||||
|
|
@ -55,27 +56,27 @@ class CallLaterOnce(Generic[_T]):
|
|||
self._func: Callable[_P, _T] = func
|
||||
self._a: tuple[Any, ...] = a
|
||||
self._kw: dict[str, Any] = kw
|
||||
self._call: DelayedCall | None = None
|
||||
self._call: CallLaterResult | None = None
|
||||
self._deferreds: list[Deferred] = []
|
||||
|
||||
def schedule(self, delay: float = 0) -> None:
|
||||
from twisted.internet import reactor
|
||||
from scrapy.utils.asyncio import call_later
|
||||
|
||||
if self._call is None:
|
||||
self._call = reactor.callLater(delay, self)
|
||||
self._call = call_later(delay, self)
|
||||
|
||||
def cancel(self) -> None:
|
||||
if self._call:
|
||||
self._call.cancel()
|
||||
|
||||
def __call__(self) -> _T:
|
||||
from twisted.internet import reactor
|
||||
from scrapy.utils.asyncio import call_later
|
||||
|
||||
self._call = None
|
||||
result = self._func(*self._a, **self._kw)
|
||||
|
||||
for d in self._deferreds:
|
||||
reactor.callLater(0, d.callback, None)
|
||||
call_later(0, d.callback, None)
|
||||
self._deferreds = []
|
||||
|
||||
return result
|
||||
|
|
|
|||
Loading…
Reference in New Issue