Implement lazy support

This commit is contained in:
Adrián Chaves 2025-03-26 19:40:53 +01:00
parent 9726538cec
commit 313f9de28d
12 changed files with 236 additions and 380 deletions

View File

@ -22,7 +22,8 @@ Backward-incompatible changes
As a result, the order in which start requests are sent may change. See
:ref:`start-requests` for details and information on how to force start
request order.
request order or pause start request iteration while there are scheduled
requests.
- In ``scrapy.core.engine.ExecutionEngine``:

View File

@ -373,16 +373,22 @@ Scrapy does not try to send :meth:`~scrapy.Spider.start` requests in order.
Instead, it prioritizes reaching :setting:`CONCURRENT_REQUESTS` and
:ref:`scheduling <topics-scheduler>` start requests.
To change that, override the :meth:`~scrapy.Spider.start` method to set
:attr:`Request.priority <scrapy.http.Request.priority>`. For example:
Forcing a start request order
-----------------------------
To force a specific **request order**, override the
:meth:`~scrapy.Spider.start` method to set :attr:`Request.priority
<scrapy.http.Request.priority>`. For example:
- To send start requests before other requests:
.. code-block:: python
async def start(self):
async for request in super().start():
yield request.replace(priority=1)
async for item_or_request in super().start():
if isinstance(item_or_request, Request):
item_or_request = item_or_request.replace(priority=1)
yield item_or_request
- To send start requests in order:
@ -390,13 +396,33 @@ To change that, override the :meth:`~scrapy.Spider.start` method to set
async def start(self):
priority = len(self.start_urls)
async for request in super().start():
yield request.replace(priority=priority)
async for item_or_request in super().start():
if isinstance(item_or_request, Request):
item_or_request = item_or_request.replace(priority=priority)
yield item_or_request
priority -= 1
You can also :ref:`customize the scheduler <topics-scheduler>` if you need
more control over request prioritization.
Delaying start request iteration
--------------------------------
You can override the :meth:`~scrapy.Spider.start` method as follows to pause
its iteration whenever there are scheduled requests:
.. code-block:: python
async def start(self):
async for item_or_request in super().start():
if self.crawler.engine.needs_backoff():
await self.crawler.signals.wait_for(signals.scheduler_empty)
yield item_or_request
This can help minimize the number of requests in the scheduler at any given
time, to minimize resource usage (memory or disk, depending on
:setting:`JOBDIR`).
.. _builtin-spiders:
Generic Spiders

View File

@ -85,7 +85,7 @@ class Command(ScrapyCommand):
crawler._apply_settings()
# The Shell class needs a persistent engine in the crawler
crawler.engine = crawler._create_engine()
crawler.engine.start()
crawler.engine.start(_start_request_processing=False)
self._start_crawler_thread()

View File

