From 72c12f10a79142342953100f99824cae27e0fdce Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Chaves?= Date: Sat, 15 Mar 2025 11:14:34 +0100 Subject: [PATCH] Good Bye, Heartbeat! --- scrapy/core/engine.py | 30 +++- tests/test_engine_seeding.py | 146 +++++++++++++++---- tests/test_scheduler.py | 6 +- tests/test_spider_yield_seeds.py | 40 ----- tests/test_spidermiddleware_process_seeds.py | 67 --------- 5 files changed, 148 insertions(+), 141 deletions(-) diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index ef3d91606..d60d05675 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -13,7 +13,6 @@ from traceback import format_exc from typing import TYPE_CHECKING, Any, TypeVar, cast from twisted.internet.defer import Deferred, inlineCallbacks, succeed -from twisted.internet.task import LoopingCall from twisted.python.failure import Failure from scrapy import signals @@ -60,7 +59,6 @@ class _Slot: self.close_if_idle: bool = close_if_idle self.nextcall: CallLaterOnce[Deferred[None]] = nextcall self.scheduler: BaseScheduler = scheduler - self.heartbeat: LoopingCall = LoopingCall(nextcall.schedule) def add_request(self, request: Request) -> None: self.inprogress.add(request) @@ -78,13 +76,12 @@ class _Slot: if self.closing is not None and not self.inprogress: if self.nextcall: self.nextcall.cancel() - if self.heartbeat.running: - self.heartbeat.stop() self.closing.callback(None) class ExecutionEngine: - _SLOT_HEARTBEAT_INTERVAL: float = 5.0 + _MIN_BACK_IN_SECONDS = 0.001 + _MAX_BACK_IN_SECONDS = 5.0 def __init__( self, @@ -113,6 +110,7 @@ class ExecutionEngine: self._load_seeding_policy() self._seeds: AsyncIterator[Any] | None = None self._waiting_for_seed: bool = False + self._back_in_seconds = self._MIN_BACK_IN_SECONDS def _load_seeding_policy(self) -> None: try: @@ -190,6 +188,10 @@ class ExecutionEngine: if self._waiting_for_seed: return self._waiting_for_seed = True + # Schedule a new call for next requests while waiting for + # self._seeds.__anext__(), so that if it takes long enough and there + # are pending scheduler requests we can process those. + self._slot.nextcall.schedule(self._MIN_BACK_IN_SECONDS) try: seed = yield deferred_from_coro(self._seeds.__anext__()) except StopAsyncIteration: @@ -245,7 +247,8 @@ class ExecutionEngine: if self._start_scheduled_request() is None: break if ( - self._seeds is not None + not self._waiting_for_seed + and self._seeds is not None and not self._needs_backout() and ( self._seeding_policy is not SeedingPolicy.idle @@ -258,7 +261,7 @@ class ExecutionEngine: SeedingPolicy.front_load, SeedingPolicy.greedy, } - if self._seeds is not None: + if not self._waiting_for_seed and self._seeds is not None: if not self._needs_backout(): yield self._process_next_seed() else: @@ -275,6 +278,18 @@ class ExecutionEngine: if self.spider_is_idle() and self._slot.close_if_idle: self._spider_idle() + elif self._needs_backout(): + self._back_in_seconds = self._MIN_BACK_IN_SECONDS + elif self._slot.scheduler.has_pending_requests(): + # If the scheduler reports having pending requests but did not + # actually return one, use exponential backoff to schedule a new + # call to this method, to see if the scheduler finally returns a + # pending request or stops reporting that it has some. + self._slot.nextcall.schedule(self._back_in_seconds) + if self._back_in_seconds != self._MAX_BACK_IN_SECONDS: + self._back_in_seconds = min( + self._back_in_seconds**2, self._MAX_BACK_IN_SECONDS + ) def _needs_backout(self) -> bool: assert self._slot is not None # typing @@ -460,7 +475,6 @@ class ExecutionEngine: self.crawler.stats.open_spider(spider) yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider) self._slot.nextcall.schedule() - self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL) def _spider_idle(self) -> None: """ diff --git a/tests/test_engine_seeding.py b/tests/test_engine_seeding.py index aee7fad33..db59fdfea 100644 --- a/tests/test_engine_seeding.py +++ b/tests/test_engine_seeding.py @@ -4,6 +4,7 @@ from collections import deque from logging import ERROR from testfixtures import LogCapture +from twisted.internet.defer import Deferred from twisted.trial.unittest import TestCase from scrapy import Request, SeedingPolicy, Spider, signals @@ -15,21 +16,15 @@ from tests.test_scheduler import MemoryScheduler, PriorityScheduler from .mockserver import MockServer -class MainTestCase(TestCase): - # If the test ends before the heartbeat, it may mean that the logic to - # re-schecule a new call of _start_next_requests under the right - # ciscumstances is not properly implemented, and the hearatbeat is working - # as a workaround for that issue. This is a performance issue and should - # be addressed. - # - # It could also happen that, on some CI runners, some tests (e.g. those - # below using a mock server) run too slow and proper handling overlaps with - # the heartbeat. If that is the case, it may be worth considering - # increasing the heartbeat time. It should be safe, since in most real live - # scenarios the heartbeat should never make a difference, and we may - # eventually remove the heartbeat altogether. - timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL +def sleep(seconds: float = ExecutionEngine._MIN_BACK_IN_SECONDS): + from twisted.internet import reactor + deferred = Deferred() + reactor.callLater(seconds, deferred.callback, None) + return maybe_deferred_to_future(deferred) + + +class MainTestCase(TestCase): @deferred_f_from_coro_f async def test_greedy(self): class TestScheduler(MemoryScheduler): @@ -55,6 +50,114 @@ class MainTestCase(TestCase): expected_urls = ["data:,a", "data:,b"] assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" + @deferred_f_from_coro_f + async def test_greedy_sleep(self): + """If the seeds sleep long enough, scheduler requests should be + processed in the meantime.""" + + class TestScheduler(MemoryScheduler): + queue = ["data:,b"] + + class TestSpider(Spider): + name = "test" + + async def yield_seeds(self): + yield Request("data:,a") + await sleep(ExecutionEngine._MIN_BACK_IN_SECONDS * 2) + yield Request("data:,c") + + def parse(self, response): + pass + + actual_urls = [] + + def track_url(request, spider): + actual_urls.append(request.url) + + settings = {"SCHEDULER": TestScheduler} + crawler = get_crawler(TestSpider, settings_dict=settings) + crawler.signals.connect(track_url, signals.request_reached_downloader) + with LogCapture() as log: + await maybe_deferred_to_future(crawler.crawl()) + assert crawler.stats.get_value("finish_reason") == "finished" + expected_urls = ["data:,a", "data:,b", "data:,c"] + assert actual_urls == expected_urls, ( + f"{actual_urls=} != {expected_urls=}\n{log}" + ) + + @deferred_f_from_coro_f + async def test_greedy_scheduler_sleep(self): + """If the scheduler sleeps but not longer than the seeds, its + processing should resume before that of the seeds, instead of being + blocked by the seeds finishing processing.""" + + class TestScheduler(MemoryScheduler): + pause = True + queue = ["data:,a"] + + class TestSpider(Spider): + name = "test" + + async def yield_seeds(self): + seconds = ExecutionEngine._MIN_BACK_IN_SECONDS + await sleep(seconds) + self.crawler.engine._slot.scheduler.pause = False + await sleep(seconds) + yield Request("data:,b") + + def parse(self, response): + pass + + actual_urls = [] + + def track_url(request, spider): + actual_urls.append(request.url) + + settings = {"SCHEDULER": TestScheduler} + crawler = get_crawler(TestSpider, settings_dict=settings) + crawler.signals.connect(track_url, signals.request_reached_downloader) + with LogCapture() as log: + await maybe_deferred_to_future(crawler.crawl()) + assert crawler.stats.get_value("finish_reason") == "finished" + expected_urls = ["data:,a", "data:,b"] + assert actual_urls == expected_urls, ( + f"{actual_urls=} != {expected_urls=}\n{log}" + ) + + @deferred_f_from_coro_f + async def test_greedy_exception(self): + """If the seeds raise an unhandled exception, scheduler requests should + still be processed.""" + + class TestScheduler(MemoryScheduler): + queue = ["data:,b"] + + class TestSpider(Spider): + name = "test" + + async def yield_seeds(self): + yield Request("data:,a") + raise RuntimeError + + def parse(self, response): + pass + + actual_urls = [] + + def track_url(request, spider): + actual_urls.append(request.url) + + settings = {"SCHEDULER": TestScheduler} + crawler = get_crawler(TestSpider, settings_dict=settings) + crawler.signals.connect(track_url, signals.request_reached_downloader) + with LogCapture() as log: + await maybe_deferred_to_future(crawler.crawl()) + assert crawler.stats.get_value("finish_reason") == "finished" + expected_urls = ["data:,a", "data:,b"] + assert actual_urls == expected_urls, ( + f"{actual_urls=} != {expected_urls=}\n{log}" + ) + @deferred_f_from_coro_f async def test_lazy(self): class TestScheduler(MemoryScheduler): @@ -81,26 +184,21 @@ class MainTestCase(TestCase): assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}" @deferred_f_from_coro_f - async def test_lazy_blocking(self): + async def test_lazy_sleep(self): """If the scheduler reports having requests but yields none, the lazy policy schedules requests from seeds.""" class TestScheduler(MemoryScheduler): - def __init__(self, *args, **kwargs): - super().__init__(*args, **kwargs) - self.stop = False - - def has_pending_requests(self) -> bool: - return not self.stop + queue = ["data:,b"] + pause = True class TestSpider(Spider): name = "test" async def yield_seeds(self): - self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,b")) + self.crawler.engine._slot.scheduler.pause = False yield Request("data:,a") yield Request("data:,c") - self.crawler.engine._slot.scheduler.stop = True def parse(self, response): pass @@ -213,8 +311,6 @@ class MainTestCase(TestCase): class MockServerTestCase(TestCase): - # See the comment on the matching line above. - timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL # If requests are too fast, test_idle will fail because the outcome will # match that of the lazy seeding policy. delay = 0.2 diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 0c61c6e7d..559bd7571 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -22,6 +22,8 @@ from tests.mockserver import MockServer class MemoryScheduler(BaseScheduler): + pause = False + def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self.queue = deque( @@ -34,9 +36,11 @@ class MemoryScheduler(BaseScheduler): return True def has_pending_requests(self) -> bool: - return bool(self.queue) + return self.pause or bool(self.queue) def next_request(self) -> Request | None: + if self.pause: + return None try: return self.queue.pop() except IndexError: diff --git a/tests/test_spider_yield_seeds.py b/tests/test_spider_yield_seeds.py index 3212144de..cb1f76121 100644 --- a/tests/test_spider_yield_seeds.py +++ b/tests/test_spider_yield_seeds.py @@ -1,41 +1,22 @@ -from asyncio import sleep - import pytest from testfixtures import LogCapture from twisted import version as TWISTED_VERSION -from twisted.internet.defer import Deferred from twisted.python.versions import Version from twisted.trial.unittest import TestCase from scrapy import Spider, signals -from scrapy.core.engine import ExecutionEngine from scrapy.exceptions import ScrapyDeprecationWarning from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future from scrapy.utils.test import get_crawler from .test_scheduler import MemoryScheduler -# These are the minimum seconds necessary to wait to reproduce the issue that -# has been solved by catching the RuntimeError exception in the -# ExecutionEngine._next_request() method. A lower value makes these tests pass -# even if we remove that exception handling, but they start failing with this -# much delay. -ASYNC_GEN_ERROR_MINIMUM_SECONDS = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL + 0.01 - ITEM_A = {"id": "a"} ITEM_B = {"id": "b"} TWISTED_KEEPS_TRACEBACKS = TWISTED_VERSION >= Version("twisted", 24, 10, 0) -def twisted_sleep(seconds): - from twisted.internet import reactor - - d = Deferred() - reactor.callLater(seconds, d.callback, None) - return d - - class MainTestCase(TestCase): # Utility methods @@ -122,27 +103,6 @@ class MainTestCase(TestCase): await self._test_spider(TestSpider, [ITEM_A]) - # Delays. - - @pytest.mark.only_asyncio - @deferred_f_from_coro_f - async def test_asyncio_delayed(self): - async def yield_seeds(spider): - await sleep(ASYNC_GEN_ERROR_MINIMUM_SECONDS) - yield ITEM_A - - await self._test_yield_seeds(yield_seeds, [ITEM_A]) - - @deferred_f_from_coro_f - async def test_twisted_delayed(self): - async def yield_seeds(spider): - await maybe_deferred_to_future( - twisted_sleep(ASYNC_GEN_ERROR_MINIMUM_SECONDS) - ) - yield ITEM_A - - await self._test_yield_seeds(yield_seeds, [ITEM_A]) - # Bad definitions. @deferred_f_from_coro_f diff --git a/tests/test_spidermiddleware_process_seeds.py b/tests/test_spidermiddleware_process_seeds.py index f97dfce49..396695c6f 100644 --- a/tests/test_spidermiddleware_process_seeds.py +++ b/tests/test_spidermiddleware_process_seeds.py @@ -1,5 +1,3 @@ -from asyncio import sleep - import pytest from testfixtures import LogCapture from twisted.trial.unittest import TestCase @@ -10,9 +8,7 @@ from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future from scrapy.utils.test import get_crawler from .test_spider_yield_seeds import ( - ASYNC_GEN_ERROR_MINIMUM_SECONDS, TWISTED_KEEPS_TRACEBACKS, - twisted_sleep, ) ITEM_A = {"id": "a"} @@ -20,36 +16,6 @@ ITEM_B = {"id": "b"} ITEM_C = {"id": "c"} ITEM_D = {"id": "d"} - -class AsyncioSleepSpiderMiddleware: - async def process_seeds(self, seeds): - await sleep(ASYNC_GEN_ERROR_MINIMUM_SECONDS) - async for seed in seeds: - yield seed - - -class NoOpSpiderMiddleware: - async def process_seeds(self, seeds): - async for seed in seeds: - yield seed - - -class TwistedSleepSpiderMiddleware: - async def process_seeds(self, seeds): - await maybe_deferred_to_future(twisted_sleep(ASYNC_GEN_ERROR_MINIMUM_SECONDS)) - async for seed in seeds: - yield seed - - -class UniversalSpiderMiddleware: - async def process_seeds(self, seeds): - async for seed in seeds: - yield seed - - def process_start_requests(self, start_requests, spider): - raise NotImplementedError - - # Spiders and spider middlewares for MainTestCase._test_wrap @@ -207,39 +173,6 @@ class MainTestCase(TestCase): ): await self._test_wrap(DeprecatedWrapSpiderMiddleware, DeprecatedWrapSpider) - # Sleep tests - - async def _test_sleep(self, spider_middlewares): - class TestSpider(Spider): - name = "test" - - async def yield_seeds(self): - yield ITEM_A - - await self._test(spider_middlewares, TestSpider, [ITEM_A]) - - @pytest.mark.only_asyncio - @deferred_f_from_coro_f - async def test_asyncio_sleep_single(self): - await self._test_sleep([AsyncioSleepSpiderMiddleware]) - - @pytest.mark.only_asyncio - @deferred_f_from_coro_f - async def test_asyncio_sleep_multiple(self): - await self._test_sleep( - [NoOpSpiderMiddleware, AsyncioSleepSpiderMiddleware, NoOpSpiderMiddleware] - ) - - @deferred_f_from_coro_f - async def test_twisted_sleep_single(self): - await self._test_sleep([TwistedSleepSpiderMiddleware]) - - @deferred_f_from_coro_f - async def test_twisted_sleep_multiple(self): - await self._test_sleep( - [NoOpSpiderMiddleware, TwistedSleepSpiderMiddleware, NoOpSpiderMiddleware] - ) - # Bad definitions. @deferred_f_from_coro_f