This commit is contained in:
Adrián Chaves 2025-03-26 20:49:44 +01:00
parent c5518bc9fe
commit d2bb2579f7
3 changed files with 27 additions and 28 deletions

View File

@ -81,14 +81,12 @@ class Command(ScrapyCommand):
spider_loader = self.crawler_process.spider_loader spider_loader = self.crawler_process.spider_loader
async def start(self): async def start(self):
for request in conman.from_spider(self, self._result): for request in conman.from_spider(self, result):
yield request yield request
with set_environ(SCRAPY_CHECK="true"): with set_environ(SCRAPY_CHECK="true"):
for spidername in args or spider_loader.list(): for spidername in args or spider_loader.list():
spidercls = spider_loader.load(spidername) spidercls = spider_loader.load(spidername)
spidercls._result = result # type: ignore[assignment,attr-defined,method-assign,return-value]
spidercls.start = start # type: ignore[assignment,method-assign,return-value] spidercls.start = start # type: ignore[assignment,method-assign,return-value]
tested_methods = conman.tested_methods_from_spidercls(spidercls) tested_methods = conman.tested_methods_from_spidercls(spidercls)

View File

@ -180,7 +180,12 @@ class ExecutionEngine:
def unpause(self) -> None: def unpause(self) -> None:
self.paused = False self.paused = False
async def _process_next_spider_start_yield(self): async def _process_start_next(self):
"""Processes the next item or request from Spider.start().
If a request, it is scheduled. If an item, it is sent to item
pipelines.
"""
try: try:
item_or_request = await self._start.__anext__() item_or_request = await self._start.__anext__()
except StopAsyncIteration: except StopAsyncIteration:
@ -201,24 +206,28 @@ class ExecutionEngine:
@deferred_f_from_coro_f @deferred_f_from_coro_f
async def start_request_processing(self) -> None: async def start_request_processing(self) -> None:
"""Process start items and requests in an asynchronous loop. """Starts consuming Spider.start() output and sending scheduled
requests."""
Items are scraped. Requests are scheduled.
"""
if self._started_request_processing: if self._started_request_processing:
raise RuntimeError("Request processing already started") raise RuntimeError("Request processing already started")
self._started_request_processing = True self._started_request_processing = True
# Starts the processing of scheduled requests, as well as a periodic
# call to that processing method for scenarios where the scheduler
# reports having pending requests but returns none.
assert self._slot is not None # typing assert self._slot is not None # typing
self._slot.nextcall.schedule() self._slot.nextcall.schedule()
self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL) self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL)
while self._start is not None: while self._start is not None:
await self._process_next_spider_start_yield() await self._process_start_next()
if not self.needs_backout(): if not self.needs_backout():
# Give room for the outcome of self._process_start_next() to be
# processed before continuing with the next iteration.
self._slot.nextcall.schedule() self._slot.nextcall.schedule()
await self._slot.nextcall.wait() await self._slot.nextcall.wait()
def _start_next_requests(self) -> None: def _start_scheduled_requests(self) -> None:
if self._slot is None or self._slot.closing is not None or self.paused: if self._slot is None or self._slot.closing is not None or self.paused:
return return
@ -403,7 +412,7 @@ class ExecutionEngine:
raise RuntimeError(f"No free spider slot when opening {spider.name!r}") raise RuntimeError(f"No free spider slot when opening {spider.name!r}")
logger.info("Spider opened", extra={"spider": spider}) logger.info("Spider opened", extra={"spider": spider})
self.spider = spider self.spider = spider
nextcall = CallLaterOnce(self._start_next_requests) nextcall = CallLaterOnce(self._start_scheduled_requests)
scheduler = build_from_crawler(self.scheduler_cls, self.crawler) scheduler = build_from_crawler(self.scheduler_cls, self.crawler)
self._slot = _Slot(close_if_idle, nextcall, scheduler) self._slot = _Slot(close_if_idle, nextcall, scheduler)
self._start = await maybe_deferred_to_future( self._start = await maybe_deferred_to_future(

View File

@ -8,7 +8,7 @@ from __future__ import annotations
import logging import logging
from collections.abc import AsyncIterable, Callable, Iterable from collections.abc import AsyncIterable, Callable, Iterable
from inspect import isasyncgenfunction, iscoroutine, iscoroutinefunction from inspect import isasyncgenfunction, iscoroutine
from itertools import islice from itertools import islice
from typing import TYPE_CHECKING, Any, TypeVar, Union, cast from typing import TYPE_CHECKING, Any, TypeVar, Union, cast
from warnings import warn from warnings import warn
@ -387,20 +387,20 @@ class SpiderMiddlewareManager(MiddlewareManager):
dfd2.addErrback(process_spider_exception) dfd2.addErrback(process_spider_exception)
return dfd2 return dfd2
@inlineCallbacks @deferred_f_from_coro_f
def process_start( async def process_start(self, spider: Spider) -> AsyncIterable[Any] | None:
self, spider: Spider
) -> Generator[Deferred[Any], Any, AsyncIterable[Any] | None]:
self._check_deprecated_start_requests_use(spider) self._check_deprecated_start_requests_use(spider)
if self._use_start_requests: if self._use_start_requests:
sync_start = iter(spider.start_requests()) sync_start = iter(spider.start_requests())
sync_start = yield self._process_chain( sync_start = await maybe_deferred_to_future(
"process_start_requests", sync_start, spider self._process_chain("process_start_requests", sync_start, spider)
) )
start = as_async_generator(sync_start) start = as_async_generator(sync_start)
else: else:
start = yield self._iter_seeds(spider) start = spider.start()
start = yield self._process_chain("process_start", start) start = await maybe_deferred_to_future(
self._process_chain("process_start", start)
)
return start return start
def _check_deprecated_start_requests_use(self, spider: Spider): def _check_deprecated_start_requests_use(self, spider: Spider):
@ -473,14 +473,6 @@ class SpiderMiddlewareManager(MiddlewareManager):
f"https://docs.scrapy.org/en/VERSION/news.html" f"https://docs.scrapy.org/en/VERSION/news.html"
) )
@staticmethod
def _iter_seeds(spider: Spider):
fn = spider.start
if isasyncgenfunction(fn):
return fn().__aiter__()
assert iscoroutinefunction(fn)
return deferred_from_coro(fn())
# This method is only needed until _async compatibility methods are removed. # This method is only needed until _async compatibility methods are removed.
@staticmethod @staticmethod
def _get_async_method_pair( def _get_async_method_pair(