@ -12,7 +12,7 @@ from time import time
from traceback import format_exc
from typing import TYPE_CHECKING, Any, TypeVar, cast
from twisted.internet.defer import Deferred, inlineCallbacks, succeed
from twisted.internet.defer import Deferred, succeed
from twisted.internet.task import LoopingCall
from twisted.python.failure import Failure
@ -20,13 +20,16 @@ from scrapy import signals
from scrapy.core.scraper import Scraper, _HandleOutputDeferred
from scrapy.exceptions import CloseSpider, DontCloseSpider, IgnoreRequest
from scrapy.http import Request, Response
from scrapy.utils.defer import deferred_from_coro
from scrapy.utils.defer import (
deferred_f_from_coro_f,
maybe_deferred_to_future,
)
from scrapy.utils.log import failure_to_exc_info, logformatter_adapter
from scrapy.utils.misc import build_from_crawler, load_object
from scrapy.utils.reactor import CallLaterOnce
if TYPE_CHECKING:
from collections.abc import AsyncIterable, Callable, Generator
from collections.abc import AsyncIterable, Callable
from scrapy.core.downloader import Downloader
from scrapy.core.scheduler import BaseScheduler
@ -99,6 +102,7 @@ class ExecutionEngine:
)
self.start_time: float | None = None
self._start: AsyncIterable[Any] | None = None
self._started_request_processing = False
downloader_cls: type[Downloader] = load_object(self.settings["DOWNLOADER"])
try:
self.scheduler_cls: type[BaseScheduler] = self._get_scheduler_class(
@ -121,22 +125,28 @@ class ExecutionEngine:
)
return scheduler_cls
@inlineCallbacks
def start(self) -> Generator[Deferred[Any], Any, None]:
@deferred_f_from_coro_f
async def start(self, _start_request_processing=True) -> None:
if self.running:
raise RuntimeError("Engine already running")
self.start_time = time()
yield self.signals.send_catch_log_deferred(signal=signals.engine_started)
await maybe_deferred_to_future(
self.signals.send_catch_log_deferred(signal=signals.engine_started)
)
self.running = True
if _start_request_processing:
self.start_request_processing()
self._closewait: Deferred[None] = Deferred()
yield self._closewait
await maybe_deferred_to_future(self._closewait)
def stop(self) -> Deferred[None]:
"""Gracefully stop the execution engine"""
@inlineCallbacks
def _finish_stopping_engine(_: Any) -> Generator[Deferred[Any], Any, None]:
yield self.signals.send_catch_log_deferred(signal=signals.engine_stopped)
@deferred_f_from_coro_f
async def _finish_stopping_engine(_: Any) -> None:
await maybe_deferred_to_future(
self.signals.send_catch_log_deferred(signal=signals.engine_stopped)
)
self._closewait.callback(None)
if not self.running:
@ -170,10 +180,9 @@ class ExecutionEngine:
def unpause(self) -> None:
self.paused = False
@inlineCallbacks
def _process_next_spider_start_yield(self):
async def _process_next_spider_start_yield(self):
try:
item_or_request = yield deferred_from_coro(self._start.__anext__())
item_or_request = await self._start.__anext__()
except StopAsyncIteration:
self._start = None
except Exception as exception:
@ -190,30 +199,37 @@ class ExecutionEngine:
self.scraper.start_itemproc(item_or_request, response=None)
self._slot.nextcall.schedule()
@inlineCallbacks
def _process_spider_start(self) -> Generator[Deferred[Any], Any, None]:
@deferred_f_from_coro_f
async def start_request_processing(self) -> None:
"""Process start items and requests in an asynchronous loop.
Items are scraped. Requests are scheduled.
"""
if self._started_request_processing:
raise RuntimeError("Request processing already started")
self._started_request_processing = True
assert self._slot is not None # typing
self._slot.nextcall.schedule()
self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL)
while self._start is not None:
yield self._process_next_spider_start_yield()
if not self._needs_backout():
await self._process_next_spider_start_yield()
if not self.needs_backout():
self._slot.nextcall.schedule()
await self._slot.nextcall.wait()
def _start_next_requests(self) -> None:
if self._slot is None or self._slot.closing is not None or self.paused:
return
while not self._needs_backout():
while not self.needs_backout():
if self._start_scheduled_request() is None:
break
if self.spider_is_idle() and self._slot.close_if_idle:
self._spider_idle()
def _needs_backout(self) -> bool:
def needs_backout(self) -> bool:
assert self._slot is not None # typing
assert self.scraper.slot is not None # typing
return (
@ -377,29 +393,30 @@ class ExecutionEngine:
dwld.addBoth(_on_complete)
return dwld
@inlineCallbacks
def open_spider(
@deferred_f_from_coro_f
async def open_spider(
self,
spider: Spider,
close_if_idle: bool = True,
) -> Generator[Deferred[Any], Any, None]:
) -> None:
if self._slot is not None:
raise RuntimeError(f"No free spider slot when opening {spider.name!r}")
logger.info("Spider opened", extra={"spider": spider})
self.spider = spider
nextcall = CallLaterOnce(self._start_next_requests)
scheduler = build_from_crawler(self.scheduler_cls, self.crawler)
self._start = yield self.scraper.spidermw.process_start(spider)
self._slot = _Slot(close_if_idle, nextcall, scheduler)
self.spider = spider
self._start = await maybe_deferred_to_future(
self.scraper.spidermw.process_start(spider)
)
if hasattr(scheduler, "open") and (d := scheduler.open(spider)):
yield d
yield self.scraper.open_spider(spider)
await maybe_deferred_to_future(d)
await maybe_deferred_to_future(self.scraper.open_spider(spider))
assert self.crawler.stats
self.crawler.stats.open_spider(spider)
yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider)
self._process_spider_start()
self._slot.nextcall.schedule()
self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL)
await maybe_deferred_to_future(
self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider)
)
def _spider_idle(self) -> None:
"""

View File

@ -136,6 +136,9 @@ class Crawler:
"Overridden settings:\n%(settings)s", {"settings": pprint.pformat(d)}
)
# Cannot use @deferred_f_from_coro_f because that relies on the reactor
# being installed already, which is done within _apply_settings(), inside
# this method.
@inlineCallbacks
def crawl(self, *args: Any, **kwargs: Any) -> Generator[Deferred[Any], Any, None]:
if self.crawling:
@ -152,7 +155,7 @@ class Crawler:
self._update_root_log_handler()
self.engine = self._create_engine()
yield self.engine.open_spider(self.spider)
yield maybeDeferred(self.engine.start)
yield self.engine.start()
except Exception:
self.crawling = False
if self.engine is not None:

View File

@ -24,6 +24,7 @@ from scrapy.spiders import Spider
from scrapy.utils.conf import get_config
from scrapy.utils.console import DEFAULT_PYTHON_SHELLS, start_python_console
from scrapy.utils.datatypes import SequenceExclude
from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future
from scrapy.utils.misc import load_object
from scrapy.utils.reactor import is_asyncio_reactor_installed, set_asyncio_event_loop
from scrapy.utils.response import open_in_browser
@ -102,25 +103,33 @@ class Shell:
# set the asyncio event loop for the current thread
event_loop_path = self.crawler.settings["ASYNCIO_EVENT_LOOP"]
set_asyncio_event_loop(event_loop_path)
spider = self._open_spider(request, spider)
def crawl_request(_):
assert self.crawler.engine is not None
self.crawler.engine.crawl(request)
d2 = self._open_spider(request, spider)
d2.addCallback(crawl_request)
d = _request_deferred(request)
d.addCallback(lambda x: (x, spider))
assert self.crawler.engine
self.crawler.engine.crawl(request)
return d
def _open_spider(self, request: Request, spider: Spider | None) -> Spider:
@deferred_f_from_coro_f
async def _open_spider(self, request: Request, spider: Spider | None) -> None:
if self.spider:
return self.spider
return
if spider is None:
spider = self.crawler.spider or self.crawler._create_spider()
self.crawler.spider = spider
assert self.crawler.engine
self.crawler.engine.open_spider(spider, close_if_idle=False)
await maybe_deferred_to_future(
self.crawler.engine.open_spider(spider, close_if_idle=False)
)
self.crawler.engine.start_request_processing()
self.spider = spider
return spider
def fetch(
self,

View File

@ -1,13 +1,12 @@
from __future__ import annotations
from typing import TYPE_CHECKING, Any
from typing import Any
from pydispatch import dispatcher
from twisted.internet.defer import Deferred
from scrapy.utils import signal as _signal
if TYPE_CHECKING:
from twisted.internet.defer import Deferred
from scrapy.utils.defer import maybe_deferred_to_future
class SignalManager:
@ -75,3 +74,14 @@ class SignalManager:
"""
kwargs.setdefault("sender", self.sender)
_signal.disconnect_all(signal, **kwargs)
async def wait_for(self, signal):
"""Await the next *signal*."""
d = Deferred()
def handle():
self.disconnect(handle, signal)
d.callback(None)
self.connect(handle, signal)
await maybe_deferred_to_future(d)

