mirror of https://github.com/scrapy/scrapy.git
WIP
This commit is contained in:
parent
b37f3b5422
commit
4ac921f389
|
|
@ -19,9 +19,6 @@ Backward-incompatible changes
|
|||
- By default, the iteration of start requests and items no longer stops once
|
||||
there are requests in the scheduler.
|
||||
|
||||
You can restore the previous behavior by setting :setting:`SEEDING_POLICY`
|
||||
to :py:enum:mem:`~scrapy.SeedingPolicy.lazy`.
|
||||
|
||||
- In ``scrapy.core.engine.ExecutionEngine``:
|
||||
|
||||
- The second parameter of ``open_spider()``, ``start_requests()``, has
|
||||
|
|
@ -69,21 +66,6 @@ New features
|
|||
|
||||
(:issue:`456`, :issue:`3477`, :issue:`4467`, :issue:`5627`, :issue:`6729`)
|
||||
|
||||
- The new :setting:`SEEDING_POLICY` setting allows customizing how start
|
||||
requests and items are iterated.
|
||||
|
||||
You can also override the active seeding policy from
|
||||
:meth:`Spider.start <scrapy.Spider.start>` and from
|
||||
:meth:`SpiderMiddleware.process_start
|
||||
<scrapy.spidermiddlewares.SpiderMiddleware.process_start>`.
|
||||
|
||||
.. note:: Some third-party spider middlewares may need to be updated for
|
||||
Scrapy VERSION support before you can use them in combination with the
|
||||
ability to override the active seeding policy.
|
||||
|
||||
(:issue:`740`, :issue:`1051`, :issue:`1443`, :issue:`3237`, :issue:`4467`,
|
||||
:issue:`5282`, :issue:`6729`)
|
||||
|
||||
Bug fixes
|
||||
~~~~~~~~~
|
||||
|
||||
|
|
|
|||
|
|
@ -127,10 +127,7 @@ Request objects
|
|||
body to bytes (if given as a string).
|
||||
:type encoding: str
|
||||
|
||||
:param priority: the priority of this request (defaults to ``0``).
|
||||
The priority is used by the scheduler to define the order used to process
|
||||
requests. Requests with a higher priority value will execute earlier.
|
||||
Negative values are allowed in order to indicate relatively low-priority.
|
||||
:param priority: sets :attr:`priority`, defaults to ``0``.
|
||||
:type priority: int
|
||||
|
||||
:param dont_filter: sets :attr:`dont_filter`, defaults to ``False``.
|
||||
|
|
@ -179,6 +176,8 @@ Request objects
|
|||
|
||||
.. autoattribute:: errback
|
||||
|
||||
.. autoattribute:: priority
|
||||
|
||||
.. attribute:: Request.cb_kwargs
|
||||
|
||||
A dictionary that contains arbitrary metadata for this request. Its contents
|
||||
|
|
@ -353,7 +352,7 @@ errors if needed:
|
|||
"https://example.invalid/", # DNS error expected
|
||||
]
|
||||
|
||||
async def start(self):
|
||||
async def yield_seeds(self):
|
||||
for u in self.start_urls:
|
||||
yield scrapy.Request(
|
||||
u,
|
||||
|
|
|
|||
|
|
@ -1733,29 +1733,6 @@ Soft limit (in bytes) for response data being processed.
|
|||
While the sum of the sizes of all responses being processed is above this value,
|
||||
Scrapy does not process new requests.
|
||||
|
||||
.. setting:: SEEDING_POLICY
|
||||
|
||||
SEEDING_POLICY
|
||||
--------------
|
||||
|
||||
.. versionadded:: VERSION
|
||||
|
||||
Default: :py:enum:mem:`SeedingPolicy.greedy <scrapy.SeedingPolicy.greedy>`
|
||||
|
||||
Determines the way :meth:`Spider.start <scrapy.Spider.start>` is
|
||||
iterated.
|
||||
|
||||
Its value may be defined as a member of the :class:`~scrapy.SeedingPolicy` enum
|
||||
(e.g. :py:enum:mem:`SeedingPolicy.lazy <scrapy.SeedingPolicy.lazy>`) or as a
|
||||
matching string (e.g. ``"lazy"``).
|
||||
|
||||
You can also override the active seeding policy from :meth:`Spider.start
|
||||
<scrapy.Spider.start>` and from :meth:`SpiderMiddleware.process_start
|
||||
<scrapy.spidermiddlewares.SpiderMiddleware.process_start>`.
|
||||
|
||||
.. autoenum:: scrapy.SeedingPolicy
|
||||
:members:
|
||||
|
||||
.. setting:: SPIDER_CONTRACTS
|
||||
|
||||
SPIDER_CONTRACTS
|
||||
|
|
|
|||
|
|
@ -85,23 +85,6 @@ one or more of these methods:
|
|||
|
||||
You may yield the same type of objects as :meth:`~scrapy.Spider.start`.
|
||||
|
||||
As with :meth:`~scrapy.Spider.start`, how this method is iterated by
|
||||
default is controlled by :setting:`SEEDING_POLICY`. It is also possible
|
||||
to yield a :class:`~scrapy.SeedingPolicy` enum or a matching string to
|
||||
change the active seeding policy, for example:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
async def process_start(self, start):
|
||||
yield "front_load"
|
||||
async for item_or_request in start:
|
||||
yield item_or_request
|
||||
yield "idle"
|
||||
|
||||
.. tip:: You can also restore the configured seeding policy by
|
||||
:ref:`reading its value <component-settings>` from the
|
||||
:setting:`SEEDING_POLICY` setting and yielding it.
|
||||
|
||||
To write spider middlewares that work on Scrapy versions lower than
|
||||
VERSION, define also a synchronous ``process_start_requests()`` method
|
||||
that returns an iterable. For example:
|
||||
|
|
|
|||
|
|
@ -364,6 +364,57 @@ used by :class:`~scrapy.downloadermiddlewares.useragent.UserAgentMiddleware`::
|
|||
Spider arguments can also be passed through the Scrapyd ``schedule.json`` API.
|
||||
See `Scrapyd documentation`_.
|
||||
|
||||
.. _spider-start:
|
||||
|
||||
Spider start
|
||||
============
|
||||
|
||||
The way :meth:`~scrapy.Spider.start` works may be counterintuitive.
|
||||
|
||||
By default, if you do not use :ref:`await <await>` in
|
||||
:meth:`~scrapy.Spider.start`, or if you instead use the
|
||||
:attr:`~scrapy.Spider.start_urls` attribute, this is what happens:
|
||||
|
||||
- The first 8-16 start requests (based on :setting:`CONCURRENT_REQUESTS` and
|
||||
:setting:`CONCURRENT_REQUESTS_PER_DOMAIN`) are sent in the order in which
|
||||
they are yielded.
|
||||
|
||||
.. note:: Responses may come in a different order.
|
||||
|
||||
- The remaining start requests are sent in reverse order, and only when there
|
||||
are not enough pending requests yielded from callbacks to reach the
|
||||
configured concurrency.
|
||||
|
||||
This is because the :ref:`scheduler <topics-scheduler>`, where pending
|
||||
requests are stored, uses a LIFO (last in, first out) queue by default,
|
||||
configured in the :setting:`SCHEDULER_MEMORY_QUEUE` and
|
||||
:setting:`SCHEDULER_DISK_QUEUE`.
|
||||
|
||||
So, provided all pending requests have the same
|
||||
:attr:`~scrapy.Request.priority`, scheduled requests are sent in reserve
|
||||
order. The first few requests are sent in order only because they are sent
|
||||
as soon as they are scheduled.
|
||||
|
||||
If you need start requests to be sent before requests yielded from spider
|
||||
callbacks, you can set a higher priority for them. For example:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
async def start(self):
|
||||
async for request in super().start():
|
||||
yield request.replace(priority=1)
|
||||
|
||||
If you also need them to be sent in order, you can assign them decreasing
|
||||
priority values. For example:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
async def start(self):
|
||||
priority = len(self.start_urls)
|
||||
async for request in super().start():
|
||||
yield request.replace(priority=priority)
|
||||
priority -= 1
|
||||
|
||||
.. _builtin-spiders:
|
||||
|
||||
Generic Spiders
|
||||
|
|
|
|||
|
|
@ -7,7 +7,6 @@ import sys
|
|||
import warnings
|
||||
|
||||
# Declare top-level shortcuts
|
||||
from scrapy.core._seeding import SeedingPolicy
|
||||
from scrapy.http import FormRequest, Request
|
||||
from scrapy.item import Field, Item
|
||||
from scrapy.selector import Selector
|
||||
|
|
@ -18,7 +17,6 @@ __all__ = [
|
|||
"FormRequest",
|
||||
"Item",
|
||||
"Request",
|
||||
"SeedingPolicy",
|
||||
"Selector",
|
||||
"Spider",
|
||||
"__version__",
|
||||
|
|
|
|||
|
|
@ -9,7 +9,6 @@ from __future__ import annotations
|
|||
from threading import Thread
|
||||
from typing import TYPE_CHECKING, Any
|
||||
|
||||
from scrapy import SeedingPolicy
|
||||
from scrapy.commands import ScrapyCommand
|
||||
from scrapy.http import Request
|
||||
from scrapy.shell import Shell
|
||||
|
|
@ -28,7 +27,6 @@ class Command(ScrapyCommand):
|
|||
"DUPEFILTER_CLASS": "scrapy.dupefilters.BaseDupeFilter",
|
||||
"KEEP_ALIVE": True,
|
||||
"LOGSTATS_INTERVAL": 0,
|
||||
"SEEDING_POLICY": SeedingPolicy.lazy,
|
||||
}
|
||||
|
||||
def syntax(self) -> str:
|
||||
|
|
|
|||
|
|
@ -1,68 +0,0 @@
|
|||
from enum import Enum
|
||||
|
||||
try:
|
||||
from enum_tools.documentation import document_enum
|
||||
except ImportError:
|
||||
|
||||
def document_enum(func): # type: ignore[misc]
|
||||
return func
|
||||
else:
|
||||
# https://github.com/domdfcoding/enum_tools/issues/29
|
||||
import enum_tools.documentation
|
||||
|
||||
enum_tools.documentation.INTERACTIVE = True
|
||||
|
||||
|
||||
@document_enum
|
||||
class SeedingPolicy(Enum):
|
||||
front_load = "front_load"
|
||||
"""The crawl does not start until all start requests have been scheduled.
|
||||
|
||||
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.
|
||||
"""
|
||||
|
||||
greedy = "greedy"
|
||||
"""Iterating start items and requests takes priority over processing
|
||||
scheduled requests.
|
||||
|
||||
Every time a start request is iterated, it is scheduled, and then the next
|
||||
request from the scheduler is sent.
|
||||
|
||||
.. note:: That request sent may not be the scheduled start request
|
||||
depending on the priority of scheduled requests, on the configured
|
||||
:setting:`SCHEDULER` and on certain scheduler settings (e.g.
|
||||
:setting:`SCHEDULER_MEMORY_QUEUE`).
|
||||
|
||||
Best used when prioritizing start requests is important.
|
||||
"""
|
||||
|
||||
idle = "idle"
|
||||
"""A single start item or request is read only when there are neither
|
||||
scheduled nor on-going requests.
|
||||
|
||||
That is, a new start item or request is not read until all requests
|
||||
triggered by the previous start request, directly or indirectly, have been
|
||||
processed.
|
||||
|
||||
Unlike :py:enum:mem:`lazy`, resource savings are prioritized over crawl
|
||||
speed.
|
||||
|
||||
It is functionally equivalent to running a spider multiple times in a row,
|
||||
one per start request.
|
||||
"""
|
||||
|
||||
lazy = "lazy"
|
||||
"""Processing scheduled requests takes priority over iterating start items
|
||||
and requests.
|
||||
|
||||
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 start request priority is not important.
|
||||
|
||||
Switching to :py:enum:mem:`idle` may lower resource usage further at the
|
||||
cost of also lowering crawl speed.
|
||||
"""
|
||||
|
|
@ -24,8 +24,6 @@ from scrapy.utils.log import failure_to_exc_info, logformatter_adapter
|
|||
from scrapy.utils.misc import build_from_crawler, load_object
|
||||
from scrapy.utils.reactor import CallLaterOnce
|
||||
|
||||
from ._seeding import SeedingPolicy
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from collections.abc import AsyncIterable, Callable, Generator
|
||||
|
||||
|
|
@ -43,10 +41,6 @@ logger = logging.getLogger(__name__)
|
|||
_T = TypeVar("_T")
|
||||
|
||||
|
||||
class _SeedingPolicyChange(Exception):
|
||||
pass
|
||||
|
||||
|
||||
class _Slot:
|
||||
def __init__(
|
||||
self,
|
||||
|
|
@ -109,20 +103,7 @@ class ExecutionEngine:
|
|||
spider_closed_callback
|
||||
)
|
||||
self.start_time: float | None = None
|
||||
self._load_seeding_policy()
|
||||
self._start: AsyncIterable[Any] | None = None
|
||||
self._waiting_for_seed: bool = False
|
||||
|
||||
def _load_seeding_policy(self) -> None:
|
||||
try:
|
||||
self._seeding_policy = SeedingPolicy(self.settings["SEEDING_POLICY"])
|
||||
except ValueError:
|
||||
supported_values = ", ".join(policy.value for policy in SeedingPolicy)
|
||||
raise ValueError(
|
||||
f"The value of the SEEDING_POLICY setting "
|
||||
f"({self.settings['SEEDING_POLICY']!r}) is not supported. "
|
||||
f"Supported values: {supported_values}."
|
||||
)
|
||||
|
||||
def _get_scheduler_class(self, settings: BaseSettings) -> type[BaseScheduler]:
|
||||
from scrapy.core.scheduler import BaseScheduler
|
||||
|
|
@ -185,10 +166,7 @@ class ExecutionEngine:
|
|||
self.paused = False
|
||||
|
||||
@inlineCallbacks
|
||||
def _process_next_seed(self):
|
||||
if self._waiting_for_seed:
|
||||
return
|
||||
self._waiting_for_seed = True
|
||||
def _process_next_spider_start_yield(self):
|
||||
try:
|
||||
item_or_request = yield deferred_from_coro(self._start.__anext__())
|
||||
except StopAsyncIteration:
|
||||
|
|
@ -203,68 +181,28 @@ class ExecutionEngine:
|
|||
else:
|
||||
if isinstance(item_or_request, Request):
|
||||
self.crawl(item_or_request)
|
||||
if (
|
||||
self._seeding_policy is not SeedingPolicy.front_load
|
||||
and not self._needs_backout()
|
||||
):
|
||||
self._start_scheduled_request()
|
||||
elif isinstance(item_or_request, (str, SeedingPolicy)):
|
||||
try:
|
||||
self._seeding_policy = SeedingPolicy(item_or_request)
|
||||
except ValueError:
|
||||
valid_policy_strings = ", ".join(
|
||||
policy.value for policy in SeedingPolicy
|
||||
)
|
||||
logger.error(
|
||||
f"Start value {item_or_request!r} has been ignored. "
|
||||
f"Start values of {str} type must be valid seeding "
|
||||
f"policies ({valid_policy_strings})."
|
||||
)
|
||||
self._slot.nextcall.schedule()
|
||||
else:
|
||||
raise _SeedingPolicyChange
|
||||
else:
|
||||
self.scraper.start_itemproc(item_or_request, response=None)
|
||||
self._slot.nextcall.schedule()
|
||||
finally:
|
||||
self._waiting_for_seed = False
|
||||
if self._seeding_policy is SeedingPolicy.front_load and self._start is None:
|
||||
self._slot.nextcall.schedule()
|
||||
|
||||
@inlineCallbacks
|
||||
def _start_next_requests(self) -> Generator[Deferred[Any], Any, None]:
|
||||
def _process_spider_start(self) -> Generator[Deferred[Any], Any, None]:
|
||||
"""Process start items and requests in an asynchronous loop.
|
||||
|
||||
Items are scraped. Requests are scheduled.
|
||||
"""
|
||||
while self._start is not None:
|
||||
yield self._process_next_spider_start_yield()
|
||||
if not self._needs_backout():
|
||||
self._slot.nextcall.schedule()
|
||||
|
||||
def _start_next_requests(self) -> None:
|
||||
if self._slot is None or self._slot.closing is not None or self.paused:
|
||||
return
|
||||
|
||||
try:
|
||||
if self._seeding_policy in {SeedingPolicy.idle, SeedingPolicy.lazy}:
|
||||
while not self._needs_backout():
|
||||
if self._start_scheduled_request() is None:
|
||||
break
|
||||
if (
|
||||
self._start 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._start 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
|
||||
except _SeedingPolicyChange:
|
||||
self._slot.nextcall.schedule()
|
||||
return
|
||||
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()
|
||||
|
|
@ -452,6 +390,7 @@ class ExecutionEngine:
|
|||
assert self.crawler.stats
|
||||
self.crawler.stats.open_spider(spider)
|
||||
yield self.signals.send_catch_log_deferred(signals.spider_opened, spider=spider)
|
||||
self._process_spider_start()
|
||||
self._slot.nextcall.schedule()
|
||||
self._slot.heartbeat.start(self._SLOT_HEARTBEAT_INTERVAL)
|
||||
|
||||
|
|
|
|||
|
|
@ -130,6 +130,16 @@ class Request(object_ref):
|
|||
self._set_body(body)
|
||||
if not isinstance(priority, int):
|
||||
raise TypeError(f"Request priority not an integer: {priority!r}")
|
||||
|
||||
#: Default: ``0``
|
||||
#:
|
||||
#: Value that the :ref:`scheduler <topics-scheduler>` may use for
|
||||
#: request prioritization.
|
||||
#:
|
||||
#: Built-in schedulers prioritize requests with a higher priority
|
||||
#: value.
|
||||
#:
|
||||
#: Negative values are allowed.
|
||||
self.priority: int = priority
|
||||
|
||||
if not (callable(callback) or callback is None):
|
||||
|
|
@ -191,7 +201,7 @@ class Request(object_ref):
|
|||
#:
|
||||
#: When defining the start URLs of a spider through
|
||||
#: :attr:`~scrapy.Spider.start_urls`, this attribute is enabled by
|
||||
#: default. See :meth:`~scrapy.Spider.start`.
|
||||
#: default. See :meth:`~scrapy.Spider.yield_seeds`.
|
||||
self.dont_filter: bool = dont_filter
|
||||
|
||||
self._meta: dict[str, Any] | None = dict(meta) if meta else None
|
||||
|
|
|
|||
|
|
@ -17,8 +17,6 @@ import sys
|
|||
from importlib import import_module
|
||||
from pathlib import Path
|
||||
|
||||
from scrapy import SeedingPolicy
|
||||
|
||||
ADDONS = {}
|
||||
|
||||
AJAXCRAWL_ENABLED = False
|
||||
|
|
@ -310,8 +308,6 @@ SCHEDULER_PRIORITY_QUEUE = "scrapy.pqueues.ScrapyPriorityQueue"
|
|||
|
||||
SCRAPER_SLOT_MAX_ACTIVE_SIZE = 5000000
|
||||
|
||||
SEEDING_POLICY = SeedingPolicy.greedy
|
||||
|
||||
SPIDER_LOADER_CLASS = "scrapy.spiderloader.SpiderLoader"
|
||||
SPIDER_LOADER_WARN_ONLY = False
|
||||
|
||||
|
|
|
|||
|
|
@ -113,20 +113,6 @@ class Spider(object_ref):
|
|||
async def start(self):
|
||||
yield {"foo": "bar"}
|
||||
|
||||
Use :setting:`SEEDING_POLICY` to set how :meth:`start` is
|
||||
iterated by default. It is also
|
||||
possible to yield a :class:`~scrapy.SeedingPolicy` enum or a matching
|
||||
string to change the active seeding policy, for example:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
async def start(self):
|
||||
yield "front_load"
|
||||
yield Request("https://a.example")
|
||||
yield Request("https://b.example")
|
||||
yield self.crawler.settings["SEEDING_POLICY"]
|
||||
yield Request("https://c.example")
|
||||
|
||||
To write spiders that work on Scrapy versions lower than VERSION,
|
||||
define also a synchronous ``start_requests()`` method that returns an
|
||||
iterable. For example:
|
||||
|
|
@ -135,6 +121,8 @@ class Spider(object_ref):
|
|||
|
||||
def start_requests(self):
|
||||
yield Request("https://toscrape.com/")
|
||||
|
||||
.. seealso:: :ref:`spider-start`
|
||||
"""
|
||||
for item_or_request in self.start_requests():
|
||||
yield item_or_request
|
||||
|
|
|
|||
|
|
@ -194,25 +194,15 @@ class TestCrawl(TestCase):
|
|||
|
||||
@defer.inlineCallbacks
|
||||
def test_start_unsupported_output(self):
|
||||
"""Anything that is not a request, a seeding policy or a string (which
|
||||
is assumed to be a seeding policy) is assumed to be an item, avoiding a
|
||||
potentially expensive call to itemadapter.is_item, and letting instead
|
||||
things fail when ItemAdapter is actually used on the corresponding
|
||||
non-item object."""
|
||||
"""Anything that is not a request is assumed to be an item, avoiding a
|
||||
potentially expensive call to itemadapter.is_item(), and letting
|
||||
instead things fail when ItemAdapter is actually used on the
|
||||
corresponding non-item object."""
|
||||
with LogCapture("scrapy", level=logging.ERROR) as log:
|
||||
crawler = get_crawler(StartGoodAndBadOutput)
|
||||
yield crawler.crawl(mockserver=self.mockserver)
|
||||
|
||||
assert len(log.records) == 1
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_start_laziness(self):
|
||||
settings = {"CONCURRENT_REQUESTS": 1, "SEEDING_POLICY": "lazy"}
|
||||
crawler = get_crawler(BrokenStartSpider, settings)
|
||||
yield crawler.crawl(mockserver=self.mockserver)
|
||||
assert crawler.spider.seedsseen.index(None) < crawler.spider.seedsseen.index(
|
||||
99
|
||||
), crawler.spider.seedsseen
|
||||
assert len(log.records) == 0
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_start_dupes(self):
|
||||
|
|
|
|||
|
|
@ -0,0 +1,130 @@
|
|||
from collections import deque
|
||||
|
||||
from twisted.internet.defer import Deferred
|
||||
from twisted.trial.unittest import TestCase
|
||||
|
||||
from scrapy import Request, Spider, signals
|
||||
from scrapy.core.engine import ExecutionEngine
|
||||
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_scheduler import MemoryScheduler
|
||||
|
||||
|
||||
def sleep(seconds: float = 0.001):
|
||||
from twisted.internet import reactor
|
||||
|
||||
deferred: Deferred[None] = Deferred()
|
||||
reactor.callLater(seconds, deferred.callback, None)
|
||||
return maybe_deferred_to_future(deferred)
|
||||
|
||||
|
||||
class MainTestCase(TestCase):
|
||||
@deferred_f_from_coro_f
|
||||
async def test_sleep(self):
|
||||
"""Neither asynchronous sleeps on Spider.start() nor the equivalent on
|
||||
the scheduler (returning no requests while also returning True from
|
||||
the has_pending_requests() method) should cause the spider to miss the
|
||||
processing of any later requests."""
|
||||
seconds = ExecutionEngine._SLOT_HEARTBEAT_INTERVAL + 0.01
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
|
||||
async def start(self):
|
||||
from twisted.internet import reactor
|
||||
|
||||
yield Request("data:,a")
|
||||
|
||||
await sleep(seconds)
|
||||
|
||||
self.crawler.engine._slot.scheduler.pause()
|
||||
self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,b"))
|
||||
|
||||
# During this time, the reactor reports having requests but
|
||||
# returns None.
|
||||
await sleep(seconds)
|
||||
|
||||
self.crawler.engine._slot.scheduler.unpause()
|
||||
|
||||
# The scheduler request is processed.
|
||||
await sleep(seconds)
|
||||
|
||||
yield Request("data:,c")
|
||||
|
||||
await sleep(seconds)
|
||||
|
||||
self.crawler.engine._slot.scheduler.pause()
|
||||
self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,d"))
|
||||
|
||||
# The last start request is processed during the time until the
|
||||
# delayed call below, proving that the start iteration can
|
||||
# finish before a scheduler “sleep” without causing the
|
||||
# scheduler to finish.
|
||||
reactor.callLater(seconds, self.crawler.engine._slot.scheduler.unpause)
|
||||
|
||||
def parse(self, response):
|
||||
pass
|
||||
|
||||
actual_urls = []
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
settings = {"SCHEDULER": MemoryScheduler}
|
||||
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", "data:,c", "data:,d"]
|
||||
assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}"
|
||||
|
||||
|
||||
class MockServerTestCase(TestCase):
|
||||
# See the comment on the matching line above.
|
||||
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_default_behavior(self):
|
||||
"""Verify the behavior that the docs claims is the default when it
|
||||
comes to start request send order."""
|
||||
seconds = 0.1
|
||||
|
||||
def _url(id, delay=seconds):
|
||||
return self.mockserver.url(f"/delay?n={seconds}&{id}")
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
start_urls = [_url("a", delay=0), _url("b"), _url("d"), _url("e")]
|
||||
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 = {"CONCURRENT_REQUESTS": 2}
|
||||
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=}"
|
||||
|
|
@ -1,358 +0,0 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from collections import defaultdict, deque
|
||||
from logging import ERROR
|
||||
|
||||
from testfixtures import LogCapture
|
||||
from twisted.trial.unittest import TestCase
|
||||
|
||||
from scrapy import Request, SeedingPolicy, Spider, signals
|
||||
from scrapy.core.engine import ExecutionEngine
|
||||
from scrapy.core.scheduler import BaseScheduler
|
||||
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_start import twisted_sleep
|
||||
|
||||
|
||||
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
|
||||
|
||||
@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}
|
||||
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_lazy(self):
|
||||
class TestScheduler(BaseScheduler):
|
||||
def __init__(self, *args, **kwargs):
|
||||
self.requests = deque((Request("data:,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 = ["data:,b"]
|
||||
|
||||
def parse(self, response):
|
||||
pass
|
||||
|
||||
actual_urls = []
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "lazy"}
|
||||
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_lazy_blocking(self):
|
||||
"""If the scheduler reports having requests but yields none, the lazy
|
||||
policy schedules start requests."""
|
||||
|
||||
class TestScheduler(BaseScheduler):
|
||||
def __init__(self, *args, **kwargs):
|
||||
self.requests = deque()
|
||||
self.stop = False
|
||||
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
self.requests.append(request)
|
||||
return True
|
||||
|
||||
def has_pending_requests(self) -> bool:
|
||||
return not self.stop
|
||||
|
||||
def next_request(self) -> Request | None:
|
||||
try:
|
||||
return self.requests.popleft()
|
||||
except IndexError:
|
||||
return None
|
||||
|
||||
sleep_seconds = 0.0001
|
||||
|
||||
class TestSpider(Spider):
|
||||
name = "test"
|
||||
|
||||
async def start(self):
|
||||
await maybe_deferred_to_future(twisted_sleep(sleep_seconds))
|
||||
yield Request("data:,a")
|
||||
await maybe_deferred_to_future(twisted_sleep(sleep_seconds))
|
||||
self.crawler.engine._slot.scheduler.enqueue_request(Request("data:,b"))
|
||||
await maybe_deferred_to_future(twisted_sleep(sleep_seconds))
|
||||
yield Request("data:,c")
|
||||
self.crawler.engine._slot.scheduler.stop = True
|
||||
|
||||
def parse(self, response):
|
||||
pass
|
||||
|
||||
actual_urls = []
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "lazy"}
|
||||
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", "data:,c"]
|
||||
assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}"
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_lazy_start_order(self):
|
||||
"""By default, start 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)
|
||||
|
||||
settings = {"SEEDING_POLICY": "lazy"}
|
||||
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", "data:,c"]
|
||||
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 start(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=}"
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_override(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 start(self):
|
||||
yield "front-load" # typo
|
||||
yield SeedingPolicy.front_load
|
||||
yield Request("data:,b", priority=1)
|
||||
yield Request("data:,a", priority=2)
|
||||
yield self.crawler.settings["SEEDING_POLICY"]
|
||||
yield Request("data:,c", priority=3)
|
||||
|
||||
def parse(self, response):
|
||||
pass
|
||||
|
||||
actual_items = []
|
||||
actual_urls = []
|
||||
|
||||
def track_item(item, response, spider):
|
||||
actual_items.append(item)
|
||||
|
||||
def track_url(request, spider):
|
||||
actual_urls.append(request.url)
|
||||
|
||||
settings = {"SCHEDULER": TestScheduler, "SEEDING_POLICY": "lazy"}
|
||||
crawler = get_crawler(TestSpider, settings_dict=settings)
|
||||
crawler.signals.connect(track_item, signals.item_scraped)
|
||||
crawler.signals.connect(track_url, signals.request_reached_downloader)
|
||||
with LogCapture(level=ERROR) as log:
|
||||
await maybe_deferred_to_future(crawler.crawl())
|
||||
assert len(log.records) == 1
|
||||
assert "must be valid seeding policies" in str(log.records[0])
|
||||
assert crawler.stats.get_value("finish_reason") == "finished"
|
||||
assert not actual_items, (
|
||||
f"{actual_items=} should be empty, policies are not items"
|
||||
)
|
||||
expected_urls = ["data:,a", "data:,b", "data:,c"]
|
||||
assert actual_urls == expected_urls, f"{actual_urls=} != {expected_urls=}"
|
||||
|
||||
|
||||
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
|
||||
|
||||
@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={self.delay}&{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=}"
|
||||
|
|
@ -3,6 +3,7 @@ from __future__ import annotations
|
|||
import shutil
|
||||
import tempfile
|
||||
from abc import ABC, abstractmethod
|
||||
from collections import deque
|
||||
from typing import Any, NamedTuple
|
||||
|
||||
import pytest
|
||||
|
|
@ -10,7 +11,7 @@ from twisted.internet import defer
|
|||
from twisted.trial.unittest import TestCase
|
||||
|
||||
from scrapy.core.downloader import Downloader
|
||||
from scrapy.core.scheduler import Scheduler
|
||||
from scrapy.core.scheduler import BaseScheduler, Scheduler
|
||||
from scrapy.crawler import Crawler
|
||||
from scrapy.http import Request
|
||||
from scrapy.spiders import Spider
|
||||
|
|
@ -20,6 +21,38 @@ from scrapy.utils.test import get_crawler
|
|||
from tests.mockserver import MockServer
|
||||
|
||||
|
||||
class MemoryScheduler(BaseScheduler):
|
||||
paused = False
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
self.queue = deque(
|
||||
Request(value) if isinstance(value, str) else value
|
||||
for value in getattr(self, "queue", [])
|
||||
)
|
||||
|
||||
def enqueue_request(self, request: Request) -> bool:
|
||||
self.queue.append(request)
|
||||
return True
|
||||
|
||||
def has_pending_requests(self) -> bool:
|
||||
return self.paused or bool(self.queue)
|
||||
|
||||
def next_request(self) -> Request | None:
|
||||
if self.paused:
|
||||
return None
|
||||
try:
|
||||
return self.queue.pop()
|
||||
except IndexError:
|
||||
return None
|
||||
|
||||
def pause(self) -> None:
|
||||
self.paused = True
|
||||
|
||||
def unpause(self) -> None:
|
||||
self.paused = False
|
||||
|
||||
|
||||
class MockEngine(NamedTuple):
|
||||
downloader: MockDownloader
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue