Good Bye, Heartbeat!

This commit is contained in:
Adrián Chaves 2025-03-15 11:14:34 +01:00
parent 1305cfb240
commit 72c12f10a7
5 changed files with 148 additions and 141 deletions

View File

@ -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:
"""

View File

@ -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

View File

@ -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:

View File

@ -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

View File

@ -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