View File

@ -7,6 +7,7 @@ from typing import TYPE_CHECKING, Any, Generic, TypeVar
from warnings import catch_warnings, filterwarnings
from twisted.internet import asyncioreactor, error
from twisted.internet.defer import Deferred
from scrapy.utils.misc import load_object
@ -54,6 +55,7 @@ class CallLaterOnce(Generic[_T]):
self._a: tuple[Any, ...] = a
self._kw: dict[str, Any] = kw
self._call: DelayedCall | None = None
self._deferreds = []
def schedule(self, delay: float = 0) -> None:
from twisted.internet import reactor
@ -66,8 +68,23 @@ class CallLaterOnce(Generic[_T]):
self._call.cancel()
def __call__(self) -> _T:
from twisted.internet import reactor
self._call = None
return self._func(*self._a, **self._kw)
result = self._func(*self._a, **self._kw)
for d in self._deferreds:
reactor.callLater(0, d.callback, None)
self._deferreds = []
return result
async def wait(self):
from scrapy.utils.defer import maybe_deferred_to_future
d = Deferred()
self._deferreds.append(d)
await maybe_deferred_to_future(d)
def set_asyncio_event_loop_policy() -> None:

View File

@ -8,11 +8,16 @@ class TestCmdlineCrawlPipeline:
args = (sys.executable, "-m", "scrapy.cmdline", "crawl", spname)
cwd = Path(__file__).resolve().parent
proc = Popen(args, stdout=PIPE, stderr=PIPE, cwd=cwd)
proc.communicate()
return proc.returncode
_, stderr = proc.communicate()
return proc.returncode, stderr
def test_open_spider_normally_in_pipeline(self):
assert self._execute("normal") == 0
returncode, stderr = self._execute("normal")
assert returncode == 0
def test_exception_at_open_spider_in_pipeline(self):
assert self._execute("exception") == 1
returncode, stderr = self._execute("exception")
assert (
returncode == 0
) # An unhandled exception in a pipeline should not stop the crawl
assert b'RuntimeError("exception")' in stderr

