mirror of https://github.com/scrapy/scrapy.git
Implement seeding policies
This commit is contained in:
parent
0ebc057344
commit
69f829fa5e
|
|
@ -1746,60 +1746,52 @@ The way :meth:`Spider.yield_seeds <scrapy.Spider.yield_seeds>` is iterated:
|
|||
|
||||
- .. _lazy-seeding:
|
||||
|
||||
``"lazy"``: Seeds are only read while the :ref:`scheduler
|
||||
<topics-scheduler>` is empty and the number of ongoing requests is
|
||||
lower than :setting:`CONCURRENT_REQUESTS`.
|
||||
``"lazy"``: Processing scheduled requests takes priority over iterating
|
||||
seeds.
|
||||
|
||||
This seeding policy aims to:
|
||||
|
||||
- Maximize crawl speed by maxing out concurrent requests as often
|
||||
as possible.
|
||||
|
||||
- Minimize the number of requests in the scheduler at any given
|
||||
time by prioritizing scheduler requests over seeds, to minimize
|
||||
resource usage (memory or disk, depending on
|
||||
:setting:`JOBDIR`).
|
||||
|
||||
This seeding policy is best used when seed request priority is not
|
||||
important. Switching to :ref:`serial <serial-seeding>` may lower
|
||||
This seeding policy aims to minimize the number of requests in the
|
||||
scheduler at any given time, to minimize resource usage (memory or disk,
|
||||
depending on :setting:`JOBDIR`). It is best used when seed request priority
|
||||
is not important. Switching to :ref:`idle <idle-seeding>` may lower
|
||||
resource usage further at the cost of also lowering crawl speed.
|
||||
|
||||
- .. _front-load-seeding:
|
||||
|
||||
``"front-load"``: The spider does not start until all seeds have
|
||||
been read and loaded into the scheduler.
|
||||
|
||||
This seeding policy aims to give the :ref:`scheduler
|
||||
<topics-scheduler>` full control over request order, at the cost of
|
||||
a higher resource usage and a delayed crawl start.
|
||||
|
||||
This seeding policy is best used when having all requests go
|
||||
through the scheduler is more important than resource usage and
|
||||
crawl speed.
|
||||
|
||||
- .. _greedy-seeding:
|
||||
|
||||
``"greedy"``: While the :ref:`scheduler <topics-scheduler>` is
|
||||
empty and the number of ongoing requests is lower than
|
||||
:setting:`CONCURRENT_REQUESTS`, seeds are read and sent directly
|
||||
(bypassing the scheduler). While the scheduler has requests, seeds
|
||||
are fed into the scheduler.
|
||||
``"greedy"``: Iterating seeds takes priority over processing scheduled
|
||||
requests.
|
||||
|
||||
This seeding policy is similar to :ref:`front-load
|
||||
<front-load-seeding>`, but it bypasses the scheduler for the first
|
||||
few requests to avoid delaying the crawl start.
|
||||
Every time a seed request is iterated, it is scheduled, and then the next
|
||||
request from the scheduler is sent.
|
||||
|
||||
- .. _serial-seeding:
|
||||
.. note:: That request sent may not be the schedueld seed request
|
||||
depending on the priority of scheduled requests, on the configured
|
||||
:setting:`SCHEDULER` and on certain scheduler settings (e.g.
|
||||
:setting:`SCHEDULER_MEMORY_QUEUE`).
|
||||
|
||||
This seeding policy is best used when prioritizing seed requests is
|
||||
important, and seed requests may be sent as they come.
|
||||
|
||||
- .. _front-load-seeding:
|
||||
|
||||
``"front-load"``: The spider does not start until all seed requests have
|
||||
been scheduled.
|
||||
|
||||
This seeding policy aims to give the :ref:`scheduler <topics-scheduler>`
|
||||
full control over request order from the start. Some custom schedulers may
|
||||
require this seeding policy to work as designed.
|
||||
|
||||
- .. _idle-seeding:
|
||||
|
||||
``"idle"``: A single seed is read only when there are neither scheduled nor
|
||||
on-going requests.
|
||||
|
||||
``"serial"``: A single seed is read whenever the :ref:`scheduler
|
||||
<topics-scheduler>` is empty and there are no ongoing requests.
|
||||
That is, a new seed is not read until all requests triggered by the
|
||||
previous seed, directly or indirectly, have been processed.
|
||||
|
||||
This seeding policy is similar to :ref:`lazy <lazy-seeding>`, but
|
||||
it prioritizes resource savings over crawl speed. It is
|
||||
functionally equivalent to running the spider multiple times in a
|
||||
row, one per seed request.
|
||||
This seeding policy is similar to :ref:`lazy <lazy-seeding>`, but it
|
||||
prioritizes resource savings over crawl speed. It is functionally
|
||||
equivalent to running your spider multiple times in a row, one per seed
|
||||
request.
|
||||
|
||||
.. setting:: SPIDER_CONTRACTS
|
||||
|
||||
|
|
|
|||
|
|
@ -223,7 +223,9 @@ markers = [
|
|||
"requires_botocore: marks tests that need botocore (but not boto3)",
|
||||
"requires_boto3: marks tests that need botocore and boto3",
|
||||
]
|
||||
filterwarnings = []
|
||||
filterwarnings = [
|
||||
"ignore::DeprecationWarning:twisted.web.static"
|
||||
]
|
||||
|
||||
[tool.ruff.lint]
|
||||
extend-select = [
|
||||
|
|
|
|||
|
|
@ -78,10 +78,10 @@ class _Slot:
|
|||
|
||||
|
||||
class _SeedingPolicy(Enum):
|
||||
lazy = "lazy"
|
||||
front_load = "front-load"
|
||||
greedy = "greedy"
|
||||
serial = "serial"
|
||||
idle = "idle"
|
||||
lazy = "lazy"
|
||||
|
||||
|
||||
class ExecutionEngine:
|
||||
|
|
@ -186,11 +186,6 @@ class ExecutionEngine:
|
|||
def unpause(self) -> None:
|
||||
self.paused = False
|
||||
|
||||
def _start_scheduled_requests(self):
|
||||
while not self._needs_backout():
|
||||
if self._start_scheduled_request() is None:
|
||||
break
|
||||
|
||||
@inlineCallbacks
|
||||
def _process_next_seed(self):
|
||||
if self._waiting_for_seed:
|
||||
|
|
@ -210,19 +205,50 @@ class ExecutionEngine:
|
|||
else:
|
||||
if isinstance(seed, Request):
|
||||
self.crawl(seed)
|
||||
if (
|
||||
self._seeding_policy is not _SeedingPolicy.front_load
|
||||
and not self._needs_backout()
|
||||
):
|
||||
self._start_scheduled_request()
|
||||
else:
|
||||
self.scraper.start_itemproc(seed, response=None)
|
||||
self._slot.nextcall.schedule()
|
||||
finally:
|
||||
self._waiting_for_seed = False
|
||||
if self._seeding_policy is _SeedingPolicy.front_load and self._seeds is None:
|
||||
self._slot.nextcall.schedule()
|
||||
|
||||
@inlineCallbacks
|
||||
def _start_next_requests(self) -> Generator[Deferred[Any], Any, None]:
|
||||
if self._slot is None or self._slot.closing is not None or self.paused:
|
||||
return
|
||||
self._start_scheduled_requests()
|
||||
if self._seeds is not None and not self._needs_backout():
|
||||
yield self._process_next_seed()
|
||||
|
||||
if self._seeding_policy in {_SeedingPolicy.idle, _SeedingPolicy.lazy}:
|
||||
while not self._needs_backout():
|
||||
if self._start_scheduled_request() is None:
|
||||
break
|
||||
if (
|
||||
self._seeds is not None
|
||||
and not self._needs_backout()
|
||||
and (
|
||||
self._seeding_policy is not _SeedingPolicy.idle
|
||||
or (not self._waiting_for_seed and not self.downloader.active)
|
||||
)
|
||||
):
|
||||
yield self._process_next_seed()
|
||||
else:
|
||||
assert self._seeding_policy in {
|
||||
_SeedingPolicy.front_load,
|
||||
_SeedingPolicy.greedy,
|
||||
}
|
||||
if self._seeds is not None:
|
||||
if not self._needs_backout():
|
||||
yield self._process_next_seed()
|
||||
else:
|
||||
while not self._needs_backout():
|
||||
if self._start_scheduled_request() is None:
|
||||
break
|
||||
|
||||
if self.spider_is_idle() and self._slot.close_if_idle:
|
||||
self._spider_idle()
|
||||
|
||||
|
|
|
|||
|
|
@ -1,29 +1,36 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from collections import deque
|
||||
from collections import defaultdict, deque
|
||||
|
||||
from twisted.internet.defer import inlineCallbacks
|
||||
from twisted.trial.unittest import TestCase
|
||||
|
||||
from scrapy import Request, Spider, signals
|
||||
from scrapy.core.engine import ExecutionEngine
|
||||
from scrapy.core.scheduler import BaseScheduler
|
||||
from scrapy.utils.defer import maybe_deferred_to_future
|
||||
from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future
|
||||
from scrapy.utils.test import get_crawler
|
||||
|
||||
from .mockserver import MockServer
|
||||
from .test_spider_yield_seeds import twisted_sleep
|
||||
|
||||
|
||||
class MainTestCase(TestCase):
|
||||
@inlineCallbacks
|
||||
def test_scheduler_priority_over_seeds_simple(self):
|
||||
"""The seeding policy is to read seeds into the scheduler while the
|
||||
scheduler is empty, but otherwise priorize requests already in the
|
||||
scheduler.
|
||||
|
||||
This test shows how, given a scheduler pre-filled with a request, that
|
||||
request is sent before sending the first seed request.
|
||||
"""
|
||||
# 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
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_lazy(self):
|
||||
class TestScheduler(BaseScheduler):
|
||||
def __init__(self, *args, **kwargs):
|
||||
self.requests = deque((Request("data:,a"),))
|
||||
|
|
@ -56,20 +63,15 @@ class MainTestCase(TestCase):
|
|||
settings = {"SCHEDULER": TestScheduler}
|
||||
crawler = get_crawler(TestSpider, settings_dict=settings)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
yield crawler.crawl()
|
||||
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=}"
|
||||
|
||||
@inlineCallbacks
|
||||
def test_scheduler_priority_over_seeds_complex(self):
|
||||
"""While the seeding policy is to read seeds into the scheduler while
|
||||
the scheduler is empty and otherwise priorize requests already in the
|
||||
scheduler, this is done in a non-blocking way.
|
||||
|
||||
That is, if the scheduler reports having requests but yields none,
|
||||
requests from seeds will be scheduled.
|
||||
"""
|
||||
@deferred_f_from_coro_f
|
||||
async def test_lazy_blocking(self):
|
||||
"""If the scheduler reports having requests but yields none, the lazy
|
||||
policy schedules requests from seeds."""
|
||||
|
||||
class TestScheduler(BaseScheduler):
|
||||
def __init__(self, *args, **kwargs):
|
||||
|
|
@ -114,7 +116,176 @@ class MainTestCase(TestCase):
|
|||
settings = {"SCHEDULER": TestScheduler}
|
||||
crawler = get_crawler(TestSpider, settings_dict=settings)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
yield crawler.crawl()
|
||||
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=}"
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_lazy_seed_order(self):
|
||||
"""By default, seed requests should be sent in the order in which they
|
||||
are iterated."""
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
start_urls = ["data:,a", "data:,b", "data:,c"]
|
||||
|
||||
def parse(self, response):
|
||||
pass
|
||||
|
||||
actual_urls = []
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
crawler = get_crawler(TestSpider)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
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=}"
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_greedy(self):
|
||||
class TestScheduler(BaseScheduler):
|
||||
def __init__(self, *args, **kwargs):
|
||||
self.requests = deque((Request("data:,b"),))
|
||||
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
self.requests.append(request)
|
||||
return True
|
||||
|
||||
def has_pending_requests(self) -> bool:
|
||||
return bool(self.requests)
|
||||
|
||||
def next_request(self) -> Request | None:
|
||||
try:
|
||||
return self.requests.pop()
|
||||
except IndexError:
|
||||
return None
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
start_urls = ["data:,a"]
|
||||
|
||||
def parse(self, response):
|
||||
pass
|
||||
|
||||
actual_urls = []
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "greedy"}
|
||||
crawler = get_crawler(TestSpider, settings_dict=settings)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
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=}"
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_front_load(self):
|
||||
class TestScheduler(BaseScheduler):
|
||||
def __init__(self, *args, **kwargs):
|
||||
self.requests = defaultdict(deque)
|
||||
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
self.requests[request.priority].append(request)
|
||||
return True
|
||||
|
||||
def has_pending_requests(self) -> bool:
|
||||
return bool(self.requests)
|
||||
|
||||
def next_request(self) -> Request | None:
|
||||
if not self.requests:
|
||||
return None
|
||||
priority = max(self.requests)
|
||||
request = self.requests[priority].popleft()
|
||||
if not self.requests[priority]:
|
||||
del self.requests[priority]
|
||||
return request
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
|
||||
async def yield_seeds(self):
|
||||
yield Request("data:,b", priority=0)
|
||||
yield Request("data:,a", priority=1)
|
||||
|
||||
def parse(self, response):
|
||||
pass
|
||||
|
||||
actual_urls = []
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "front-load"}
|
||||
crawler = get_crawler(TestSpider, settings_dict=settings)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
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=}"
|
||||
|
||||
|
||||
class MockServerTestCase(TestCase):
|
||||
timeout = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL
|
||||
|
||||
@classmethod
|
||||
def setUpClass(cls):
|
||||
cls.mockserver = MockServer()
|
||||
cls.mockserver.__enter__()
|
||||
|
||||
@classmethod
|
||||
def tearDownClass(cls):
|
||||
cls.mockserver.__exit__(None, None, None)
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_idle(self):
|
||||
def _url(id):
|
||||
return self.mockserver.url(f"/delay?n=0.1&{id}")
|
||||
|
||||
class TestScheduler(BaseScheduler):
|
||||
def __init__(self, *args, **kwargs):
|
||||
self.requests = deque((Request(_url("a")),))
|
||||
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
self.requests.append(request)
|
||||
return True
|
||||
|
||||
def has_pending_requests(self) -> bool:
|
||||
return bool(self.requests)
|
||||
|
||||
def next_request(self) -> Request | None:
|
||||
try:
|
||||
return self.requests.popleft()
|
||||
except IndexError:
|
||||
return None
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
start_urls = [_url("b"), _url("d")]
|
||||
queue = deque((Request(_url("c")),))
|
||||
|
||||
def parse(self, response):
|
||||
try:
|
||||
request = self.queue.popleft()
|
||||
except IndexError:
|
||||
pass
|
||||
else:
|
||||
yield request
|
||||
|
||||
actual_urls = []
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "idle"}
|
||||
crawler = get_crawler(TestSpider, settings_dict=settings)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
await maybe_deferred_to_future(crawler.crawl())
|
||||
assert crawler.stats.get_value("finish_reason") == "finished"
|
||||
expected_urls = [_url(letter) for letter in "abcd"]
|
||||
assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}"
|
||||
|
|
|
|||
Loading…
Reference in New Issue