diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 881e39017..ccec3a123 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -18,14 +18,13 @@ from scrapy.utils.defer import deferred_from_coro, _isasyncgen from scrapy.utils.misc import load_object from scrapy.utils.reactor import CallLaterOnce from scrapy.utils.log import logformatter_adapter, failure_to_exc_info -from scrapy.spiders import WaitUntilQueueEmpty logger = logging.getLogger(__name__) class Slot: - def __init__(self, start_requests, close_if_idle, nextcall, scheduler): + def __init__(self, start_requests, close_if_idle, nextcall, scheduler, new_queue_behavior=False): self.closing = False self.inprogress = set() # requests in progress if _isasyncgen(start_requests): @@ -33,6 +32,7 @@ class Slot: else: self.start_requests = iter(start_requests) self.close_if_idle = close_if_idle + self.new_queue_behavior = new_queue_behavior self.nextcall = nextcall self.scheduler = scheduler self.heartbeat = task.LoopingCall(nextcall.schedule) @@ -118,6 +118,27 @@ class ExecutionEngine: """Resume the execution engine""" self.paused = False + def _send_reqs_to_downloader(self, spider): + while not self._needs_backout(spider): + if not self._next_request_from_scheduler(spider): + break + + @defer.inlineCallbacks + def _schedule_next_req(self, spider, slot): + try: + if _isasyncgen(slot.start_requests): + request = yield deferred_from_coro(slot.start_requests.__anext__()) + else: + request = next(slot.start_requests) + except (StopIteration, StopAsyncIteration): + slot.start_requests = None + except Exception: + slot.start_requests = None + logger.error('Error while obtaining start requests', + exc_info=True, extra={'spider': spider}) + else: + self.crawl(request, spider) + @defer.inlineCallbacks def _next_request(self, spider): slot = self.slot @@ -131,26 +152,14 @@ class ExecutionEngine: return self._waiting_for_request = True try: - while not self._needs_backout(spider) and slot.start_requests: - try: - if _isasyncgen(slot.start_requests): - request = yield deferred_from_coro(slot.start_requests.__anext__()) - else: - request = next(slot.start_requests) - if request == WaitUntilQueueEmpty: - break - except (StopIteration, StopAsyncIteration): - slot.start_requests = None - except Exception: - slot.start_requests = None - logger.error('Error while obtaining start requests', - exc_info=True, extra={'spider': spider}) - else: - self.crawl(request, spider) - - while not self._needs_backout(spider): - if not self._next_request_from_scheduler(spider): - break + if slot.new_queue_behavior: + while not self._needs_backout(spider) and slot.start_requests: + yield self._schedule_next_req(spider, slot) + self._send_reqs_to_downloader(spider) + else: + self._send_reqs_to_downloader(spider) + if not self._needs_backout(spider) and slot.start_requests: + yield self._schedule_next_req(spider, slot) if self.spider_is_idle(spider) and slot.close_if_idle: self._spider_idle(spider) @@ -279,14 +288,14 @@ class ExecutionEngine: return dwld @defer.inlineCallbacks - def open_spider(self, spider, start_requests=(), close_if_idle=True): + def open_spider(self, spider, start_requests=(), close_if_idle=True, new_queue_behavior=False): if not self.has_capacity(): raise RuntimeError("No free spider slot when opening %r" % spider.name) logger.info("Spider opened", extra={'spider': spider}) nextcall = CallLaterOnce(self._next_request, spider) scheduler = self.scheduler_cls.from_crawler(self.crawler) start_requests = yield self.scraper.spidermw.process_start_requests(start_requests, spider) - slot = Slot(start_requests, close_if_idle, nextcall, scheduler) + slot = Slot(start_requests, close_if_idle, nextcall, scheduler, new_queue_behavior) self.slot = slot self.spider = spider yield scheduler.open(spider) diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 7313752b3..be4332142 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -88,7 +88,8 @@ class Crawler: self.spider = self._create_spider(*args, **kwargs) self.engine = self._create_engine() start_requests = yield self.call_start_requests() - yield self.engine.open_spider(self.spider, start_requests) + new_queue_behavior = self.is_start_requests_async(self.spider.start_requests) + yield self.engine.open_spider(self.spider, start_requests, new_queue_behavior=new_queue_behavior) yield defer.maybeDeferred(self.engine.start) except Exception: self.crawling = False @@ -103,7 +104,15 @@ class Crawler: elif inspect.iscoroutinefunction(self.spider.start_requests): return deferred_from_coro(self.spider.start_requests()) else: - return iter(self.spider.start_requests_with_control()) + return iter(self.spider.start_requests()) + + @staticmethod + def is_start_requests_async(start_requests_function): + if hasattr(inspect, 'isasyncgenfunction') and inspect.isasyncgenfunction(start_requests_function): + return True + elif inspect.iscoroutinefunction(start_requests_function): + return True + return False def _create_spider(self, *args, **kwargs): return self.spidercls.from_crawler(self, *args, **kwargs) diff --git a/scrapy/spiders/__init__.py b/scrapy/spiders/__init__.py index 1e30d7b3f..00d7fc708 100644 --- a/scrapy/spiders/__init__.py +++ b/scrapy/spiders/__init__.py @@ -11,10 +11,6 @@ from scrapy.http import Request from scrapy.utils.trackref import object_ref from scrapy.utils.url import url_is_from_spider from scrapy.utils.deprecate import method_is_overridden -from scrapy.utils.python import iflatten - - -WaitUntilQueueEmpty = object() class Spider(object_ref): @@ -81,11 +77,6 @@ class Spider(object_ref): for url in self.start_urls: yield Request(url, dont_filter=True) - def start_requests_with_control(self): - sig = WaitUntilQueueEmpty - gen_ = ((r, sig) if r != sig else (r,) for r in self.start_requests()) - return iflatten(gen_) - def make_requests_from_url(self, url): """ This method is deprecated. """ warnings.warn( diff --git a/tests/middlewares.py b/tests/middlewares.py index 983147642..2f8f7dc09 100644 --- a/tests/middlewares.py +++ b/tests/middlewares.py @@ -1,4 +1,5 @@ from scrapy.http import Request +from scrapy.utils.defer import _isasyncgen class RequestInOrderMiddleware: @@ -6,7 +7,11 @@ class RequestInOrderMiddleware: return (self._preserve_in_order(r, spider) for r in result or ()) def process_start_requests(self, start_requests, spider): - return (self._preserve_in_order(r, spider) for r in start_requests or ()) + if _isasyncgen(start_requests): + from tests.py36.middlewares import RequestInOrderMiddleware_process_start_requests + return RequestInOrderMiddleware_process_start_requests(self._preserve_in_order, start_requests, spider) + else: + return (self._preserve_in_order(r, spider) for r in start_requests or ()) def _preserve_in_order(self, smth, spider): self.__preserve_in_order(smth, spider) diff --git a/tests/py36/_test_crawl.py b/tests/py36/_test_crawl.py index 162a53760..678b19647 100644 --- a/tests/py36/_test_crawl.py +++ b/tests/py36/_test_crawl.py @@ -1,7 +1,7 @@ import asyncio from scrapy import Request -from tests.spiders import SimpleSpider +from tests.spiders import SimpleSpider, YieldingRequestsSpider class AsyncDefAsyncioGenSpider(SimpleSpider): @@ -55,3 +55,9 @@ class AsyncDefAsyncioGenComplexSpider(SimpleSpider): async def parse2(self, response): await asyncio.sleep(0.1) yield {'index2': response.meta['index']} + + +class EagerAsyncGenSpider(YieldingRequestsSpider): + async def start_requests(self): + for r in super().start_requests(): + yield r diff --git a/tests/py36/middlewares.py b/tests/py36/middlewares.py new file mode 100644 index 000000000..6529943d3 --- /dev/null +++ b/tests/py36/middlewares.py @@ -0,0 +1,2 @@ +def RequestInOrderMiddleware_process_start_requests(f, start_requests, spider): + return (f(r, spider) async for r in start_requests or ()) diff --git a/tests/test_crawl.py b/tests/test_crawl.py index 1511b0b30..56b286e03 100644 --- a/tests/test_crawl.py +++ b/tests/test_crawl.py @@ -186,18 +186,14 @@ class CrawlTestCase(TestCase): return 100 @defer.inlineCallbacks - def test_start_requests_eagerness(self): - class EagerSpider(YieldingRequestsSpider): - def start_requests_with_control(self): - yield from self.start_requests() - + def _test_start_requests_eagerness(self, spider_cls): settings = { "SPIDER_MIDDLEWARES": { "scrapy.spidermiddlewares.depth.DepthMiddleware": 0, "tests.middlewares.RequestInOrderMiddleware": 1, } } - crawler = CrawlerRunner(settings).create_crawler(EagerSpider) + crawler = CrawlerRunner(settings).create_crawler(spider_cls) yield crawler.crawl( mockserver=self.mockserver, number_of_start_requests=self.number_of_start_requests_for_eagernees_testing(), @@ -206,6 +202,20 @@ class CrawlTestCase(TestCase): self.is_eager(crawler.spider.requests_in_order_of_scheduling) ) + @defer.inlineCallbacks + def test_start_requests_eagerness_asyncdef(self): + class EagerAsyncDefSpider(YieldingRequestsSpider): + async def start_requests(self): + return list(super().start_requests()) + + yield self._test_start_requests_eagerness(EagerAsyncDefSpider) + + @mark.skipif(sys.version_info < (3, 6), reason="Async generators require Python 3.6 or higher") + @defer.inlineCallbacks + def test_start_requests_eagerness_asyncgen(self): + from tests.py36._test_crawl import EagerAsyncGenSpider + yield self._test_start_requests_eagerness(EagerAsyncGenSpider) + @defer.inlineCallbacks def test_start_requests_lazyness(self): """ diff --git a/tests/test_request_cb_kwargs.py b/tests/test_request_cb_kwargs.py index cfd044252..bd49179aa 100644 --- a/tests/test_request_cb_kwargs.py +++ b/tests/test_request_cb_kwargs.py @@ -59,7 +59,7 @@ class KeywordArgumentsSpider(MockServerSpider): checks = list() - def start_requests_with_control(self): + def start_requests(self): data = {'key': 'value', 'number': 123} yield Request(self.mockserver.url('/first'), self.parse_first, cb_kwargs=data) yield Request(self.mockserver.url('/general_with'), self.parse_general, cb_kwargs=data)