View File

@ -12,6 +12,11 @@ from scrapy.core.downloader.middleware import DownloaderMiddlewareManager
from scrapy.exceptions import _InvalidOutput
from scrapy.http import Request, Response
from scrapy.spiders import Spider
from scrapy.utils.defer import (
deferred_f_from_coro_f,
deferred_to_future,
maybe_deferred_to_future,
)
from scrapy.utils.python import to_bytes
from scrapy.utils.test import get_crawler, get_from_asyncio_queue
@ -29,7 +34,7 @@ class TestManagerBase(TestCase):
def tearDown(self):
return self.crawler.engine.close_spider(self.spider)
def _download(self, request, response=None):
async def _download(self, request, response=None):
"""Executes downloader mw manager's download method and returns
the result (Request or Response) or raise exception in case of
failure.
@ -44,7 +49,7 @@ class TestManagerBase(TestCase):
# catch deferred result and return the value
results = []
dfd.addBoth(results.append)
self._wait(dfd)
await maybe_deferred_to_future(dfd)
ret = results[0]
if isinstance(ret, Failure):
ret.raiseException()
@ -54,13 +59,15 @@ class TestManagerBase(TestCase):
class TestDefaults(TestManagerBase):
"""Tests default behavior with default settings"""
def test_request_response(self):
@deferred_f_from_coro_f
async def test_request_response(self):
req = Request("http://example.com/index.html")
resp = Response(req.url, status=200)
ret = self._download(req, resp)
ret = await self._download(req, resp)
assert isinstance(ret, Response), "Non-response returned"
def test_3xx_and_invalid_gzipped_body_must_redirect(self):
@deferred_f_from_coro_f
async def test_3xx_and_invalid_gzipped_body_must_redirect(self):
"""Regression test for a failure when redirecting a compressed
request.
@ -85,13 +92,14 @@ class TestDefaults(TestManagerBase):
"Location": "http://example.com/login",
},
)
ret = self._download(request=req, response=resp)
ret = await self._download(request=req, response=resp)
assert isinstance(ret, Request), f"Not redirected: {ret!r}"
assert to_bytes(ret.url) == resp.headers["Location"], (
"Not redirected to location header"
)
def test_200_and_invalid_gzipped_body_must_fail(self):
@deferred_f_from_coro_f
async def test_200_and_invalid_gzipped_body_must_fail(self):
req = Request("http://example.com")
body = b"<p>You are being redirected</p>"
resp = Response(
@ -106,13 +114,14 @@ class TestDefaults(TestManagerBase):
},
)
with pytest.raises(BadGzipFile):
self._download(request=req, response=resp)
await self._download(request=req, response=resp)
class TestResponseFromProcessRequest(TestManagerBase):
"""Tests middleware returning a response from process_request."""
def test_download_func_not_called(self):
@deferred_f_from_coro_f
async def test_download_func_not_called(self):
resp = Response("http://example.com/index.html")
class ResponseMiddleware:
@ -126,7 +135,7 @@ class TestResponseFromProcessRequest(TestManagerBase):
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
self._wait(dfd)
await maybe_deferred_to_future(dfd)
assert results[0] is resp
assert not download_func.called
@ -195,7 +204,8 @@ class TestProcessExceptionInvalidOutput(TestManagerBase):
class TestMiddlewareUsingDeferreds(TestManagerBase):
"""Middlewares using Deferreds should work"""
def test_deferred(self):
@deferred_f_from_coro_f
async def test_deferred(self):
resp = Response("http://example.com/index.html")
class DeferredMiddleware:
@ -214,7 +224,7 @@ class TestMiddlewareUsingDeferreds(TestManagerBase):
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
self._wait(dfd)
await maybe_deferred_to_future(dfd)
assert results[0] is resp
assert not download_func.called
@ -224,7 +234,8 @@ class TestMiddlewareUsingDeferreds(TestManagerBase):
class TestMiddlewareUsingCoro(TestManagerBase):
"""Middlewares using asyncio coroutines should work"""
def test_asyncdef(self):
@deferred_f_from_coro_f
async def test_asyncdef(self):
resp = Response("http://example.com/index.html")
class CoroMiddleware:
@ -238,13 +249,14 @@ class TestMiddlewareUsingCoro(TestManagerBase):
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
self._wait(dfd)
await maybe_deferred_to_future(dfd)
assert results[0] is resp
assert not download_func.called
@pytest.mark.only_asyncio
def test_asyncdef_asyncio(self):
@deferred_f_from_coro_f
async def test_asyncdef_asyncio(self):
resp = Response("http://example.com/index.html")
class CoroMiddleware:
@ -258,7 +270,7 @@ class TestMiddlewareUsingCoro(TestManagerBase):
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
self._wait(dfd)
await deferred_to_future(dfd)
assert results[0] is resp
assert not download_func.called

