mirror of https://github.com/scrapy/scrapy.git
Remove WaitUntilQueueEmpty handling, enable new queue behavior for async start_requests.
This commit is contained in:
parent
b56eba7952
commit
94ce930743
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -0,0 +1,2 @@
|
|||
def RequestInOrderMiddleware_process_start_requests(f, start_requests, spider):
|
||||
return (f(r, spider) async for r in start_requests or ())
|
||||
|
|
@ -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):
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
Loading…
Reference in New Issue