View File

@ -150,7 +150,6 @@ class CrawlerRun:
"""A class to run the crawler and keep track of events occurred"""
def __init__(self, spider_class):
self.spider = None
self.respplug = []
self.reqplug = []
self.reqdropped = []
@ -191,7 +190,6 @@ class CrawlerRun:
self.response_downloaded, signals.response_downloaded
)
self.crawler.crawl(start_urls=start_urls)
self.spider = self.crawler.spider
self.deferred = defer.Deferred()
dispatcher.connect(self.stop, signals.engine_stopped)
@ -297,7 +295,7 @@ class TestEngineBase(unittest.TestCase):
assert len(run.itemerror) == 2
for item, response, spider, failure in run.itemerror:
assert failure.value.__class__ is ZeroDivisionError
assert spider == run.spider
assert spider == run.crawler.spider
assert item["url"] == response.url
if "item1.html" in item["url"]:
@ -378,11 +376,14 @@ class TestEngineBase(unittest.TestCase):
assert signals.spider_closed in run.signals_caught
assert signals.headers_received in run.signals_caught
assert {"spider": run.spider} == run.signals_caught[signals.spider_opened]
assert {"spider": run.spider} == run.signals_caught[signals.spider_idle]
assert {"spider": run.spider, "reason": "finished"} == run.signals_caught[
signals.spider_closed
assert {"spider": run.crawler.spider} == run.signals_caught[
signals.spider_opened
]
assert {"spider": run.crawler.spider} == run.signals_caught[signals.spider_idle]
assert {
"spider": run.crawler.spider,
"reason": "finished",
} == run.signals_caught[signals.spider_closed]
class TestEngine(TestEngineBase):
@ -420,9 +421,10 @@ class TestEngine(TestEngineBase):
def test_crawler_change_close_reason_on_idle(self):
run = CrawlerRun(ChangeCloseReasonSpider)
yield run.run()
assert {"spider": run.spider, "reason": "custom_reason"} == run.signals_caught[
signals.spider_closed
]
assert {
"spider": run.crawler.spider,
"reason": "custom_reason",
} == run.signals_caught[signals.spider_closed]
@defer.inlineCallbacks
def test_close_downloader(self):
@ -472,7 +474,7 @@ class TestEngine(TestEngineBase):
finally:
timer.cancel()
assert b"Traceback" not in stderr
assert b"Traceback" not in stderr, stderr
def test_request_scheduled_signal(caplog):

View File

@ -1,6 +1,5 @@
from collections import deque
import pytest
from twisted.internet.defer import Deferred
from twisted.trial.unittest import TestCase
@ -13,12 +12,12 @@ from .mockserver import MockServer
from .test_scheduler import MemoryScheduler
def sleep(seconds: float = 0.001):
async def sleep(seconds: float = 0.001) -> None:
from twisted.internet import reactor
deferred: Deferred[None] = Deferred()
reactor.callLater(seconds, deferred.callback, None)
return maybe_deferred_to_future(deferred)
await maybe_deferred_to_future(deferred)
class MainTestCase(TestCase):
@ -87,9 +86,7 @@ class RequestSendOrderTestCase(TestCase):
callback requests have the same priority.
It is a very unintuitive behavior, documented as undefined so that we may
change it in the future without breaking the contract.
For the asyncio reactor:
change it in the future without breaking the contract:
1. First, the first CONCURRENT_REQUESTS start requests are sent in order.
@ -105,16 +102,16 @@ class RequestSendOrderTestCase(TestCase):
but only when there are not enough pending requests yielded from
callbacks to reach the configured concurrency.
For the default Twisted reactor, step 1 sends the last CONCURRENT_REQUESTS
start requests in reverse order instead.
The reverse order is because the scheduler uses a LIFO queue by default
(SCHEDULER_MEMORY_QUEUE, SCHEDULER_DISK_QUEUE). The order of the first few
requests is unnaffected because they are sent as soon as they are
scheduled, and the last start requests sent before callback requests are
those that can be sent before the first callback requests are scheduled.
scheduled. The last start requests sent before callback requests are those
that can be sent before the first callback requests are scheduled.
"""
# Error out if any tests relies on the heartbeat.
timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL
@classmethod
def setUpClass(cls):
cls.mockserver = MockServer()
@ -177,11 +174,8 @@ class RequestSendOrderTestCase(TestCase):
expected_nums = sorted(start_nums + cb_nums)
assert actual_nums == expected_nums, f"{actual_nums=} != {expected_nums=}"
# Asyncio reactor behavior
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_default(self):
async def test_default(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[
@ -215,9 +209,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_conc1(self):
async def test_conc1(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 4, 2],
@ -226,9 +219,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_conc2(self):
async def test_conc2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 6, 4, 3],
@ -237,9 +229,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_conc8(self):
async def test_conc8(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 4, 5, 6, 7, 8, 18, 16, 15, 14, 13, 12, 11, 10, 9],
@ -248,9 +239,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_conc16(self):
async def test_conc16(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[
@ -293,9 +283,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_conc3_ds2(self):
async def test_conc3_ds2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 8, 6, 5, 4],
@ -307,9 +296,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_tconc3_dconc2(self):
async def test_tconc3_dconc2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 7, 5, 4],
@ -321,9 +309,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_tconc5_dconc3(self):
async def test_tconc5_dconc3(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 4, 5, 10, 8, 7, 6],
@ -335,9 +322,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_tconc5_dconc2_ds3(self):
async def test_tconc5_dconc2_ds3(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 4, 5, 12, 10, 9, 8, 7, 6],
@ -350,9 +336,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_tconc5_dconc3_ds2(self):
async def test_tconc5_dconc3_ds2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 4, 5, 12, 10, 9, 8, 7, 6],
@ -365,9 +350,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_tconc7_dconc2_ds3(self):
async def test_tconc7_dconc2_ds3(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 4, 5, 6, 7, 15, 13, 12, 11, 10, 9, 8],
@ -380,9 +364,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_tconc7_dconc3_ds2(self):
async def test_tconc7_dconc3_ds2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[1, 2, 3, 4, 5, 6, 7, 15, 13, 12, 11, 10, 9, 8],
@ -395,9 +378,8 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_fast(self):
async def test_fast(self):
"""Very fast responses may increase the number of start requests sent
in reverse order before the first callback request."""
await maybe_deferred_to_future(
@ -409,240 +391,6 @@ class RequestSendOrderTestCase(TestCase):
)
)
# Default reactor behavior
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_default(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[
26,
24,
23,
22,
21,
20,
19,
18,
17,
16,
15,
14,
13,
12,
11,
10,
9,
8,
7,
6,
5,
4,
3,
2,
1,
],
cb_nums=[25],
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_conc1(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[4, 2, 1],
cb_nums=[3],
settings={"CONCURRENT_REQUESTS": 1},
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_conc2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[6, 4, 3, 2, 1],
cb_nums=[5],
settings={"CONCURRENT_REQUESTS": 2},
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_conc8(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[18, 16, 15, 14, 13, 12, 11, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1],
cb_nums=[17],
settings={"CONCURRENT_REQUESTS": 8},
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_conc16(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[
34,
32,
31,
30,
29,
28,
27,
26,
25,
24,
23,
22,
21,
20,
19,
18,
17,
16,
15,
14,
13,
12,
11,
10,
9,
8,
7,
6,
5,
4,
3,
2,
1,
],
cb_nums=[33],
settings={"CONCURRENT_REQUESTS_PER_DOMAIN": 16},
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_conc3_ds2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[8, 6, 5, 4, 3, 2, 1],
cb_nums=[7],
settings={
"CONCURRENT_REQUESTS": 3,
},
download_slots=2,
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_tconc3_dconc2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[7, 5, 4, 3, 2, 1],
cb_nums=[6],
settings={
"CONCURRENT_REQUESTS": 3,
"CONCURRENT_REQUESTS_PER_DOMAIN": 2,
},
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_tconc5_dconc3(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[10, 8, 7, 6, 5, 4, 3, 2, 1],
cb_nums=[9],
settings={
"CONCURRENT_REQUESTS": 5,
"CONCURRENT_REQUESTS_PER_DOMAIN": 3,
},
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_tconc5_dconc2_ds3(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[12, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1],
cb_nums=[11],
settings={
"CONCURRENT_REQUESTS": 5,
"CONCURRENT_REQUESTS_PER_DOMAIN": 2,
},
download_slots=3,
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_tconc5_dconc3_ds2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[12, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1],
cb_nums=[11],
settings={
"CONCURRENT_REQUESTS": 5,
"CONCURRENT_REQUESTS_PER_DOMAIN": 3,
},
download_slots=2,
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_tconc7_dconc2_ds3(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[15, 13, 12, 11, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1],
cb_nums=[14],
settings={
"CONCURRENT_REQUESTS": 7,
"CONCURRENT_REQUESTS_PER_DOMAIN": 2,
},
download_slots=3,
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_tconc7_dconc3_ds2(self):
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[15, 13, 12, 11, 10, 9, 8, 7, 6, 5, 4, 3, 2, 1],
cb_nums=[14],
settings={
"CONCURRENT_REQUESTS": 7,
"CONCURRENT_REQUESTS_PER_DOMAIN": 3,
},
download_slots=2,
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_fast(self):
"""Very fast responses may increase the number of start requests sent
in reverse order before the first callback request."""
await maybe_deferred_to_future(
self._test_request_order(
start_nums=[3, 2, 1],
cb_nums=[4],
settings={"CONCURRENT_REQUESTS": 1},
response_seconds=self.fast_seconds,
)
)
# Behavior shared by both reactors
@deferred_f_from_coro_f
async def test_await(self):
"""Awaiting slow operations in Spider.start() may lower the number of
@ -672,9 +420,8 @@ class RequestSendOrderTestCase(TestCase):
# Examples from the “Start requests” section of the documentation about
# spiders.
@pytest.mark.only_asyncio
@deferred_f_from_coro_f
async def test_ar_start_requests_first(self):
async def test_start_requests_first(self):
start_nums = [1, 3, 2]
cb_nums = [4]
response_seconds = self.slow_seconds
@ -695,29 +442,6 @@ class RequestSendOrderTestCase(TestCase):
)
)
@pytest.mark.only_not_asyncio
@deferred_f_from_coro_f
async def test_dr_start_requests_first(self):
start_nums = [3, 2, 1]
cb_nums = [4]
response_seconds = self.slow_seconds
download_slots = 1
async def start(spider):
for num in start_nums:
request = self._request(num, response_seconds, download_slots)
yield request.replace(priority=1)
await maybe_deferred_to_future(
self._test_request_order(
start_nums=start_nums,
cb_nums=cb_nums,
settings={"CONCURRENT_REQUESTS": 1},
response_seconds=response_seconds,
start_fn=start,
)
)
@deferred_f_from_coro_f
async def test_start_requests_first_sorted(self):
start_nums = [1, 2, 3]
@ -741,3 +465,33 @@ class RequestSendOrderTestCase(TestCase):
start_fn=start,
)
)
@deferred_f_from_coro_f
async def test_lazy(self):
start_nums = [1, 2, 4]
cb_nums = [3]
response_seconds = self.slow_seconds
download_slots = 1
async def start(spider):
for num in start_nums:
if spider.crawler.engine.needs_backout():
await spider.crawler.signals.wait_for(signals.scheduler_empty)
request = self._request(num, response_seconds, download_slots)
yield request
await maybe_deferred_to_future(
self._test_request_order(
start_nums=start_nums,
cb_nums=cb_nums,
settings={
"CONCURRENT_REQUESTS": 1,
# Without the lazy approach, using the FIFO queue would
# yield a different result, with start requests not being
# sorted.
"SCHEDULER_MEMORY_QUEUE": "scrapy.squeues.FifoMemoryQueue",
},
response_seconds=response_seconds,
start_fn=start,
)
)