Merge remote-tracking branch 'scrapy/master' into async-seeds

This commit is contained in:
Adrián Chaves 2025-05-07 21:16:13 +02:00
commit 51c39308d8
41 changed files with 1649 additions and 753 deletions

View File

@ -38,6 +38,9 @@ jobs:
- python-version: pypy3.10
env:
TOXENV: pypy3
- python-version: pypy3.11
env:
TOXENV: pypy3
# pinned deps
- python-version: "3.9.21"
@ -59,7 +62,7 @@ jobs:
- python-version: "3.13"
env:
TOXENV: extra-deps
- python-version: pypy3.10
- python-version: pypy3.11
env:
TOXENV: pypy3-extra-deps
- python-version: "3.13"

View File

@ -1,22 +0,0 @@
include CODE_OF_CONDUCT.md
include CONTRIBUTING.md
include INSTALL.md
include NEWS
include SECURITY.md
include scrapy/VERSION
include scrapy/mime.types
include scrapy/py.typed
include codecov.yml
include conftest.py
include tox.ini
recursive-include scrapy/templates *
recursive-include docs *
prune docs/build
recursive-include extras *
recursive-include tests *
global-exclude __pycache__ *.py[cod]

23
docs/_templates/layout.html vendored Normal file
View File

@ -0,0 +1,23 @@
{% extends "!layout.html" %}
{# Overriden to include a link to scrapy.org, not just to the docs root #}
{%- block sidebartitle %}
{# the logo helper function was removed in Sphinx 6 and deprecated since Sphinx 4 #}
{# the master_doc variable was renamed to root_doc in Sphinx 4 (master_doc still exists in later Sphinx versions) #}
{%- set _logo_url = logo_url|default(pathto('_static/' + (logo or ""), 1)) %}
{%- set _root_doc = root_doc|default(master_doc) %}
<a href="https://scrapy.org">scrapy.org</a> / <a href="{{ pathto(_root_doc) }}">docs</a>
{%- if READTHEDOCS or DEBUG %}
{%- if theme_version_selector or theme_language_selector %}
<div class="switch-menus">
<div class="version-switch"></div>
<div class="language-switch"></div>
</div>
{%- endif %}
{%- endif %}
{%- include "searchbox.html" %}
{%- endblock %}

View File

@ -96,26 +96,14 @@ How can I simulate a user login in my spider?
See :ref:`topics-request-response-ref-request-userlogin`.
.. _faq-bfo-dfo:
Does Scrapy crawl in breadth-first or depth-first order?
--------------------------------------------------------
By default, Scrapy uses a `LIFO`_ queue for storing pending requests, which
basically means that it crawls in `DFO order`_. This order is more convenient
in most cases.
:ref:`DFO by default, but other orders are possible <request-order>`.
To crawl in `BFO order`_:
- Set :setting:`DEPTH_PRIORITY` to ``1``.
- Set :setting:`SCHEDULER_MEMORY_QUEUE` to
:class:`~scrapy.squeues.FifoMemoryQueue`.
- Set :setting:`SCHEDULER_DISK_QUEUE` to
:class:`~scrapy.squeues.PickleFifoDiskQueue`.
- :ref:`Send start requests before other requests <start-requests-order>`.
My Scrapy crawler has memory leaks. What can I do?
--------------------------------------------------
@ -431,6 +419,3 @@ See :issue:`2680`.
.. _Python standard library modules: https://docs.python.org/3/py-modindex.html
.. _Python package: https://pypi.org/
.. _user agents: https://en.wikipedia.org/wiki/User_agent
.. _LIFO: https://en.wikipedia.org/wiki/Stack_(abstract_data_type)
.. _DFO order: https://en.wikipedia.org/wiki/Depth-first_search
.. _BFO order: https://en.wikipedia.org/wiki/Breadth-first_search

View File

@ -23,13 +23,19 @@ Modified requirements
Backward-incompatible changes
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
- The ``from_settings()`` method of
:class:`~scrapy.spidermiddlewares.urllength.UrlLengthMiddleware` is removed
without a deprecation period (this was needed because after the
introduction of the
:class:`~scrapy.spidermiddlewares.base.BaseSpiderMiddleware` base class and
switching built-in spider middlewares to it those middlewares need the
:class:`~scrapy.crawler.Crawler` instance at run time). Please use
``from_crawler()`` instead.
- The iteration of start requests and items no longer stops once there are
requests in the scheduler, and instead runs continuously until all start
requests have been scheduled.
As a result, the order in which start requests are sent may change. See
:ref:`start-requests` for details.
To restore the previous behavior, :ref:`use lazy start request scheduling
<start-requests-lazy>`.
@ -63,6 +69,12 @@ Backward-incompatible changes
instead of being defined as a generator, is now executed *after* the
:ref:`scheduler <topics-scheduler>` instance has been created.
- When using :setting:`JOBDIR`, :ref:`start requests <start-requests>` are
now serialized into their own, ``s``-suffixed priority folders. You can set
:setting:`SCHEDULER_START_DISK_QUEUE` to ``None`` or ``""`` to change that,
but the side effects may be undesirable. See
:setting:`SCHEDULER_START_DISK_QUEUE` for details.
Deprecations
~~~~~~~~~~~~
@ -79,6 +91,12 @@ Deprecations
(:issue:`456`, :issue:`3477`, :issue:`4467`, :issue:`5627`, :issue:`6729`)
- The ``__init__`` method of priority queue classes (see
:setting:`SCHEDULER_PRIORITY_QUEUE`) should now support a keyword-only
``start_queue_cls`` parameter.
(:issue:`6752`)
New features
~~~~~~~~~~~~
@ -133,6 +151,13 @@ New features
- Added new signals: :signal:`spider_start_blocking`,
:signal:`scheduler_empty`.
- Added new settings: :setting:`SCHEDULER_START_DISK_QUEUE` and
:setting:`SCHEDULER_START_MEMORY_QUEUE`.
- Added :class:`~scrapy.spidermiddlewares.start.StartSpiderMiddleware`, which
sets :reqmeta:`is_start_request` to ``True`` on :ref:`start requests
<start-requests>`.
- In :class:`~scrapy.core.engine.ExecutionEngine`:
- Added a :attr:`~scrapy.core.engine.ExecutionEngine.scheduler`

View File

@ -37,6 +37,10 @@ That includes the classes that you may assign to the following settings:
- :setting:`SCHEDULER_PRIORITY_QUEUE`
- :setting:`SCHEDULER_START_DISK_QUEUE`
- :setting:`SCHEDULER_START_MEMORY_QUEUE`
- :setting:`SPIDER_MIDDLEWARES`
Third-party Scrapy components may also let you define additional Scrapy

View File

@ -73,18 +73,103 @@ In addition to native coroutine APIs Scrapy has some APIs that return a
:class:`~twisted.internet.defer.Deferred` object or take a user-supplied
function that returns a :class:`~twisted.internet.defer.Deferred` object. These
APIs are also asynchronous but don't yet support native ``async def`` syntax.
For example:
In the future we plan to add support for the ``async def`` syntax to these APIs
or replace them with other APIs where changing the existing ones is
possible.
- The :meth:`ExecutionEngine.download` method returns a
:class:`~twisted.internet.defer.Deferred` object.
- A custom download handler needs to define a ``download_request()`` method that
returns a :class:`~twisted.internet.defer.Deferred` object.
The following Scrapy methods return :class:`~twisted.internet.defer.Deferred`
objects (this list is not complete as it only includes methods that we think
may be useful for user code):
- :class:`scrapy.crawler.Crawler`:
- :meth:`~scrapy.crawler.Crawler.crawl`
- :meth:`~scrapy.crawler.Crawler.stop`
- :class:`scrapy.crawler.CrawlerRunner` (also inherited by
:class:`scrapy.crawler.CrawlerProcess`):
- :meth:`~scrapy.crawler.CrawlerRunner.crawl`
- :meth:`~scrapy.crawler.CrawlerRunner.stop`
- :meth:`~scrapy.crawler.CrawlerRunner.join`
- :class:`scrapy.core.engine.ExecutionEngine`:
- :meth:`~scrapy.core.engine.ExecutionEngine.download`
- :class:`scrapy.signalmanager.SignalManager`:
- :meth:`~scrapy.signalmanager.SignalManager.send_catch_log_deferred`
- :class:`~scrapy.mail.MailSender`
- :meth:`~scrapy.mail.MailSender.send`
The following user-supplied methods can return
:class:`~twisted.internet.defer.Deferred` objects (the methods that can also
return coroutines are listed in :ref:`coroutine-support`):
- Custom download handlers (see :setting:`DOWNLOAD_HANDLERS`):
- ``download_request()``
- ``close()``
- Custom downloader implementations (see :setting:`DOWNLOADER`):
- ``fetch()``
- Custom scheduler implementations (see :setting:`SCHEDULER`):
- :meth:`~scrapy.core.scheduler.BaseScheduler.open`
- :meth:`~scrapy.core.scheduler.BaseScheduler.close`
- Custom dupefilters (see :setting:`DUPEFILTER_CLASS`):
- ``open()``
- ``close()``
- Custom feed storages (see :setting:`FEED_STORAGES`):
- ``store()``
- Subclasses of :class:`scrapy.pipelines.media.MediaPipeline`:
- ``media_to_download()``
- ``item_completed()``
- Custom storages used by subclasses of
:class:`scrapy.pipelines.files.FilesPipeline`:
- ``persist_file()``
- ``stat_file()``
In most cases you can use these APIs in code that otherwise uses coroutines, by
wrapping a :class:`~twisted.internet.defer.Deferred` object into a
:class:`~asyncio.Future` object or vice versa. See :ref:`asyncio-await-dfd` for
more information about this.
For example:
- The :meth:`ExecutionEngine.download()
<scrapy.core.engine.ExecutionEngine.download>` method returns a
:class:`~twisted.internet.defer.Deferred` object that fires with the
downloaded response. You can use this object directly in Deferred-based
code or convert it into a :class:`~asyncio.Future` object with
:func:`~scrapy.utils.defer.maybe_deferred_to_future`.
- A custom download handler needs to define a ``download_request()`` method
that returns a :class:`~twisted.internet.defer.Deferred` object. You can
write a method that works with Deferreds and returns one directly, or you
can write a coroutine and convert it into a function that returns a
Deferred with :func:`~scrapy.utils.defer.deferred_f_from_coro_f`.
General usage
=============

View File

@ -517,18 +517,18 @@ as a fallback value if that key is not provided for a specific feed definition:
FEED_EXPORT_ENCODING
--------------------
Default: ``None``
Default: ``"utf-8"`` (:ref:`fallback <default-settings>`: ``None``)
The encoding to be used for the feed.
If unset or set to ``None`` (default) it uses UTF-8 for everything except JSON output,
which uses safe numeric encoding (``\uXXXX`` sequences) for historic reasons.
If set to ``None``, it uses UTF-8 for everything except JSON output, which uses
safe numeric encoding (``\uXXXX`` sequences) for historic reasons.
Use ``utf-8`` if you want UTF-8 for JSON too.
Use ``"utf-8"`` if you want UTF-8 for JSON too.
.. versionchanged:: 2.8
The :command:`startproject` command now sets this setting to
``utf-8`` in the generated ``settings.py`` file.
``"utf-8"`` in the generated ``settings.py`` file.
.. setting:: FEED_EXPORT_FIELDS

View File

@ -46,9 +46,9 @@ Keeping persistent state between batches
Sometimes you'll want to keep some persistent spider state between pause/resume
batches. You can use the ``spider.state`` attribute for that, which should be a
dict. There's :ref:`a built-in extension <topics-extensions-ref-spiderstate>` that takes care of serializing, storing and
loading that attribute from the job directory, when the spider starts and
stops.
dict. There's :ref:`a built-in extension <topics-extensions-ref-spiderstate>`
that takes care of serializing, storing and loading that attribute from the job
directory, when the spider starts and stops.
Here's an example of a callback that uses the spider state (other spider code
is omitted for brevity):

View File

@ -647,6 +647,7 @@ Those are:
* ``ftp_user`` (See :setting:`FTP_USER` for more info)
* :reqmeta:`handle_httpstatus_all`
* :reqmeta:`handle_httpstatus_list`
* :reqmeta:`is_start_request`
* :reqmeta:`max_retry_times`
* :reqmeta:`proxy`
* :reqmeta:`redirect_reasons`

View File

@ -54,8 +54,8 @@ Built-in scheduler
------------------
.. autoclass:: Scheduler()
:members: __len__, pause, unpause
:member-order: bysource
:members:
:special-members: __init__, __len__, pause, unpause
.. _priority-queues:
@ -119,6 +119,7 @@ Writing a priority queue
:special-members: __init__, __len__
:member-order: bysource
.. _custom-internal-queue:
Writing an internal queue

View File

@ -162,8 +162,17 @@ Those command-specific default settings are specified in the
6. Default global settings
--------------------------
The global defaults are located in the ``scrapy.settings.default_settings``
module and documented in the :ref:`topics-settings-ref` section.
The ``scrapy.settings.default_settings`` module defines global default values
for some :ref:`built-in settings <topics-settings-ref>`.
.. note:: :command:`startproject` generates a ``settings.py`` file that sets
some settings to different values.
The reference documentation of settings indicates the default value if one
exists. If :command:`startproject` sets a value, that value is documented
as default, and the value from ``scrapy.settings.default_settings`` is
documented as “fallback”.
Compatibility with pickle
=========================
@ -461,7 +470,7 @@ Note that the event loop class must inherit from :class:`asyncio.AbstractEventLo
BOT_NAME
--------
Default: ``'scrapybot'``
Default: ``<project name>`` (:ref:`fallback <default-settings>`: ``'scrapybot'``)
The name of the bot implemented by this Scrapy project (also known as the
project name). This name will be used for the logging too.
@ -1320,6 +1329,7 @@ Default: ``{}``
A dict containing the pipelines enabled by default in Scrapy. You should never
modify this setting in your project, modify :setting:`ITEM_PIPELINES` instead.
.. setting:: JOBDIR
JOBDIR
@ -1330,6 +1340,7 @@ Default: ``None``
A string indicating the directory for storing the state of a crawl when
:ref:`pausing and resuming crawls <topics-jobs>`.
.. setting:: LOG_ENABLED
LOG_ENABLED
@ -1566,7 +1577,7 @@ email notifying about it. If zero, no warning will be produced.
NEWSPIDER_MODULE
----------------
Default: ``''``
Default: ``"<project name>.spiders"`` (:ref:`fallback <default-settings>`: ``""``)
Module where to create new spiders using the :command:`genspider` command.
@ -1625,9 +1636,7 @@ Adjust redirect request priority relative to original request:
ROBOTSTXT_OBEY
--------------
Default: ``False``
Scope: ``scrapy.downloadermiddlewares.robotstxt``
Default: ``True`` (:ref:`fallback <default-settings>`: ``False``)
If enabled, Scrapy will respect robots.txt policies. For more information see
:ref:`topics-dlmw-robots`.
@ -1737,6 +1746,51 @@ The default scheduler, :class:`~scrapy.core.scheduler.Scheduler`, uses this
component.
.. setting:: SCHEDULER_START_DISK_QUEUE
SCHEDULER_START_DISK_QUEUE
--------------------------
Default: :class:`~scrapy.squeues.PickleFifoDiskQueue`
Type of disk queue (see :setting:`JOBDIR`) that the :ref:`scheduler
<topics-scheduler>` uses for :ref:`start requests <start-requests>`.
For available choices, see :ref:`disk-queues` and :ref:`custom-internal-queue`.
.. queue-common-starts
Use ``None`` or ``""`` to disable these separate queues entirely, and instead
have start requests share the same queues as other requests.
.. note::
Disabling separate start request queues makes :ref:`start request order
<start-request-order>` unintuitive: start requests will be sent in order
only until :setting:`CONCURRENT_REQUESTS` is reached, then remaining start
requests will be sent in reverse order.
.. queue-common-ends
.. setting:: SCHEDULER_START_MEMORY_QUEUE
SCHEDULER_START_MEMORY_QUEUE
----------------------------
Default: :class:`~scrapy.squeues.FifoMemoryQueue`
Type of in-memory queue that the :ref:`scheduler <topics-scheduler>` uses for
:ref:`start requests <start-requests>`.
For available choices, see :ref:`memory-queues` and
:ref:`custom-internal-queue`.
.. include:: settings.rst
:start-after: queue-common-starts
:end-before: queue-common-ends
.. setting:: SCRAPER_SLOT_MAX_ACTIVE_SIZE
SCRAPER_SLOT_MAX_ACTIVE_SIZE
@ -1856,7 +1910,7 @@ the spider. For more info see :ref:`topics-spider-middleware-setting`.
SPIDER_MODULES
--------------
Default: ``[]``
Default: ``["<project name>.spiders"]`` (:ref:`fallback <default-settings>`: ``[]``)
A list of modules where Scrapy will look for spiders.

View File

@ -190,6 +190,19 @@ one or more of these methods:
:param spider: the spider which raised the exception
:type spider: :class:`~scrapy.Spider` object
Base class for custom spider middlewares
----------------------------------------
Scrapy provides a base class for custom spider middlewares. It's not required
to use it but it can help with simplifying middleware implementations and
reducing the amount of boilerplate code in :ref:`universal middlewares
<universal-spider-middleware>`.
.. module:: scrapy.spidermiddlewares.base
.. autoclass:: BaseSpiderMiddleware
:members:
.. _topics-spider-middleware-ref:
Built-in spider middleware reference
@ -405,6 +418,14 @@ String value Class name (as a string)
.. _"unsafe-url": https://www.w3.org/TR/referrer-policy/#referrer-policy-unsafe-url
StartSpiderMiddleware
---------------------
.. module:: scrapy.spidermiddlewares.start
.. autoclass:: StartSpiderMiddleware
UrlLengthMiddleware
-------------------

View File

@ -369,81 +369,12 @@ See `Scrapyd documentation`_.
Start requests
==============
**Start requests** are the :ref:`requests <request>` that serve as a starting
point for a spider.
**Start requests** are :class:`~scrapy.Request` objects yielded from the
:meth:`~scrapy.Spider.start` method of a spider or from the
:meth:`~scrapy.spidermiddlewares.SpiderMiddleware.process_start` method of a
:ref:`spider middleware <topics-spider-middleware>`.
They are yielded by the :meth:`~scrapy.Spider.start` method, and may be
modified by :ref:`spider middlewares <topics-spider-middleware>`.
They are not necessarily the *first* requests sent, and may not be sent in
order; reaching :setting:`CONCURRENT_REQUESTS` and :ref:`scheduling
<topics-scheduler>` start requests is prioritized by default. But you can
change that.
..
The request send order when all start requests and callback requests have
the same priority is rather unintuitive:
1. First, the first CONCURRENT_REQUESTS start requests are sent in order.
Awaiting slow operations in Spider.start() can lower that.
2. Then, assuming an even domain distribution in start requests (i.e.
ABCABC, not AABBCC), the last N start requests are sent in reverse
order, where N is:
min(CONCURRENT_REQUESTS, CONCURRENT_REQUESTS_PER_DOMAIN * domain_count)
3. Finally, the remaining start requests are also sent in reverse order,
but only when there are not enough pending requests yielded from
callbacks to reach the configured concurrency.
The reverse order is because the scheduler uses a LIFO queue by default
(SCHEDULER_MEMORY_QUEUE, SCHEDULER_DISK_QUEUE). The order of the first few
requests is unnaffected because they are sent as soon as they are
scheduled. The last start requests sent before callback requests are those
that can be sent before the first callback requests are scheduled.
We do not document this behavior, so that we may change it in the future
without breaking the contract.
.. _start-requests-order:
Start request order
-------------------
To force a specific request order, override the :meth:`~scrapy.Spider.start`
method to set :attr:`Request.priority <scrapy.http.Request.priority>`.
For example:
- To send start requests before other requests:
.. code-block:: python
async def start(self):
async for item_or_request in super().start():
if isinstance(item_or_request, Request):
item_or_request = item_or_request.replace(priority=1)
yield item_or_request
- To send start requests in yield order:
.. code-block:: python
async def start(self):
priority = len(self.start_urls)
async for item_or_request in super().start():
if isinstance(item_or_request, Request):
item_or_request = item_or_request.replace(priority=priority)
yield item_or_request
priority -= 1
You can also :ref:`customize the scheduler <topics-scheduler>` if you need
more control over request prioritization.
.. note:: By default, :ref:`the first few requests are sent in yield order
<start-requests-front-load>` regardless, for performance reasons.
.. seealso:: :ref:`start-request-order`
.. _start-requests-front-load:
@ -452,7 +383,7 @@ Start request front loading
By default, the first few requests yielded by :meth:`~scrapy.Spider.start` are
sent in yield order, even if you :ref:`try to sort them differently
<start-requests-order>`. This is because, until :setting:`CONCURRENT_REQUESTS`
<start-request-order>`. This is because, until :setting:`CONCURRENT_REQUESTS`
is reached, start requests are removed from the scheduler immediately after
they are scheduled.

View File

@ -1,6 +1,6 @@
[build-system]
requires = ["setuptools >= 61.0"]
build-backend = "setuptools.build_meta"
requires = ["hatchling>=1.27.0"]
build-backend = "hatchling.build"
[project]
name = "Scrapy"
@ -10,29 +10,28 @@ dependencies = [
"Twisted>=21.7.0",
"cryptography>=37.0.0",
"cssselect>=0.9.1",
"defusedxml>=0.7.1",
"itemadapter>=0.1.0",
"itemloaders>=1.0.1",
"lxml>=4.9.3",
"packaging",
"parsel>=1.5.0",
"protego>=0.1.15",
"pyOpenSSL>=22.0.0",
"queuelib>=1.4.2",
"service_identity>=18.1.0",
"tldextract",
"w3lib>=1.17.0",
"zope.interface>=5.1.0",
"protego>=0.1.15",
"itemadapter>=0.1.0",
"packaging",
"tldextract",
"lxml>=4.9.3",
"defusedxml>=0.7.1",
# Platform-specific dependencies
'PyDispatcher>=2.0.5; platform_python_implementation == "CPython"',
'PyPyDispatcher>=2.1.0; platform_python_implementation == "PyPy"',
]
classifiers = [
"Framework :: Scrapy",
"Development Status :: 5 - Production/Stable",
"Environment :: Console",
"Framework :: Scrapy",
"Intended Audience :: Developers",
"License :: OSI Approved :: BSD License",
"Operating System :: OS Independent",
"Programming Language :: Python",
"Programming Language :: Python :: 3",
@ -47,6 +46,8 @@ classifiers = [
"Topic :: Software Development :: Libraries :: Application Frameworks",
"Topic :: Software Development :: Libraries :: Python Modules",
]
license = "BSD-3-Clause"
license-files = ["LICENSE", "AUTHORS"]
readme = "README.rst"
requires-python = ">=3.9"
authors = [{ name = "Scrapy developers", email = "pablo@pablohoffman.com" }]
@ -63,12 +64,26 @@ releasenotes = "https://docs.scrapy.org/en/latest/news.html"
[project.scripts]
scrapy = "scrapy.cmdline:execute"
[tool.setuptools.packages.find]
where = ["."]
include = ["scrapy", "scrapy.*",]
[tool.hatch.build.targets.sdist]
include = [
"/docs",
"/extras",
"/scrapy",
"/tests",
"/tests_typing",
"/CODE_OF_CONDUCT.md",
"/CONTRIBUTING.md",
"/INSTALL.md",
"/NEWS",
"/SECURITY.md",
"/codecov.yml",
"/conftest.py",
"/tox.ini",
]
[tool.setuptools.dynamic]
version = {file = "./scrapy/VERSION"}
[tool.hatch.version]
path = "scrapy/VERSION"
pattern = "^(?P<version>.+)$"
[tool.mypy]
ignore_missing_imports = true

View File

@ -12,12 +12,12 @@ from time import time
from traceback import format_exc
from typing import TYPE_CHECKING, Any, TypeVar, cast
from twisted.internet.defer import Deferred, succeed
from twisted.internet.defer import Deferred, inlineCallbacks, succeed
from twisted.python.failure import Failure
from scrapy import signals
from scrapy.core.scheduler import BaseScheduler
from scrapy.core.scraper import Scraper, _HandleOutputDeferred
from scrapy.core.scraper import Scraper
from scrapy.exceptions import CloseSpider, DontCloseSpider, IgnoreRequest
from scrapy.http import Request, Response
from scrapy.utils.defer import (
@ -30,7 +30,7 @@ from scrapy.utils.python import global_object_name
from scrapy.utils.reactor import CallLaterOnce
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Callable
from collections.abc import AsyncIterator, Callable, Generator
from scrapy.core.downloader import Downloader
from scrapy.crawler import Crawler
@ -344,11 +344,10 @@ class ExecutionEngine:
)
return True
@inlineCallbacks
def _handle_downloader_output(
self, result: Request | Response | Failure, request: Request
) -> _HandleOutputDeferred | None:
assert self.spider is not None # typing
) -> Generator[Deferred[Any], Any, None]:
if not isinstance(result, (Request, Response, Failure)):
raise TypeError(
f"Incorrect type: expected Request, Response or Failure, got {type(result)}: {result!r}"
@ -357,17 +356,17 @@ class ExecutionEngine:
# downloader middleware can return requests (for example, redirects)
if isinstance(result, Request):
self.crawl(result)
return None
return
d = self.scraper.enqueue_scrape(result, request, self.spider)
d.addErrback(
lambda f: logger.error(
"Error while enqueuing downloader output",
exc_info=failure_to_exc_info(f),
try:
yield self.scraper.enqueue_scrape(result, request)
except Exception:
assert self.spider is not None
logger.error(
"Error while enqueuing scrape",
exc_info=True,
extra={"spider": self.spider},
)
)
return d
def spider_is_idle(self) -> bool:
if self._slot is None:
@ -387,20 +386,20 @@ class ExecutionEngine:
"""Inject the request into the spider <-> downloader pipeline"""
if self.spider is None:
raise RuntimeError(f"No open spider to crawl: {request}")
self._schedule_request(request, self.spider)
self._schedule_request(request)
self._slot.nextcall.schedule() # type: ignore[union-attr]
def _schedule_request(self, request: Request, spider: Spider) -> None:
assert self.scheduler is not None # typing
def _schedule_request(self, request: Request) -> None:
request_scheduled_result = self.signals.send_catch_log(
signals.request_scheduled,
request=request,
spider=spider,
spider=self.spider,
dont_log=IgnoreRequest,
)
for handler, result in request_scheduled_result:
if isinstance(result, Failure) and isinstance(result.value, IgnoreRequest):
return
assert self.scheduler is not None
try:
request_was_enqueued = self.scheduler.enqueue_request(request)
except Exception as exception:
@ -411,7 +410,7 @@ class ExecutionEngine:
request_was_enqueued = False
if not request_was_enqueued:
self.signals.send_catch_log(
signals.request_dropped, request=request, spider=spider
signals.request_dropped, request=request, spider=self.spider
)
def download(self, request: Request) -> Deferred[Response]:
@ -481,9 +480,7 @@ class ExecutionEngine:
nextcall = CallLaterOnce(self._start_scheduled_requests)
self.scheduler = build_from_crawler(self.scheduler_cls, self.crawler)
self._slot = _Slot(close_if_idle, nextcall)
self._start = await maybe_deferred_to_future(
self.scraper.spidermw.process_start(spider)
)
self._start = await self.scraper.spidermw.process_start(spider)
if hasattr(self.scheduler, "open") and (d := self.scheduler.open(spider)):
await maybe_deferred_to_future(d)
await maybe_deferred_to_future(self.scraper.open_spider(spider))
@ -543,7 +540,7 @@ class ExecutionEngine:
dfd.addBoth(lambda _: self.downloader.close())
dfd.addErrback(log_failure("Downloader close failure"))
dfd.addBoth(lambda _: self.scraper.close_spider(spider))
dfd.addBoth(lambda _: self.scraper.close_spider())
dfd.addErrback(log_failure("Scraper close failure"))
if hasattr(self.scheduler, "close"):

View File

@ -5,13 +5,16 @@ import logging
from abc import abstractmethod
from pathlib import Path
from typing import TYPE_CHECKING, Any, cast
from warnings import warn
# working around https://github.com/sphinx-doc/sphinx/issues/10400
from twisted.internet.defer import Deferred # noqa: TC002
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.spiders import Spider # noqa: TC001
from scrapy.utils.job import job_dir
from scrapy.utils.misc import build_from_crawler, load_object
from scrapy.utils.python import global_object_name
if TYPE_CHECKING:
# requires queuelib >= 1.6.2
@ -117,13 +120,101 @@ class BaseScheduler(metaclass=BaseSchedulerMeta):
class Scheduler(BaseScheduler):
"""Default :ref:`scheduler <topics-scheduler>`.
Requests are stored in memory by default. Set :setting:`JOBDIR` to switch
to disk storage.
Requests are stored into priority queues
(:setting:`SCHEDULER_PRIORITY_QUEUE`) that sort requests by
:attr:`~scrapy.http.Request.priority`.
Requests are dropped if :attr:`~scrapy.Request.dont_filter` is ``False``
and :setting:`DUPEFILTER_CLASS` flags them as duplicate requests.
By default, a single, memory-based priority queue is used for all requests.
When using :setting:`JOBDIR`, a disk-based priority queue is also created,
and only unserializable requests are stored in the memory-based priority
queue. For a given priority value, requests in memory take precedence over
requests in disk.
:setting:`SCHEDULER_PRIORITY_QUEUE` handles request prioritization.
Each priority queue stores requests in separate internal queues, one per
priority value. The memory priority queue uses
:setting:`SCHEDULER_MEMORY_QUEUE` queues, while the disk priority queue
uses :setting:`SCHEDULER_DISK_QUEUE` queues. The internal queues determine
:ref:`request order <request-order>` when requests have the same priority.
:ref:`Start requests <start-requests>` are stored into separate internal
queues by default, and :ref:`ordered differently <start-request-order>`.
Duplicate requests are filtered out with an instance of
:setting:`DUPEFILTER_CLASS` if :attr:`~scrapy.Request.dont_filter` is
``False``.
.. seealso:: :ref:`topics-jobs`
.. _request-order:
Request order
=============
With default settings, pending requests are stored in a LIFO_ queue
(:ref:`except for start requests <start-request-order>`). As a result,
crawling happens in `DFO order`_, which is usually the most convenient
crawl order. However, you can enforce :ref:`BFO <bfo>` or :ref:`a custom
order <custom-request-order>` (:ref:`except for the first few requests
<concurrency-v-order>`).
.. _LIFO: https://en.wikipedia.org/wiki/Stack_(abstract_data_type)
.. _DFO order: https://en.wikipedia.org/wiki/Depth-first_search
.. _start-request-order:
Start request order
-------------------
:ref:`Start requests <start-requests>` are sent in the order they are
yielded from :meth:`~scrapy.Spider.start`, and given the same
:attr:`~scrapy.http.Request.priority`, start requests take precedence over
other requests.
You can set :setting:`SCHEDULER_START_MEMORY_QUEUE` and
:setting:`SCHEDULER_START_DISK_QUEUE` to ``None`` to handle start requests
the same as other requests when it comes to order and priority.
.. _bfo:
Crawling in BFO order
---------------------
If you do want to crawl in `BFO order`_, you can do it by setting the
following :ref:`settings <topics-settings>`:
| :setting:`DEPTH_PRIORITY` = ``1``
| :setting:`SCHEDULER_DISK_QUEUE` = ``"scrapy.squeues.PickleFifoDiskQueue"``
| :setting:`SCHEDULER_MEMORY_QUEUE` = ``"scrapy.squeues.FifoMemoryQueue"``
.. _BFO order: https://en.wikipedia.org/wiki/Breadth-first_search
.. _custom-request-order:
Crawling in a custom order
--------------------------
You can manually set :attr:`~scrapy.http.Request.priority` on requests to
force a specific request order.
.. _concurrency-v-order:
Concurrency affects order
-------------------------
While pending requests are below the configured values of
:setting:`CONCURRENT_REQUESTS`, :setting:`CONCURRENT_REQUESTS_PER_DOMAIN`
or :setting:`CONCURRENT_REQUESTS_PER_IP`, those requests are sent
concurrently.
As a result, the first few requests of a crawl may not follow the desired
order. Lowering those settings to ``1`` enforces the desired order except
for the very first request, but it significantly slows down the crawl as a
whole.
Stats
=====
The following stats are generated:
@ -141,31 +232,8 @@ class Scheduler(BaseScheduler):
enabling :setting:`SCHEDULER_DEBUG` to log a warning message with details
about the first unserializable request, to try and figure out how to make
it serializable.
.. seealso:: :ref:`topics-jobs`
"""
def __init__(
self,
dupefilter: BaseDupeFilter,
jobdir: str | None = None,
dqclass: type[BaseQueue] | None = None,
mqclass: type[BaseQueue] | None = None,
logunser: bool = False,
stats: StatsCollector | None = None,
pqclass: type[PriorityQueueProtocol] | None = None,
crawler: Crawler | None = None,
):
self.df: BaseDupeFilter = dupefilter
self.dqdir: str | None = self._dqdir(jobdir)
self.pqclass: type[PriorityQueueProtocol] | None = pqclass
self.dqclass: type[BaseQueue] | None = dqclass
self.mqclass: type[BaseQueue] | None = mqclass
self.logunser: bool = logunser
self.stats: StatsCollector | None = stats
self.crawler: Crawler | None = crawler
self._paused = False
@classmethod
def from_crawler(cls, crawler: Crawler) -> Self:
dupefilter_cls = load_object(crawler.settings["DUPEFILTER_CLASS"])
@ -180,6 +248,79 @@ class Scheduler(BaseScheduler):
crawler=crawler,
)
def __init__(
self,
dupefilter: BaseDupeFilter,
jobdir: str | None = None,
dqclass: type[BaseQueue] | None = None,
mqclass: type[BaseQueue] | None = None,
logunser: bool = False,
stats: StatsCollector | None = None,
pqclass: type[PriorityQueueProtocol] | None = None,
crawler: Crawler | None = None,
):
"""Initialize the scheduler.
:param dupefilter: An object responsible for checking and filtering duplicate requests.
The value for the :setting:`DUPEFILTER_CLASS` setting is used by default.
:type dupefilter: :class:`scrapy.dupefilters.BaseDupeFilter` instance or similar:
any class that implements the `BaseDupeFilter` interface
:param jobdir: The path of a directory to be used for persisting the crawl's state.
The value for the :setting:`JOBDIR` setting is used by default.
See :ref:`topics-jobs`.
:type jobdir: :class:`str` or ``None``
:param dqclass: A class to be used as persistent request queue.
The value for the :setting:`SCHEDULER_DISK_QUEUE` setting is used by default.
:type dqclass: class
:param mqclass: A class to be used as non-persistent request queue.
The value for the :setting:`SCHEDULER_MEMORY_QUEUE` setting is used by default.
:type mqclass: class
:param logunser: A boolean that indicates whether or not unserializable requests should be logged.
The value for the :setting:`SCHEDULER_DEBUG` setting is used by default.
:type logunser: bool
:param stats: A stats collector object to record stats about the request scheduling process.
The value for the :setting:`STATS_CLASS` setting is used by default.
:type stats: :class:`scrapy.statscollectors.StatsCollector` instance or similar:
any class that implements the `StatsCollector` interface
:param pqclass: A class to be used as priority queue for requests.
The value for the :setting:`SCHEDULER_PRIORITY_QUEUE` setting is used by default.
:type pqclass: class
:param crawler: The crawler object corresponding to the current crawl.
:type crawler: :class:`scrapy.crawler.Crawler`
"""
self.df: BaseDupeFilter = dupefilter
self.dqdir: str | None = self._dqdir(jobdir)
self.pqclass: type[PriorityQueueProtocol] | None = pqclass
self.dqclass: type[BaseQueue] | None = dqclass
self.mqclass: type[BaseQueue] | None = mqclass
self.logunser: bool = logunser
self.stats: StatsCollector | None = stats
self.crawler: Crawler | None = crawler
self._paused = False
self._sdqclass: type[BaseQueue] | None = self._get_start_queue_cls(
crawler, "DISK"
)
self._smqclass: type[BaseQueue] | None = self._get_start_queue_cls(
crawler, "MEMORY"
)
def _get_start_queue_cls(
self, crawler: Crawler | None, queue: str
) -> type[BaseQueue] | None:
if crawler is None:
return None
cls = crawler.settings[f"SCHEDULER_START_{queue}_QUEUE"]
if not cls:
return None
return load_object(cls)
def has_pending_requests(self) -> bool:
return len(self) > 0
@ -278,12 +419,27 @@ class Scheduler(BaseScheduler):
"""Create a new priority queue instance, with in-memory storage"""
assert self.crawler
assert self.pqclass
return build_from_crawler(
self.pqclass,
self.crawler,
downstream_queue_cls=self.mqclass,
key="",
)
try:
return build_from_crawler(
self.pqclass,
self.crawler,
downstream_queue_cls=self.mqclass,
key="",
start_queue_cls=self._smqclass,
)
except TypeError:
warn(
f"The __init__ method of {global_object_name(self.pqclass)} "
f"does not support a `start_queue_cls` keyword-only "
f"parameter.",
ScrapyDeprecationWarning,
)
return build_from_crawler(
self.pqclass,
self.crawler,
downstream_queue_cls=self.mqclass,
key="",
)
def _dq(self) -> PriorityQueueProtocol:
"""Create a new priority queue instance, with disk storage"""
@ -291,13 +447,29 @@ class Scheduler(BaseScheduler):
assert self.dqdir
assert self.pqclass
state = self._read_dqs_state(self.dqdir)
q = build_from_crawler(
self.pqclass,
self.crawler,
downstream_queue_cls=self.dqclass,
key=self.dqdir,
startprios=state,
)
try:
q = build_from_crawler(
self.pqclass,
self.crawler,
downstream_queue_cls=self.dqclass,
key=self.dqdir,
startprios=state,
start_queue_cls=self._sdqclass,
)
except TypeError:
warn(
f"The __init__ method of {global_object_name(self.pqclass)} "
f"does not support a `start_queue_cls` keyword-only "
f"parameter.",
ScrapyDeprecationWarning,
)
q = build_from_crawler(
self.pqclass,
self.crawler,
downstream_queue_cls=self.dqclass,
key=self.dqdir,
startprios=state,
)
if q:
logger.info(
"Resuming crawl (%(queuesize)d requests scheduled)",

View File

@ -4,22 +4,30 @@ extracts information from them"""
from __future__ import annotations
import logging
import warnings
from collections import deque
from collections.abc import AsyncIterator, Iterator
from typing import TYPE_CHECKING, Any, TypeVar, Union, cast
from collections.abc import AsyncIterator
from typing import TYPE_CHECKING, Any, TypeVar, Union
from twisted.internet.defer import Deferred, inlineCallbacks
from twisted.internet.defer import Deferred, inlineCallbacks, maybeDeferred
from twisted.python.failure import Failure
from scrapy import Spider, signals
from scrapy.core.spidermw import SpiderMiddlewareManager
from scrapy.exceptions import CloseSpider, DropItem, IgnoreRequest
from scrapy.exceptions import (
CloseSpider,
DropItem,
IgnoreRequest,
ScrapyDeprecationWarning,
)
from scrapy.http import Request, Response
from scrapy.utils.defer import (
_defer_sleep,
aiter_errback,
defer_fail,
defer_succeed,
deferred_f_from_coro_f,
deferred_from_coro,
iter_errback,
maybe_deferred_to_future,
parallel,
parallel_async,
)
@ -40,9 +48,7 @@ logger = logging.getLogger(__name__)
_T = TypeVar("_T")
_ParallelResult = list[tuple[bool, Iterator[Any]]]
_HandleOutputDeferred = Deferred[Union[_ParallelResult, None]]
QueueTuple = tuple[Union[Response, Failure], Request, _HandleOutputDeferred]
QueueTuple = tuple[Union[Response, Failure], Request, Deferred[None]]
class Slot:
@ -60,8 +66,9 @@ class Slot:
def add_response_request(
self, result: Response | Failure, request: Request
) -> _HandleOutputDeferred:
deferred: _HandleOutputDeferred = Deferred()
) -> Deferred[None]:
# this Deferred will be awaited in enqueue_scrape()
deferred: Deferred[None] = Deferred()
self.queue.append((result, request, deferred))
if isinstance(result, Response):
self.active_size += max(len(result.body), self.MIN_RESPONSE_SIZE)
@ -70,9 +77,9 @@ class Slot:
return deferred
def next_response_request_deferred(self) -> QueueTuple:
response, request, deferred = self.queue.popleft()
result, request, deferred = self.queue.popleft()
self.active.add(request)
return response, request, deferred
return result, request, deferred
def finish_response(self, result: Response | Failure, request: Request) -> None:
self.active.remove(request)
@ -110,145 +117,192 @@ class Scraper:
self.slot = Slot(self.crawler.settings.getint("SCRAPER_SLOT_MAX_ACTIVE_SIZE"))
yield self.itemproc.open_spider(spider)
def close_spider(self, spider: Spider) -> Deferred[Spider]:
def close_spider(self, spider: Spider | None = None) -> Deferred[Spider]:
"""Close a spider being scraped and release its resources"""
if spider is not None:
warnings.warn(
"Passing a 'spider' argument to Scraper.close_spider() is deprecated.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
if self.slot is None:
raise RuntimeError("Scraper slot not assigned")
self.slot.closing = Deferred()
self.slot.closing.addCallback(self.itemproc.close_spider)
self._check_if_closing(spider)
self._check_if_closing()
return self.slot.closing
def is_idle(self) -> bool:
"""Return True if there isn't any more spiders to process"""
return not self.slot
def _check_if_closing(self, spider: Spider) -> None:
def _check_if_closing(self) -> None:
assert self.slot is not None # typing
assert self.crawler.spider
if self.slot.closing and self.slot.is_idle():
self.slot.closing.callback(spider)
assert self.crawler.spider
self.slot.closing.callback(self.crawler.spider)
@inlineCallbacks
def enqueue_scrape(
self, result: Response | Failure, request: Request, spider: Spider
) -> _HandleOutputDeferred:
self, result: Response | Failure, request: Request, spider: Spider | None = None
) -> Generator[Deferred[Any], Any, None]:
if spider is not None:
warnings.warn(
"Passing a 'spider' argument to Scraper.enqueue_scrape() is deprecated.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
if self.slot is None:
raise RuntimeError("Scraper slot not assigned")
dfd = self.slot.add_response_request(result, request)
def finish_scraping(_: _T) -> _T:
assert self.slot is not None
self.slot.finish_response(result, request)
self._check_if_closing(spider)
self._scrape_next(spider)
return _
dfd.addBoth(finish_scraping)
dfd.addErrback(
lambda f: logger.error(
self._scrape_next()
try:
yield dfd
except Exception:
logger.error(
"Scraper bug processing %(request)s",
{"request": request},
exc_info=failure_to_exc_info(f),
extra={"spider": spider},
exc_info=True,
extra={"spider": self.crawler.spider},
)
)
self._scrape_next(spider)
return dfd
finally:
self.slot.finish_response(result, request)
self._check_if_closing()
self._scrape_next()
def _scrape_next(self, spider: Spider) -> None:
def _scrape_next(self) -> None:
assert self.slot is not None # typing
while self.slot.queue:
response, request, deferred = self.slot.next_response_request_deferred()
self._scrape(response, request, spider).chainDeferred(deferred)
result, request, deferred = self.slot.next_response_request_deferred()
self._scrape(result, request).chainDeferred(deferred)
def _scrape(
self, result: Response | Failure, request: Request, spider: Spider
) -> _HandleOutputDeferred:
"""
Handle the downloaded response or failure through the spider callback/errback
"""
@deferred_f_from_coro_f
async def _scrape(self, result: Response | Failure, request: Request) -> None:
"""Handle the downloaded response or failure through the spider callback/errback."""
if not isinstance(result, (Response, Failure)):
raise TypeError(
f"Incorrect type: expected Response or Failure, got {type(result)}: {result!r}"
)
dfd: Deferred[Iterable[Any] | AsyncIterator[Any]] = self._scrape2(
result, request, spider
) # returns spider's processed output
dfd.addErrback(self.handle_spider_error, request, result, spider)
dfd2: _HandleOutputDeferred = dfd.addCallback(
self.handle_spider_output, request, cast(Response, result), spider
)
return dfd2
def _scrape2(
self, result: Response | Failure, request: Request, spider: Spider
) -> Deferred[Iterable[Any] | AsyncIterator[Any]]:
"""
Handle the different cases of request's result been a Response or a Failure
"""
assert self.crawler.spider
if isinstance(result, Response):
# Deferreds are invariant so Mutable*Chain isn't matched to *Iterable
return self.spidermw.scrape_response( # type: ignore[return-value]
self.call_spider, result, request, spider
)
# else result is a Failure
dfd = self.call_spider(result, request, spider)
dfd.addErrback(self._log_download_errors, result, request, spider)
return dfd
try:
# call the spider middlewares and the request callback with the response
output = await maybe_deferred_to_future(
self.spidermw.scrape_response(
self.call_spider, result, request, self.crawler.spider
)
)
except Exception:
self.handle_spider_error(Failure(), request, result)
else:
await self.handle_spider_output_async(output, request, result)
return
try:
# call the request errback with the downloader error
await self.call_spider_async(result, request)
except Exception as spider_exc:
# the errback didn't silence the exception
if not result.check(IgnoreRequest):
logkws = self.logformatter.download_error(
result, request, self.crawler.spider
)
logger.log(
*logformatter_adapter(logkws),
extra={"spider": self.crawler.spider},
exc_info=failure_to_exc_info(result),
)
if spider_exc is not result.value:
# the errback raised a different exception, handle it
self.handle_spider_error(Failure(), request, result)
def call_spider(
self, result: Response | Failure, request: Request, spider: Spider
self, result: Response | Failure, request: Request, spider: Spider | None = None
) -> Deferred[Iterable[Any] | AsyncIterator[Any]]:
dfd: Deferred[Any]
if spider is not None:
warnings.warn(
"Passing a 'spider' argument to Scraper.call_spider() is deprecated.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
return deferred_from_coro(self.call_spider_async(result, request))
async def call_spider_async(
self, result: Response | Failure, request: Request
) -> Iterable[Any] | AsyncIterator[Any]:
"""Call the request callback or errback with the response or failure."""
await maybe_deferred_to_future(_defer_sleep())
assert self.crawler.spider
if isinstance(result, Response):
if getattr(result, "request", None) is None:
result.request = request
assert result.request
callback = result.request.callback or spider._parse
warn_on_generator_with_return_value(spider, callback)
dfd = defer_succeed(result)
dfd.addCallbacks(
callback=callback, callbackKeywords=result.request.cb_kwargs
)
callback = result.request.callback or self.crawler.spider._parse
warn_on_generator_with_return_value(self.crawler.spider, callback)
output = callback(result, **result.request.cb_kwargs)
else: # result is a Failure
# TODO: properly type adding this attribute to a Failure
result.request = request # type: ignore[attr-defined]
dfd = defer_fail(result)
if request.errback:
warn_on_generator_with_return_value(spider, request.errback)
dfd.addErrback(request.errback)
dfd2: Deferred[Iterable[Any] | AsyncIterator[Any]] = dfd.addCallback(
iterate_spider_output
if not request.errback:
result.raiseException()
warn_on_generator_with_return_value(self.crawler.spider, request.errback)
output = request.errback(result)
if isinstance(output, Failure):
output.raiseException()
# else the errback returned actual output (like a callback),
# which needs to be passed to iterate_spider_output()
return await maybe_deferred_to_future(
maybeDeferred(iterate_spider_output, output)
)
return dfd2
def handle_spider_error(
self,
_failure: Failure,
request: Request,
response: Response | Failure,
spider: Spider,
spider: Spider | None = None,
) -> None:
"""Handle an exception raised by a spider callback or errback."""
if spider is not None:
warnings.warn(
"Passing a 'spider' argument to Scraper.handle_spider_error() is deprecated.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
assert self.crawler.spider
exc = _failure.value
if isinstance(exc, CloseSpider):
assert self.crawler.engine is not None # typing
self.crawler.engine.close_spider(spider, exc.reason or "cancelled")
self.crawler.engine.close_spider(
self.crawler.spider, exc.reason or "cancelled"
)
return
logkws = self.logformatter.spider_error(_failure, request, response, spider)
logkws = self.logformatter.spider_error(
_failure, request, response, self.crawler.spider
)
logger.log(
*logformatter_adapter(logkws),
exc_info=failure_to_exc_info(_failure),
extra={"spider": spider},
extra={"spider": self.crawler.spider},
)
self.signals.send_catch_log(
signal=signals.spider_error,
failure=_failure,
response=response,
spider=spider,
spider=self.crawler.spider,
)
assert self.crawler.stats
self.crawler.stats.inc_value("spider_exceptions/count", spider=spider)
self.crawler.stats.inc_value(
f"spider_exceptions/{_failure.value.__class__.__name__}", spider=spider
"spider_exceptions/count", spider=self.crawler.spider
)
self.crawler.stats.inc_value(
f"spider_exceptions/{_failure.value.__class__.__name__}",
spider=self.crawler.spider,
)
def handle_spider_output(
@ -256,55 +310,65 @@ class Scraper:
result: Iterable[_T] | AsyncIterator[_T],
request: Request,
response: Response,
spider: Spider,
) -> _HandleOutputDeferred:
if not result:
return defer_succeed(None)
it: Iterable[_T] | AsyncIterator[_T]
dfd: Deferred[_ParallelResult]
if isinstance(result, AsyncIterator):
it = aiter_errback(
result, self.handle_spider_error, request, response, spider
spider: Spider | None = None,
) -> Deferred[None]:
"""Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel."""
if spider is not None:
warnings.warn(
"Passing a 'spider' argument to Scraper.handle_spider_output() is deprecated.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
dfd = parallel_async(
it,
self.concurrent_items,
self._process_spidermw_output,
request,
response,
spider,
)
else:
it = iter_errback(
result, self.handle_spider_error, request, response, spider
)
dfd = parallel(
it,
self.concurrent_items,
self._process_spidermw_output,
request,
response,
spider,
)
# returning Deferred[_ParallelResult] instead of Deferred[Union[_ParallelResult, None]]
return dfd # type: ignore[return-value]
return deferred_from_coro(
self.handle_spider_output_async(result, request, response)
)
def _process_spidermw_output(
self, output: Any, request: Request, response: Response, spider: Spider
) -> Deferred[Any] | None:
async def handle_spider_output_async(
self,
result: Iterable[_T] | AsyncIterator[_T],
request: Request,
response: Response,
) -> None:
"""Pass items/requests produced by a callback to ``_process_spidermw_output()`` in parallel."""
if isinstance(result, AsyncIterator):
ait = aiter_errback(result, self.handle_spider_error, request, response)
await maybe_deferred_to_future(
parallel_async(
ait,
self.concurrent_items,
self._process_spidermw_output,
response,
)
)
return
it = iter_errback(result, self.handle_spider_error, request, response)
await maybe_deferred_to_future(
parallel(
it,
self.concurrent_items,
self._process_spidermw_output,
response,
)
)
@deferred_f_from_coro_f
async def _process_spidermw_output(self, output: Any, response: Response) -> None:
"""Process each Request/Item (given in the output parameter) returned
from the given spider
from the given spider.
Items are sent to the item pipelines, requests are scheduled.
"""
if isinstance(output, Request):
assert self.crawler.engine is not None # typing
self.crawler.engine.crawl(request=output)
elif output is None:
pass
else:
return self.start_itemproc(output, response=response)
return None
return
if output is not None:
await maybe_deferred_to_future(
self.start_itemproc(output, response=response)
)
def start_itemproc(self, item: Any, *, response: Response | None) -> Deferred[Any]:
@deferred_f_from_coro_f
async def start_itemproc(self, item: Any, *, response: Response | None) -> None:
"""Send *item* to the item pipelines for processing.
*response* is the source of the item data. If the item does not come
@ -313,85 +377,56 @@ class Scraper:
assert self.slot is not None # typing
assert self.crawler.spider is not None # typing
self.slot.itemproc_size += 1
dfd = self.itemproc.process_item(item, self.crawler.spider)
dfd.addBoth(self._itemproc_finished, item, response, self.crawler.spider)
return dfd
def _log_download_errors(
self,
spider_failure: Failure,
download_failure: Failure,
request: Request,
spider: Spider,
) -> Failure | None:
"""Log and silence errors that come from the engine (typically download
errors that got propagated thru here).
spider_failure: the value passed into the errback of self.call_spider()
download_failure: the value passed into _scrape2() from
ExecutionEngine._handle_downloader_output() as "result"
"""
if not download_failure.check(IgnoreRequest):
if download_failure.frames:
logkws = self.logformatter.download_error(
download_failure, request, spider
)
try:
output = await maybe_deferred_to_future(
self.itemproc.process_item(item, self.crawler.spider)
)
except DropItem as ex:
logkws = self.logformatter.dropped(item, ex, response, self.crawler.spider)
if logkws is not None:
logger.log(
*logformatter_adapter(logkws),
extra={"spider": spider},
exc_info=failure_to_exc_info(download_failure),
*logformatter_adapter(logkws), extra={"spider": self.crawler.spider}
)
else:
errmsg = download_failure.getErrorMessage()
if errmsg:
logkws = self.logformatter.download_error(
download_failure, request, spider, errmsg
)
logger.log(
*logformatter_adapter(logkws),
extra={"spider": spider},
)
if spider_failure is not download_failure:
return spider_failure
return None
def _itemproc_finished(
self, output: Any, item: Any, response: Response | None, spider: Spider
) -> Deferred[Any]:
"""ItemProcessor finished for the given ``item`` and returned ``output``"""
assert self.slot is not None # typing
self.slot.itemproc_size -= 1
if isinstance(output, Failure):
ex = output.value
if isinstance(ex, DropItem):
logkws = self.logformatter.dropped(item, ex, response, spider)
if logkws is not None:
logger.log(*logformatter_adapter(logkws), extra={"spider": spider})
return self.signals.send_catch_log_deferred(
await maybe_deferred_to_future(
self.signals.send_catch_log_deferred(
signal=signals.item_dropped,
item=item,
response=response,
spider=spider,
exception=output.value,
spider=self.crawler.spider,
exception=ex,
)
assert ex
logkws = self.logformatter.item_error(item, ex, response, spider)
)
except Exception as ex:
logkws = self.logformatter.item_error(
item, ex, response, self.crawler.spider
)
logger.log(
*logformatter_adapter(logkws),
extra={"spider": spider},
exc_info=failure_to_exc_info(output),
extra={"spider": self.crawler.spider},
exc_info=True,
)
return self.signals.send_catch_log_deferred(
signal=signals.item_error,
item=item,
response=response,
spider=spider,
failure=output,
await maybe_deferred_to_future(
self.signals.send_catch_log_deferred(
signal=signals.item_error,
item=item,
response=response,
spider=self.crawler.spider,
failure=Failure(),
)
)
logkws = self.logformatter.scraped(output, response, spider)
if logkws is not None:
logger.log(*logformatter_adapter(logkws), extra={"spider": spider})
return self.signals.send_catch_log_deferred(
signal=signals.item_scraped, item=output, response=response, spider=spider
)
else:
logkws = self.logformatter.scraped(output, response, self.crawler.spider)
if logkws is not None:
logger.log(
*logformatter_adapter(logkws), extra={"spider": self.crawler.spider}
)
await maybe_deferred_to_future(
self.signals.send_catch_log_deferred(
signal=signals.item_scraped,
item=output,
response=response,
spider=self.crawler.spider,
)
)
finally:
self.slot.itemproc_size -= 1

View File

@ -41,7 +41,8 @@ logger = logging.getLogger(__name__)
_T = TypeVar("_T")
ScrapeFunc = Callable[
[Union[Response, Failure], Request, Spider], Union[Iterable[_T], AsyncIterator[_T]]
[Union[Response, Failure], Request],
Deferred[Union[Iterable[_T], AsyncIterator[_T]]],
]
@ -77,16 +78,15 @@ class SpiderMiddlewareManager(MiddlewareManager):
]
if deprecated_middlewares and modern_middlewares:
raise ValueError(
f"You are trying to combine spider middlewares that only "
f"define the deprecated process_start_requests() method "
f"({deprecated_middlewares}) with spider middlewares that "
f"only define the process_start() method "
f"({modern_middlewares}). This is not possible. You must "
f"either disable or make universal 1 of those 2 sets of "
f"spider middlewares. Making a spider middleware universal "
f"means having it define both methods. See the release notes "
f"of Scrapy VERSION for details: "
f"https://docs.scrapy.org/en/VERSION/news.html"
"You are trying to combine spider middlewares that only "
"define the deprecated process_start_requests() method () "
"with spider middlewares that only define the "
"process_start() method (). This is not possible. You must "
"either disable or make universal 1 of those 2 sets of "
"spider middlewares. Making a spider middleware universal "
"means having it define both methods. See the release notes "
"of Scrapy VERSION for details: "
"https://docs.scrapy.org/en/VERSION/news.html"
)
self._use_start_requests = bool(deprecated_middlewares)
@ -137,7 +137,7 @@ class SpiderMiddlewareManager(MiddlewareManager):
response: Response,
request: Request,
spider: Spider,
) -> Iterable[_T] | AsyncIterator[_T]:
) -> Deferred[Iterable[_T] | AsyncIterator[_T]]:
for method in self.methods["process_spider_input"]:
method = cast(Callable, method)
try:
@ -151,8 +151,8 @@ class SpiderMiddlewareManager(MiddlewareManager):
except _InvalidOutput:
raise
except Exception:
return scrape_func(Failure(), request, spider)
return scrape_func(response, request, spider)
return scrape_func(Failure(), request)
return scrape_func(response, request)
def _evaluate_iterable(
self,
@ -388,20 +388,18 @@ class SpiderMiddlewareManager(MiddlewareManager):
dfd2.addErrback(process_spider_exception)
return dfd2
@deferred_f_from_coro_f
async def process_start(self, spider: Spider) -> AsyncIterator[Any] | None:
try:
self._check_deprecated_start_requests_use(spider)
except ValueError as exception:
logger.error(exception)
return None
start: AsyncIterator[Any]
if self._use_start_requests:
sync_start = iter(spider.start_requests())
sync_start = await maybe_deferred_to_future(
self._process_chain("process_start_requests", sync_start, spider)
)
start = as_async_generator(sync_start)
start: AsyncIterator[Any] = as_async_generator(sync_start)
else:
error_found = False
for fn in (spider.start, *self.methods["process_start"]):

View File

@ -160,8 +160,16 @@ class ScrapyPriorityQueue:
downstream_queue_cls: type[QueueProtocol],
key: str,
startprios: Iterable[int] = (),
*,
start_queue_cls: type[QueueProtocol] | None = None,
) -> Self:
return cls(crawler, downstream_queue_cls, key, startprios)
return cls(
crawler,
downstream_queue_cls,
key,
startprios,
start_queue_cls=start_queue_cls,
)
def __init__(
self,
@ -169,11 +177,15 @@ class ScrapyPriorityQueue:
downstream_queue_cls: type[QueueProtocol],
key: str,
startprios: Iterable[int] = (),
*,
start_queue_cls: type[QueueProtocol] | None = None,
):
self.crawler: Crawler = crawler
self.downstream_queue_cls: type[QueueProtocol] = downstream_queue_cls
self._start_queue_cls: type[QueueProtocol] | None = start_queue_cls
self.key: str = key
self.queues: dict[int, QueueProtocol] = {}
self._start_queues: dict[int, QueueProtocol] = {}
self.curprio: int | None = None
self.init_prios(startprios)
@ -182,7 +194,13 @@ class ScrapyPriorityQueue:
return
for priority in startprios:
self.queues[priority] = self.qfactory(priority)
q = self.qfactory(priority)
if q:
self.queues[priority] = q
if self._start_queue_cls:
q = self._sqfactory(priority)
if q:
self._start_queues[priority] = q
self.curprio = min(startprios)
@ -193,29 +211,66 @@ class ScrapyPriorityQueue:
self.key + "/" + str(key),
)
def _sqfactory(self, key: int) -> QueueProtocol:
assert self._start_queue_cls is not None
return build_from_crawler(
self._start_queue_cls,
self.crawler,
f"{self.key}/{key}s",
)
def priority(self, request: Request) -> int:
return -request.priority
def push(self, request: Request) -> None:
priority = self.priority(request)
if priority not in self.queues:
self.queues[priority] = self.qfactory(priority)
q = self.queues[priority]
is_start_request = request.meta.get("is_start_request", False)
if is_start_request and self._start_queue_cls:
if priority not in self._start_queues:
self._start_queues[priority] = self._sqfactory(priority)
q = self._start_queues[priority]
else:
if priority not in self.queues:
self.queues[priority] = self.qfactory(priority)
q = self.queues[priority]
q.push(request) # this may fail (eg. serialization error)
if self.curprio is None or priority < self.curprio:
self.curprio = priority
def pop(self) -> Request | None:
if self.curprio is None:
return None
q = self.queues[self.curprio]
m = q.pop()
if not q:
del self.queues[self.curprio]
q.close()
prios = [p for p, q in self.queues.items() if q]
self.curprio = min(prios) if prios else None
return m
while self.curprio is not None:
if self._start_queues:
try:
q = self._start_queues[self.curprio]
except KeyError:
pass
else:
m = q.pop()
if not q:
del self._start_queues[self.curprio]
q.close()
return m
try:
q = self.queues[self.curprio]
except KeyError:
self._update_curprio()
else:
m = q.pop()
if not q:
del self.queues[self.curprio]
q.close()
self._update_curprio()
return m
return None
def _update_curprio(self) -> None:
prios = {
p
for queues in (self.queues, self._start_queues)
for p, q in queues.items()
if q
}
self.curprio = min(prios) if prios else None
def peek(self) -> Request | None:
"""Returns the next object to be returned by :meth:`pop`,
@ -226,19 +281,31 @@ class ScrapyPriorityQueue:
"""
if self.curprio is None:
return None
queue = self.queues[self.curprio]
try:
queue = self._start_queues[self.curprio]
except KeyError:
queue = self.queues[self.curprio]
# Protocols can't declare optional members
return cast(Request, queue.peek()) # type: ignore[attr-defined]
def close(self) -> list[int]:
active: list[int] = []
for p, q in self.queues.items():
active.append(p)
q.close()
return active
active: set[int] = set()
for queues in (self.queues, self._start_queues):
for p, q in queues.items():
active.add(p)
q.close()
return list(active)
def __len__(self) -> int:
return sum(len(x) for x in self.queues.values()) if self.queues else 0
return (
sum(
len(x)
for queues in (self.queues, self._start_queues)
for x in queues.values()
)
if self.queues or self._start_queues
else 0
)
class DownloaderInterface:
@ -281,8 +348,16 @@ class DownloaderAwarePriorityQueue:
downstream_queue_cls: type[QueueProtocol],
key: str,
startprios: dict[str, Iterable[int]] | None = None,
*,
start_queue_cls: type[QueueProtocol] | None = None,
) -> Self:
return cls(crawler, downstream_queue_cls, key, startprios)
return cls(
crawler,
downstream_queue_cls,
key,
startprios,
start_queue_cls=start_queue_cls,
)
def __init__(
self,
@ -290,6 +365,8 @@ class DownloaderAwarePriorityQueue:
downstream_queue_cls: type[QueueProtocol],
key: str,
slot_startprios: dict[str, Iterable[int]] | None = None,
*,
start_queue_cls: type[QueueProtocol] | None = None,
):
if crawler.settings.getint("CONCURRENT_REQUESTS_PER_IP") != 0:
raise ValueError(
@ -309,6 +386,7 @@ class DownloaderAwarePriorityQueue:
self._downloader_interface: DownloaderInterface = DownloaderInterface(crawler)
self.downstream_queue_cls: type[QueueProtocol] = downstream_queue_cls
self._start_queue_cls: type[QueueProtocol] | None = start_queue_cls
self.key: str = key
self.crawler: Crawler = crawler
@ -324,6 +402,7 @@ class DownloaderAwarePriorityQueue:
self.downstream_queue_cls,
self.key + "/" + _path_safe(slot),
startprios,
start_queue_cls=self._start_queue_cls,
)
def pop(self) -> Request | None:

View File

@ -305,6 +305,8 @@ SCHEDULER = "scrapy.core.scheduler.Scheduler"
SCHEDULER_DISK_QUEUE = "scrapy.squeues.PickleLifoDiskQueue"
SCHEDULER_MEMORY_QUEUE = "scrapy.squeues.LifoMemoryQueue"
SCHEDULER_PRIORITY_QUEUE = "scrapy.pqueues.ScrapyPriorityQueue"
SCHEDULER_START_DISK_QUEUE = "scrapy.squeues.PickleFifoDiskQueue"
SCHEDULER_START_MEMORY_QUEUE = "scrapy.squeues.FifoMemoryQueue"
SCRAPER_SLOT_MAX_ACTIVE_SIZE = 5000000
@ -315,6 +317,7 @@ SPIDER_MIDDLEWARES = {}
SPIDER_MIDDLEWARES_BASE = {
# Engine side
"scrapy.spidermiddlewares.start.StartSpiderMiddleware": 25,
"scrapy.spidermiddlewares.httperror.HttpErrorMiddleware": 50,
"scrapy.spidermiddlewares.referer.RefererMiddleware": 700,
"scrapy.spidermiddlewares.urllength.UrlLengthMiddleware": 800,

View File

@ -0,0 +1,108 @@
from __future__ import annotations
from typing import TYPE_CHECKING, Any
from scrapy import Request, Spider
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Iterable
# typing.Self requires Python 3.11
from typing_extensions import Self
from scrapy.crawler import Crawler
from scrapy.http import Response
class BaseSpiderMiddleware:
"""Optional base class for spider middlewares.
This class provides helper methods for asynchronous
``process_spider_output()`` and ``process_start()`` methods. Middlewares
that don't have either of these methods don't need to use this class.
You can override the
:meth:`~scrapy.spidermiddlewares.base.BaseSpiderMiddleware.get_processed_request`
method to add processing code for requests and the
:meth:`~scrapy.spidermiddlewares.base.BaseSpiderMiddleware.get_processed_item`
method to add processing code for items. These methods take a single
request or item from the spider output iterable and return a request or
item (the same or a new one), or ``None`` to remove this request or item
from the processing.
"""
def __init__(self, crawler: Crawler):
self.crawler: Crawler = crawler
@classmethod
def from_crawler(cls, crawler: Crawler) -> Self:
return cls(crawler)
def process_start_requests(
self, start: Iterable[Any], spider: Spider
) -> Iterable[Any]:
for o in start:
if (o := self._get_processed(o, None)) is not None:
yield o
async def process_start(self, start: AsyncIterator[Any]) -> AsyncIterator[Any]:
async for o in start:
if (o := self._get_processed(o, None)) is not None:
yield o
def process_spider_output(
self, response: Response, result: Iterable[Any], spider: Spider
) -> Iterable[Any]:
for o in result:
if (o := self._get_processed(o, response)) is not None:
yield o
async def process_spider_output_async(
self, response: Response, result: AsyncIterator[Any], spider: Spider
) -> AsyncIterator[Any]:
async for o in result:
if (o := self._get_processed(o, response)) is not None:
yield o
def _get_processed(self, o: Any, response: Response | None) -> Any:
if isinstance(o, Request):
return self.get_processed_request(o, response)
return self.get_processed_item(o, response)
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
"""Return a processed request from the spider output.
This method is called with a single request from the start seeds or the
spider output. It should return the same or a different request, or
``None`` to ignore it.
:param request: the input request
:type request: :class:`~scrapy.Request` object
:param response: the response being processed
:type response: :class:`~scrapy.http.Response` object or ``None`` for
start seeds
:return: the processed request or ``None``
"""
return request
def get_processed_item(self, item: Any, response: Response | None) -> Any:
"""Return a processed item from the spider output.
This method is called with a single item from the start seeds or the
spider output. It should return the same or a different item, or
``None`` to ignore it.
:param item: the input item
:type item: item object
:param response: the response being processed
:type response: :class:`~scrapy.http.Response` object or ``None`` for
start seeds
:return: the processed item or ``None``
"""
return item

View File

@ -9,7 +9,7 @@ from __future__ import annotations
import logging
from typing import TYPE_CHECKING, Any
from scrapy.http import Request, Response
from scrapy.spidermiddlewares.base import BaseSpiderMiddleware
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Iterable
@ -19,14 +19,17 @@ if TYPE_CHECKING:
from scrapy import Spider
from scrapy.crawler import Crawler
from scrapy.http import Request, Response
from scrapy.statscollectors import StatsCollector
logger = logging.getLogger(__name__)
class DepthMiddleware:
def __init__(
class DepthMiddleware(BaseSpiderMiddleware):
crawler: Crawler
def __init__( # pylint: disable=super-init-not-called
self,
maxdepth: int,
stats: StatsCollector,
@ -45,21 +48,22 @@ class DepthMiddleware:
verbose = settings.getbool("DEPTH_STATS_VERBOSE")
prio = settings.getint("DEPTH_PRIORITY")
assert crawler.stats
return cls(maxdepth, crawler.stats, verbose, prio)
o = cls(maxdepth, crawler.stats, verbose, prio)
o.crawler = crawler
return o
def process_spider_output(
self, response: Response, result: Iterable[Any], spider: Spider
) -> Iterable[Any]:
self._init_depth(response, spider)
return (r for r in result if self._filter(r, response, spider))
yield from super().process_spider_output(response, result, spider)
async def process_spider_output_async(
self, response: Response, result: AsyncIterator[Any], spider: Spider
) -> AsyncIterator[Any]:
self._init_depth(response, spider)
async for r in result:
if self._filter(r, response, spider):
yield r
async for o in super().process_spider_output_async(response, result, spider):
yield o
def _init_depth(self, response: Response, spider: Spider) -> None:
# base case (depth=0)
@ -68,9 +72,12 @@ class DepthMiddleware:
if self.verbose_stats:
self.stats.inc_value("request_depth_count/0", spider=spider)
def _filter(self, request: Any, response: Response, spider: Spider) -> bool:
if not isinstance(request, Request):
return True
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
if response is None:
# start requests
return request
depth = response.meta["depth"] + 1
request.meta["depth"] = depth
if self.prio:
@ -79,10 +86,12 @@ class DepthMiddleware:
logger.debug(
"Ignoring link (depth > %(maxdepth)d): %(requrl)s ",
{"maxdepth": self.maxdepth, "requrl": request.url},
extra={"spider": spider},
extra={"spider": self.crawler.spider},
)
return False
return None
if self.verbose_stats:
self.stats.inc_value(f"request_depth_count/{depth}", spider=spider)
self.stats.max_value("request_depth_max", depth, spider=spider)
return True
self.stats.inc_value(
f"request_depth_count/{depth}", spider=self.crawler.spider
)
self.stats.max_value("request_depth_max", depth, spider=self.crawler.spider)
return request

View File

@ -9,11 +9,11 @@ from __future__ import annotations
import logging
import re
import warnings
from typing import TYPE_CHECKING, Any
from typing import TYPE_CHECKING
from scrapy import Spider, signals
from scrapy.exceptions import ScrapyDeprecationWarning
from scrapy.http import Request, Response
from scrapy.spidermiddlewares.base import BaseSpiderMiddleware
from scrapy.utils.httpobj import urlparse_cached
warnings.warn(
@ -23,61 +23,55 @@ warnings.warn(
)
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Iterable
# typing.Self requires Python 3.11
from typing_extensions import Self
from scrapy.crawler import Crawler
from scrapy.http import Request, Response
from scrapy.statscollectors import StatsCollector
logger = logging.getLogger(__name__)
class OffsiteMiddleware:
def __init__(self, stats: StatsCollector):
class OffsiteMiddleware(BaseSpiderMiddleware):
crawler: Crawler
def __init__(self, stats: StatsCollector): # pylint: disable=super-init-not-called
self.stats: StatsCollector = stats
@classmethod
def from_crawler(cls, crawler: Crawler) -> Self:
assert crawler.stats
o = cls(crawler.stats)
o.crawler = crawler
crawler.signals.connect(o.spider_opened, signal=signals.spider_opened)
return o
def process_spider_output(
self, response: Response, result: Iterable[Any], spider: Spider
) -> Iterable[Any]:
return (r for r in result if self._filter(r, spider))
async def process_spider_output_async(
self, response: Response, result: AsyncIterator[Any], spider: Spider
) -> AsyncIterator[Any]:
async for r in result:
if self._filter(r, spider):
yield r
def _filter(self, request: Any, spider: Spider) -> bool:
if not isinstance(request, Request):
return True
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
if response is None:
# skip start requests for backward compatibility
return request
assert self.crawler.spider
if (
request.dont_filter
or request.meta.get("allow_offsite")
or self.should_follow(request, spider)
or self.should_follow(request, self.crawler.spider)
):
return True
return request
domain = urlparse_cached(request).hostname
if domain and domain not in self.domains_seen:
self.domains_seen.add(domain)
logger.debug(
"Filtered offsite request to %(domain)r: %(request)s",
{"domain": domain, "request": request},
extra={"spider": spider},
extra={"spider": self.crawler.spider},
)
self.stats.inc_value("offsite/domains", spider=spider)
self.stats.inc_value("offsite/filtered", spider=spider)
return False
self.stats.inc_value("offsite/domains", spider=self.crawler.spider)
self.stats.inc_value("offsite/filtered", spider=self.crawler.spider)
return None
def should_follow(self, request: Request, spider: Spider) -> bool:
regex = self.host_regex

View File

@ -6,7 +6,7 @@ originated it.
from __future__ import annotations
import warnings
from typing import TYPE_CHECKING, Any, cast
from typing import TYPE_CHECKING, cast
from urllib.parse import urlparse
from w3lib.url import safe_url_string
@ -14,13 +14,12 @@ from w3lib.url import safe_url_string
from scrapy import Spider, signals
from scrapy.exceptions import NotConfigured
from scrapy.http import Request, Response
from scrapy.spidermiddlewares.base import BaseSpiderMiddleware
from scrapy.utils.misc import load_object
from scrapy.utils.python import to_unicode
from scrapy.utils.url import strip_url
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Iterable
# typing.Self requires Python 3.11
from typing_extensions import Self
@ -327,8 +326,8 @@ def _load_policy_class(
return None
class RefererMiddleware:
def __init__(self, settings: BaseSettings | None = None):
class RefererMiddleware(BaseSpiderMiddleware):
def __init__(self, settings: BaseSettings | None = None): # pylint: disable=super-init-not-called
self.default_policy: type[ReferrerPolicy] = DefaultReferrerPolicy
if settings is not None:
settings_policy = _load_policy_class(settings.get("REFERRER_POLICY"))
@ -370,23 +369,16 @@ class RefererMiddleware:
cls = _load_policy_class(policy_name, warning_only=True)
return cls() if cls else self.default_policy()
def process_spider_output(
self, response: Response, result: Iterable[Any], spider: Spider
) -> Iterable[Any]:
return (self._set_referer(r, response) for r in result)
async def process_spider_output_async(
self, response: Response, result: AsyncIterator[Any], spider: Spider
) -> AsyncIterator[Any]:
async for r in result:
yield self._set_referer(r, response)
def _set_referer(self, r: Any, response: Response) -> Any:
if isinstance(r, Request):
referrer = self.policy(response, r).referrer(response.url, r.url)
if referrer is not None:
r.headers.setdefault("Referer", referrer)
return r
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
if response is None:
# start requests
return request
referrer = self.policy(response, request).referrer(response.url, request.url)
if referrer is not None:
request.headers.setdefault("Referer", referrer)
return request
def request_scheduled(self, request: Request, spider: Spider) -> None:
# check redirected request to patch "Referer" header if necessary

View File

@ -0,0 +1,31 @@
from __future__ import annotations
from typing import TYPE_CHECKING
from .base import BaseSpiderMiddleware
if TYPE_CHECKING:
from scrapy.http import Request
from scrapy.http.response import Response
class StartSpiderMiddleware(BaseSpiderMiddleware):
"""Set :reqmeta:`is_start_request`.
.. reqmeta:: is_start_request
is_start_request
----------------
:attr:`~scrapy.Request.meta` key that is set to ``True`` in :ref:`start
requests <start-requests>`, allowing you to tell start requests apart from
other requests, e.g. in :ref:`downloader middlewares
<topics-downloader-middleware>`.
"""
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
if response is None:
request.meta.setdefault("is_start_request", True)
return request

View File

@ -7,72 +7,49 @@ See documentation in docs/topics/spider-middleware.rst
from __future__ import annotations
import logging
import warnings
from typing import TYPE_CHECKING, Any
from typing import TYPE_CHECKING
from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning
from scrapy.http import Request, Response
from scrapy.exceptions import NotConfigured
from scrapy.spidermiddlewares.base import BaseSpiderMiddleware
if TYPE_CHECKING:
from collections.abc import AsyncIterator, Iterable
# typing.Self requires Python 3.11
from typing_extensions import Self
from scrapy import Spider
from scrapy.crawler import Crawler
from scrapy.settings import BaseSettings
from scrapy.http import Request, Response
logger = logging.getLogger(__name__)
class UrlLengthMiddleware:
def __init__(self, maxlength: int):
class UrlLengthMiddleware(BaseSpiderMiddleware):
crawler: Crawler
def __init__(self, maxlength: int): # pylint: disable=super-init-not-called
self.maxlength: int = maxlength
@classmethod
def from_settings(cls, settings: BaseSettings) -> Self:
warnings.warn(
f"{cls.__name__}.from_settings() is deprecated, use from_crawler() instead.",
category=ScrapyDeprecationWarning,
stacklevel=2,
)
return cls._from_settings(settings)
@classmethod
def from_crawler(cls, crawler: Crawler) -> Self:
return cls._from_settings(crawler.settings)
@classmethod
def _from_settings(cls, settings: BaseSettings) -> Self:
maxlength = settings.getint("URLLENGTH_LIMIT")
maxlength = crawler.settings.getint("URLLENGTH_LIMIT")
if not maxlength:
raise NotConfigured
return cls(maxlength)
o = cls(maxlength)
o.crawler = crawler
return o
def process_spider_output(
self, response: Response, result: Iterable[Any], spider: Spider
) -> Iterable[Any]:
return (r for r in result if self._filter(r, spider))
async def process_spider_output_async(
self, response: Response, result: AsyncIterator[Any], spider: Spider
) -> AsyncIterator[Any]:
async for r in result:
if self._filter(r, spider):
yield r
def _filter(self, request: Any, spider: Spider) -> bool:
if isinstance(request, Request) and len(request.url) > self.maxlength:
logger.info(
"Ignoring link (url length > %(maxlength)d): %(url)s ",
{"maxlength": self.maxlength, "url": request.url},
extra={"spider": spider},
)
assert spider.crawler.stats
spider.crawler.stats.inc_value(
"urllength/request_ignored_count", spider=spider
)
return False
return True
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
if len(request.url) <= self.maxlength:
return request
logger.info(
"Ignoring link (url length > %(maxlength)d): %(url)s ",
{"maxlength": self.maxlength, "url": request.url},
extra={"spider": self.crawler.spider},
)
assert self.crawler.stats
self.crawler.stats.inc_value(
"urllength/request_ignored_count", spider=self.crawler.spider
)
return None

View File

@ -177,8 +177,6 @@ PickleFifoDiskQueue = _scrapy_serialization_queue(_PickleFifoSerializationDiskQu
#: LIFO_ disk queue that serializes :ref:`requests <request>` using
#: :mod:`pickle`.
#:
#: .. _LIFO: https://en.wikipedia.org/wiki/LIFO_(computing)
PickleLifoDiskQueue = _scrapy_serialization_queue(_PickleLifoSerializationDiskQueue)
#: FIFO_ disk queue that serializes :ref:`requests <request>` using

View File

@ -14,7 +14,11 @@ from types import CoroutineType
from typing import TYPE_CHECKING, Any, Generic, TypeVar, Union, cast, overload
from twisted.internet import defer
from twisted.internet.defer import Deferred, DeferredList, ensureDeferred
from twisted.internet.defer import (
Deferred,
DeferredList,
ensureDeferred,
)
from twisted.internet.task import Cooperator
from twisted.python import failure
@ -36,6 +40,9 @@ _T = TypeVar("_T")
_T2 = TypeVar("_T2")
_DEFER_DELAY = 0.1
def defer_fail(_failure: Failure) -> Deferred[Any]:
"""Same as twisted.internet.defer.fail but delay calling errback until
next reactor loop
@ -46,7 +53,7 @@ def defer_fail(_failure: Failure) -> Deferred[Any]:
from twisted.internet import reactor
d: Deferred[Any] = Deferred()
reactor.callLater(0.1, d.errback, _failure)
reactor.callLater(_DEFER_DELAY, d.errback, _failure)
return d
@ -60,7 +67,16 @@ def defer_succeed(result: _T) -> Deferred[_T]:
from twisted.internet import reactor
d: Deferred[_T] = Deferred()
reactor.callLater(0.1, d.callback, result)
reactor.callLater(_DEFER_DELAY, d.callback, result)
return d
def _defer_sleep() -> Deferred[None]:
"""Like ``defer_succeed`` and ``defer_fail`` but doesn't call any real callbacks."""
from twisted.internet import reactor
d: Deferred[None] = Deferred()
reactor.callLater(_DEFER_DELAY, d.callback, None)
return d
@ -338,7 +354,7 @@ async def aiter_errback(
**kw: _P.kwargs,
) -> AsyncIterator[_T]:
"""Wraps an async iterable calling an errback if an error is caught while
iterating it. Similar to scrapy.utils.defer.iter_errback()
iterating it. Similar to :func:`scrapy.utils.defer.iter_errback`.
"""
it = aiterable.__aiter__()
while True:

View File

@ -1,3 +1,5 @@
import sys
from twisted.internet.defer import Deferred
import scrapy
@ -14,7 +16,7 @@ class SleepingSpider(scrapy.Spider):
from twisted.internet import reactor
d = Deferred()
reactor.callLater(int(self.sleep), d.callback, None)
reactor.callLater(int(sys.argv[1]), d.callback, None)
await maybe_deferred_to_future(d)

View File

@ -1,3 +1,5 @@
from __future__ import annotations
import json
import logging
import unittest
@ -14,7 +16,7 @@ from twisted.trial.unittest import TestCase
from scrapy import signals
from scrapy.crawler import CrawlerRunner
from scrapy.exceptions import StopDownload
from scrapy.exceptions import CloseSpider, StopDownload
from scrapy.http import Request
from scrapy.http.response import Response
from scrapy.utils.python import to_unicode
@ -186,11 +188,18 @@ class TestCrawl(TestCase):
@defer.inlineCallbacks
def test_start_items(self):
items = []
def _on_item_scraped(item):
items.append(item)
with LogCapture("scrapy", level=logging.ERROR) as log:
crawler = get_crawler(StartItemSpider)
crawler.signals.connect(_on_item_scraped, signals.item_scraped)
yield crawler.crawl(mockserver=self.mockserver)
assert len(log.records) == 0
assert items == [{"name": "test item"}]
@defer.inlineCallbacks
def test_start_unsupported_output(self):
@ -198,11 +207,20 @@ class TestCrawl(TestCase):
potentially expensive call to itemadapter.is_item(), and letting
instead things fail when ItemAdapter is actually used on the
corresponding non-item object."""
items = []
def _on_item_scraped(item):
items.append(item)
with LogCapture("scrapy", level=logging.ERROR) as log:
crawler = get_crawler(StartGoodAndBadOutput)
crawler.signals.connect(_on_item_scraped, signals.item_scraped)
yield crawler.crawl(mockserver=self.mockserver)
assert len(log.records) == 0
assert len(items) == 3
assert not any(isinstance(item, Request) for item in items)
@defer.inlineCallbacks
def test_start_dupes(self):
@ -692,3 +710,100 @@ class TestCrawlSpider(TestCase):
assert crawler.spider.meta[
"failure"
].value.response.headers == crawler.spider.meta.get("headers_received")
@defer.inlineCallbacks
def test_spider_errback(self):
failures = []
def eb(failure: Failure) -> Failure:
failures.append(failure)
return failure
crawler = get_crawler(SingleRequestSpider)
with LogCapture() as log:
yield crawler.crawl(
seed=self.mockserver.url("/status?n=400"), errback_func=eb
)
assert len(failures) == 1
assert "HTTP status code is not handled or not allowed" in str(log)
assert "Spider error processing" not in str(log)
@defer.inlineCallbacks
def test_spider_errback_silence(self):
failures = []
def eb(failure: Failure) -> None:
failures.append(failure)
crawler = get_crawler(SingleRequestSpider)
with LogCapture() as log:
yield crawler.crawl(
seed=self.mockserver.url("/status?n=400"), errback_func=eb
)
assert len(failures) == 1
assert "HTTP status code is not handled or not allowed" not in str(log)
assert "Spider error processing" not in str(log)
@defer.inlineCallbacks
def test_spider_errback_exception(self):
def eb(failure: Failure) -> None:
raise ValueError("foo")
crawler = get_crawler(SingleRequestSpider)
with LogCapture() as log:
yield crawler.crawl(
seed=self.mockserver.url("/status?n=400"), errback_func=eb
)
assert "Spider error processing" in str(log)
@defer.inlineCallbacks
def test_spider_errback_downloader_error(self):
failures = []
def eb(failure: Failure) -> Failure:
failures.append(failure)
return failure
crawler = get_crawler(SingleRequestSpider)
with LogCapture() as log:
yield crawler.crawl(
seed=self.mockserver.url("/drop?abort=1"), errback_func=eb
)
assert len(failures) == 1
assert "Error downloading" in str(log)
assert "Spider error processing" not in str(log)
@defer.inlineCallbacks
def test_spider_errback_exception_downloader_error(self):
def eb(failure: Failure) -> None:
raise ValueError("foo")
crawler = get_crawler(SingleRequestSpider)
with LogCapture() as log:
yield crawler.crawl(
seed=self.mockserver.url("/drop?abort=1"), errback_func=eb
)
assert "Error downloading" in str(log)
assert "Spider error processing" in str(log)
@defer.inlineCallbacks
def test_raise_closespider(self):
def cb(response):
raise CloseSpider
crawler = get_crawler(SingleRequestSpider)
with LogCapture() as log:
yield crawler.crawl(seed=self.mockserver.url("/"), callback_func=cb)
assert "Closing spider (cancelled)" in str(log)
assert "Spider error processing" not in str(log)
@defer.inlineCallbacks
def test_raise_closespider_reason(self):
def cb(response):
raise CloseSpider("my_reason")
crawler = get_crawler(SingleRequestSpider)
with LogCapture() as log:
yield crawler.crawl(seed=self.mockserver.url("/"), callback_func=cb)
assert "Closing spider (my_reason)" in str(log)
assert "Spider error processing" not in str(log)

View File

@ -891,7 +891,7 @@ class TestCrawlerProcessSubprocess(ScriptRunnerMixin, unittest.TestCase):
def test_shutdown_graceful(self):
sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK
args = self.get_script_args("sleeping.py", "-a", "sleep=3")
args = self.get_script_args("sleeping.py", "3")
p = PopenSpawn(args, timeout=5)
p.expect_exact("Spider opened")
p.expect_exact("Crawled (200)")
@ -905,7 +905,7 @@ class TestCrawlerProcessSubprocess(ScriptRunnerMixin, unittest.TestCase):
from twisted.internet import reactor
sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK
args = self.get_script_args("sleeping.py", "-a", "sleep=10")
args = self.get_script_args("sleeping.py", "10")
p = PopenSpawn(args, timeout=5)
p.expect_exact("Spider opened")
p.expect_exact("Crawled (200)")

View File

@ -1,22 +1,18 @@
from __future__ import annotations
import asyncio
from gzip import BadGzipFile
from unittest import mock
import pytest
from twisted.internet import defer
from twisted.internet.defer import Deferred
from twisted.python.failure import Failure
from twisted.internet.defer import Deferred, succeed
from twisted.trial.unittest import TestCase
from scrapy.core.downloader.middleware import DownloaderMiddlewareManager
from scrapy.exceptions import _InvalidOutput
from scrapy.http import Request, Response
from scrapy.spiders import Spider
from scrapy.utils.defer import (
deferred_f_from_coro_f,
deferred_to_future,
maybe_deferred_to_future,
)
from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future
from scrapy.utils.python import to_bytes
from scrapy.utils.test import get_crawler, get_from_asyncio_queue
@ -34,26 +30,22 @@ class TestManagerBase(TestCase):
def tearDown(self):
return self.crawler.engine.close_spider(self.spider)
async def _download(self, request, response=None):
async def _download(
self, request: Request, response: Response | None = None
) -> Response | Request:
"""Executes downloader mw manager's download method and returns
the result (Request or Response) or raise exception in case of
the result (Request or Response) or raises exception in case of
failure.
"""
if not response:
response = Response(request.url)
def download_func(request, spider):
return response
def download_func(request: Request, spider: Spider) -> Deferred[Response]:
return succeed(response)
dfd = self.mwman.download(download_func, request, self.spider)
# catch deferred result and return the value
results = []
dfd.addBoth(results.append)
await maybe_deferred_to_future(dfd)
ret = results[0]
if isinstance(ret, Failure):
ret.raiseException()
return ret
return await maybe_deferred_to_future(
self.mwman.download(download_func, request, self.spider)
)
class TestDefaults(TestManagerBase):
@ -92,7 +84,7 @@ class TestDefaults(TestManagerBase):
"Location": "http://example.com/login",
},
)
ret = await self._download(request=req, response=resp)
ret = await self._download(req, resp)
assert isinstance(ret, Request), f"Not redirected: {ret!r}"
assert to_bytes(ret.url) == resp.headers["Location"], (
"Not redirected to location header"
@ -114,7 +106,7 @@ class TestDefaults(TestManagerBase):
},
)
with pytest.raises(BadGzipFile):
await self._download(request=req, response=resp)
await self._download(req, resp)
class TestResponseFromProcessRequest(TestManagerBase):
@ -132,19 +124,17 @@ class TestResponseFromProcessRequest(TestManagerBase):
req = Request("http://example.com/index.html")
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
await maybe_deferred_to_future(dfd)
assert results[0] is resp
result = await maybe_deferred_to_future(
self.mwman.download(download_func, req, self.spider)
)
assert result is resp
assert not download_func.called
class TestProcessRequestInvalidOutput(TestManagerBase):
"""Invalid return value for process_request method should raise an exception"""
def test_invalid_process_request(self):
class TestInvalidOutput(TestManagerBase):
@deferred_f_from_coro_f
async def test_invalid_process_request(self):
"""Invalid return value for process_request method should raise an exception"""
req = Request("http://example.com/index.html")
class InvalidProcessRequestMiddleware:
@ -152,18 +142,12 @@ class TestProcessRequestInvalidOutput(TestManagerBase):
return 1
self.mwman._add_middleware(InvalidProcessRequestMiddleware())
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
assert isinstance(results[0], Failure)
assert isinstance(results[0].value, _InvalidOutput)
with pytest.raises(_InvalidOutput):
await self._download(req)
class TestProcessResponseInvalidOutput(TestManagerBase):
"""Invalid return value for process_response method should raise an exception"""
def test_invalid_process_response(self):
@deferred_f_from_coro_f
async def test_invalid_process_response(self):
"""Invalid return value for process_response method should raise an exception"""
req = Request("http://example.com/index.html")
class InvalidProcessResponseMiddleware:
@ -171,18 +155,12 @@ class TestProcessResponseInvalidOutput(TestManagerBase):
return 1
self.mwman._add_middleware(InvalidProcessResponseMiddleware())
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
assert isinstance(results[0], Failure)
assert isinstance(results[0].value, _InvalidOutput)
with pytest.raises(_InvalidOutput):
await self._download(req)
class TestProcessExceptionInvalidOutput(TestManagerBase):
"""Invalid return value for process_exception method should raise an exception"""
def test_invalid_process_exception(self):
@deferred_f_from_coro_f
async def test_invalid_process_exception(self):
"""Invalid return value for process_exception method should raise an exception"""
req = Request("http://example.com/index.html")
class InvalidProcessExceptionMiddleware:
@ -193,12 +171,8 @@ class TestProcessExceptionInvalidOutput(TestManagerBase):
return 1
self.mwman._add_middleware(InvalidProcessExceptionMiddleware())
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
assert isinstance(results[0], Failure)
assert isinstance(results[0].value, _InvalidOutput)
with pytest.raises(_InvalidOutput):
await self._download(req)
class TestMiddlewareUsingDeferreds(TestManagerBase):
@ -221,12 +195,10 @@ class TestMiddlewareUsingDeferreds(TestManagerBase):
self.mwman._add_middleware(DeferredMiddleware())
req = Request("http://example.com/index.html")
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
await maybe_deferred_to_future(dfd)
assert results[0] is resp
result = await maybe_deferred_to_future(
self.mwman.download(download_func, req, self.spider)
)
assert result is resp
assert not download_func.called
@ -240,18 +212,16 @@ class TestMiddlewareUsingCoro(TestManagerBase):
class CoroMiddleware:
async def process_request(self, request, spider):
await defer.succeed(42)
await succeed(42)
return resp
self.mwman._add_middleware(CoroMiddleware())
req = Request("http://example.com/index.html")
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
await maybe_deferred_to_future(dfd)
assert results[0] is resp
result = await maybe_deferred_to_future(
self.mwman.download(download_func, req, self.spider)
)
assert result is resp
assert not download_func.called
@pytest.mark.only_asyncio
@ -267,10 +237,8 @@ class TestMiddlewareUsingCoro(TestManagerBase):
self.mwman._add_middleware(CoroMiddleware())
req = Request("http://example.com/index.html")
download_func = mock.MagicMock()
dfd = self.mwman.download(download_func, req, self.spider)
results = []
dfd.addBoth(results.append)
await deferred_to_future(dfd)
assert results[0] is resp
result = await maybe_deferred_to_future(
self.mwman.download(download_func, req, self.spider)
)
assert result is resp
assert not download_func.called

View File

@ -482,8 +482,7 @@ def test_request_scheduled_signal(caplog):
if "drop" in request.url:
raise IgnoreRequest
spider = MySpider()
crawler = get_crawler(spider.__class__)
crawler = get_crawler(MySpider)
engine = ExecutionEngine(crawler, lambda _: None)
engine.downloader._slot_gc_loop.stop()
scheduler = MemoryScheduler()
@ -497,10 +496,10 @@ def test_request_scheduled_signal(caplog):
engine._slot = _Slot(False, Mock())
crawler.signals.connect(signal_handler, request_scheduled)
keep_request = Request("https://keep.example")
engine._schedule_request(keep_request, spider)
engine._schedule_request(keep_request)
drop_request = Request("https://drop.example")
caplog.set_level(DEBUG)
engine._schedule_request(drop_request, spider)
engine._schedule_request(drop_request)
assert list(scheduler.queue) == [keep_request], (
f"{list(scheduler.queue)!r} != [{keep_request!r}]"
)

View File

@ -10,6 +10,7 @@ from twisted.trial.unittest import TestCase
from scrapy import Request, Spider, signals
from scrapy.core.engine import ExecutionEngine
from scrapy.squeues import LifoMemoryQueue
from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future
from scrapy.utils.test import get_crawler
@ -208,10 +209,10 @@ class RequestSendOrderTestCase(TestCase):
def tearDownClass(cls):
cls.mockserver.__exit__(None, None, None) # increase if flaky
def request(self, num, response_seconds, download_slots=1):
def request(self, num, response_seconds, download_slots=1, priority=0):
url = self.mockserver.url(f"/delay?n={response_seconds}&{num}")
meta = {"download_slot": str(num % download_slots)}
return Request(url, meta=meta)
return Request(url, meta=meta, priority=priority)
def get_num(self, request_or_response: Request | Response):
return int(request_or_response.url.rsplit("&", maxsplit=1)[1])
@ -264,55 +265,156 @@ class RequestSendOrderTestCase(TestCase):
expected_nums = sorted(start_nums + cb_nums)
assert actual_nums == expected_nums, f"{actual_nums=} != {expected_nums=}"
@deferred_f_from_coro_f
async def test_default(self):
"""By default, start requests take priority over callback requests and
are sent in order. Priority matters, but given the same priority, a
start request takes precedence."""
nums = [1, 2, 3, 4, 5, 6]
response_seconds = 0
download_slots = 1
def _request(num, priority=0):
return self.request(
num, response_seconds, download_slots, priority=priority
)
async def start(spider):
# The first CONCURRENT_REQUESTS start requests are sent
# immediately.
yield _request(1)
for request in (
_request(4, priority=1),
_request(6),
):
spider.crawler.engine.scheduler.enqueue_request(request)
yield _request(5)
yield _request(2, priority=1)
yield _request(3, priority=1)
def parse(spider, response):
return
yield
await maybe_deferred_to_future(
self._test_request_order(
start_nums=nums,
settings={"CONCURRENT_REQUESTS": 1},
response_seconds=response_seconds,
start_fn=start,
parse_fn=parse,
)
)
@deferred_f_from_coro_f
async def test_lifo_start(self):
"""Changing the queues of start requests to LIFO, matching the queues
of non-start requests, does not cause all requests to be stored in the
same queue objects, it only affects the order of start requests."""
nums = [1, 2, 3, 4, 5, 6]
response_seconds = 0
download_slots = 1
def _request(num, priority=0):
return self.request(
num, response_seconds, download_slots, priority=priority
)
async def start(spider):
# The first CONCURRENT_REQUESTS start requests are sent
# immediately.
yield _request(1)
for request in (
_request(4, priority=1),
_request(6),
):
spider.crawler.engine.scheduler.enqueue_request(request)
yield _request(5)
yield _request(3, priority=1)
yield _request(2, priority=1)
def parse(spider, response):
return
yield
await maybe_deferred_to_future(
self._test_request_order(
start_nums=nums,
settings={
"CONCURRENT_REQUESTS": 1,
"SCHEDULER_START_MEMORY_QUEUE": "scrapy.squeues.LifoMemoryQueue",
},
response_seconds=response_seconds,
start_fn=start,
parse_fn=parse,
)
)
@deferred_f_from_coro_f
async def test_shared_queues(self):
"""If SCHEDULER_START_*_QUEUE is falsy, start requests and other
requests share the same queue, i.e. start requests are not priorized
over other requests if their priority matches."""
nums = list(range(1, 14))
response_seconds = 0
download_slots = 1
def _request(num, priority=0):
return self.request(
num, response_seconds, download_slots, priority=priority
)
async def start(spider):
# The first CONCURRENT_REQUESTS start requests are sent
# immediately.
yield _request(1)
# Below, priority 1 requests are sent first, and requests are sent
# in LIFO order.
for request in (
_request(7, priority=1),
_request(6, priority=1),
_request(13),
_request(12),
):
spider.crawler.engine.scheduler.enqueue_request(request)
yield _request(11)
yield _request(10)
yield _request(5, priority=1)
yield _request(4, priority=1)
for request in (
_request(3, priority=1),
_request(2, priority=1),
_request(9),
_request(8),
):
spider.crawler.engine.scheduler.enqueue_request(request)
def parse(spider, response):
return
yield
await maybe_deferred_to_future(
self._test_request_order(
start_nums=nums,
settings={
"CONCURRENT_REQUESTS": 1,
"SCHEDULER_START_MEMORY_QUEUE": None,
},
response_seconds=response_seconds,
start_fn=start,
parse_fn=parse,
)
)
# Examples from the “Start requests” section of the documentation about
# spiders.
@deferred_f_from_coro_f
async def test_start_requests_first(self):
start_nums = [1, 3, 2]
cb_nums = [4]
response_seconds = self.seconds
download_slots = 1
async def start(spider):
for num in start_nums:
request = self.request(num, response_seconds, download_slots)
yield request.replace(priority=1)
await maybe_deferred_to_future(
self._test_request_order(
start_nums=start_nums,
cb_nums=cb_nums,
settings={"CONCURRENT_REQUESTS": 1},
response_seconds=response_seconds,
start_fn=start,
)
)
@deferred_f_from_coro_f
async def test_start_requests_first_sorted(self):
start_nums = [1, 2, 3]
cb_nums = [4]
response_seconds = self.seconds
download_slots = 1
async def start(spider):
priority = len(start_nums)
for num in start_nums:
request = self.request(num, response_seconds, download_slots)
yield request.replace(priority=priority)
priority -= 1
await maybe_deferred_to_future(
self._test_request_order(
start_nums=start_nums,
cb_nums=cb_nums,
settings={"CONCURRENT_REQUESTS": 1},
response_seconds=response_seconds,
start_fn=start,
)
)
@deferred_f_from_coro_f
async def test_front_load(self):
start_nums = [2, 1]
@ -325,8 +427,8 @@ class RequestSendOrderTestCase(TestCase):
assert isinstance(spider.crawler.engine.scheduler, MemoryScheduler)
spider.crawler.engine.scheduler.pause()
# By pausing the scheduler, a is scheduled before b is sent,
# and since the scheduler uses a LIFO queue, a is sent first.
# By pausing the scheduler, 1 is scheduled before 2 is sent,
# and provided the scheduler uses a LIFO queue, 1 is sent first.
yield self.request(2, response_seconds, download_slots)
yield self.request(1, response_seconds, download_slots)
spider.crawler.engine.scheduler.unpause()
@ -336,6 +438,9 @@ class RequestSendOrderTestCase(TestCase):
start_nums=start_nums,
response_seconds=response_seconds,
start_fn=start,
settings={
"SCHEDULER_START_MEMORY_QUEUE": LifoMemoryQueue,
},
)
)

View File

@ -2,12 +2,12 @@ from __future__ import annotations
from collections.abc import AsyncIterator, Iterable
from inspect import isasyncgen
from typing import Any
from unittest import mock
import pytest
from testfixtures import LogCapture
from twisted.internet import defer
from twisted.python.failure import Failure
from twisted.trial.unittest import TestCase
from scrapy.core.spidermw import SpiderMiddlewareManager
@ -31,53 +31,51 @@ class TestSpiderMiddleware(TestCase):
self.spider = self.crawler._create_spider("foo")
self.mwman = SpiderMiddlewareManager.from_crawler(self.crawler)
def _scrape_response(self):
async def _scrape_response(self) -> Any:
"""Execute spider mw manager's scrape_response method and return the result.
Raise exception in case of failure.
"""
scrape_func = mock.MagicMock()
dfd = self.mwman.scrape_response(
scrape_func, self.response, self.request, self.spider
return await maybe_deferred_to_future(
self.mwman.scrape_response(
scrape_func, self.response, self.request, self.spider
)
)
# catch deferred result and return the value
results = []
dfd.addBoth(results.append)
self._wait(dfd)
return results[0]
class TestProcessSpiderInputInvalidOutput(TestSpiderMiddleware):
"""Invalid return value for process_spider_input method"""
def test_invalid_process_spider_input(self):
@deferred_f_from_coro_f
async def test_invalid_process_spider_input(self):
class InvalidProcessSpiderInputMiddleware:
def process_spider_input(self, response, spider):
return 1
self.mwman._add_middleware(InvalidProcessSpiderInputMiddleware())
result = self._scrape_response()
assert isinstance(result, Failure)
assert isinstance(result.value, _InvalidOutput)
with pytest.raises(_InvalidOutput):
await self._scrape_response()
class TestProcessSpiderOutputInvalidOutput(TestSpiderMiddleware):
"""Invalid return value for process_spider_output method"""
def test_invalid_process_spider_output(self):
@deferred_f_from_coro_f
async def test_invalid_process_spider_output(self):
class InvalidProcessSpiderOutputMiddleware:
def process_spider_output(self, response, result, spider):
return 1
self.mwman._add_middleware(InvalidProcessSpiderOutputMiddleware())
result = self._scrape_response()
assert isinstance(result, Failure)
assert isinstance(result.value, _InvalidOutput)
with pytest.raises(_InvalidOutput):
await self._scrape_response()
class TestProcessSpiderExceptionInvalidOutput(TestSpiderMiddleware):
"""Invalid return value for process_spider_exception method"""
def test_invalid_process_spider_exception(self):
@deferred_f_from_coro_f
async def test_invalid_process_spider_exception(self):
class InvalidProcessSpiderOutputExceptionMiddleware:
def process_spider_exception(self, response, exception, spider):
return 1
@ -88,15 +86,15 @@ class TestProcessSpiderExceptionInvalidOutput(TestSpiderMiddleware):
self.mwman._add_middleware(InvalidProcessSpiderOutputExceptionMiddleware())
self.mwman._add_middleware(RaiseExceptionProcessSpiderOutputMiddleware())
result = self._scrape_response()
assert isinstance(result, Failure)
assert isinstance(result.value, _InvalidOutput)
with pytest.raises(_InvalidOutput):
await self._scrape_response()
class TestProcessSpiderExceptionReRaise(TestSpiderMiddleware):
"""Re raise the exception by returning None"""
def test_process_spider_exception_return_none(self):
@deferred_f_from_coro_f
async def test_process_spider_exception_return_none(self):
class ProcessSpiderExceptionReturnNoneMiddleware:
def process_spider_exception(self, response, exception, spider):
return None
@ -107,9 +105,8 @@ class TestProcessSpiderExceptionReRaise(TestSpiderMiddleware):
self.mwman._add_middleware(ProcessSpiderExceptionReturnNoneMiddleware())
self.mwman._add_middleware(RaiseExceptionProcessSpiderOutputMiddleware())
result = self._scrape_response()
assert isinstance(result, Failure)
assert isinstance(result.value, ZeroDivisionError)
with pytest.raises(ZeroDivisionError):
await self._scrape_response()
class TestBaseAsyncSpiderMiddleware(TestSpiderMiddleware):
@ -350,7 +347,7 @@ class TestProcessStartSimple(TestBaseAsyncSpiderMiddleware):
)
self.spider = self.crawler._create_spider()
self.mwman = SpiderMiddlewareManager.from_crawler(self.crawler)
return await maybe_deferred_to_future(self.mwman.process_start(self.spider))
return await self.mwman.process_start(self.spider)
@deferred_f_from_coro_f
async def test_simple(self):

View File

@ -0,0 +1,132 @@
from __future__ import annotations
from typing import TYPE_CHECKING, Any
import pytest
from scrapy import Request, Spider
from scrapy.http import Response
from scrapy.spidermiddlewares.base import BaseSpiderMiddleware
from scrapy.utils.test import get_crawler
if TYPE_CHECKING:
from scrapy.crawler import Crawler
@pytest.fixture
def crawler() -> Crawler:
return get_crawler(Spider)
def test_trivial(crawler):
class TrivialSpiderMiddleware(BaseSpiderMiddleware):
pass
mw = TrivialSpiderMiddleware.from_crawler(crawler)
assert hasattr(mw, "crawler")
assert mw.crawler is crawler
test_req = Request("data:,")
spider_output = [test_req, {"foo": "bar"}]
for processed in [
list(
mw.process_spider_output(Response("data:,"), spider_output, crawler.spider)
),
list(mw.process_start_requests(spider_output, crawler.spider)),
]:
assert processed == [test_req, {"foo": "bar"}]
def test_processed_request(crawler):
class ProcessReqSpiderMiddleware(BaseSpiderMiddleware):
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
if request.url == "data:2,":
return None
if request.url == "data:3,":
return Request("data:30,")
return request
mw = ProcessReqSpiderMiddleware.from_crawler(crawler)
test_req1 = Request("data:1,")
test_req2 = Request("data:2,")
test_req3 = Request("data:3,")
spider_output = [test_req1, {"foo": "bar"}, test_req2, test_req3]
for processed in [
list(
mw.process_spider_output(Response("data:,"), spider_output, crawler.spider)
),
list(mw.process_start_requests(spider_output, crawler.spider)),
]:
assert len(processed) == 3
assert isinstance(processed[0], Request)
assert processed[0].url == "data:1,"
assert processed[1] == {"foo": "bar"}
assert isinstance(processed[2], Request)
assert processed[2].url == "data:30,"
def test_processed_item(crawler):
class ProcessItemSpiderMiddleware(BaseSpiderMiddleware):
def get_processed_item(self, item: Any, response: Response | None) -> Any:
if item["foo"] == 2:
return None
if item["foo"] == 3:
item["foo"] = 30
return item
mw = ProcessItemSpiderMiddleware.from_crawler(crawler)
test_req = Request("data:,")
spider_output = [{"foo": 1}, {"foo": 2}, test_req, {"foo": 3}]
for processed in [
list(
mw.process_spider_output(Response("data:,"), spider_output, crawler.spider)
),
list(mw.process_start_requests(spider_output, crawler.spider)),
]:
assert processed == [{"foo": 1}, test_req, {"foo": 30}]
def test_processed_both(crawler):
class ProcessBothSpiderMiddleware(BaseSpiderMiddleware):
def get_processed_request(
self, request: Request, response: Response | None
) -> Request | None:
if request.url == "data:2,":
return None
if request.url == "data:3,":
return Request("data:30,")
return request
def get_processed_item(self, item: Any, response: Response | None) -> Any:
if item["foo"] == 2:
return None
if item["foo"] == 3:
item["foo"] = 30
return item
mw = ProcessBothSpiderMiddleware.from_crawler(crawler)
test_req1 = Request("data:1,")
test_req2 = Request("data:2,")
test_req3 = Request("data:3,")
spider_output = [
test_req1,
{"foo": 1},
{"foo": 2},
test_req2,
{"foo": 3},
test_req3,
]
for processed in [
list(
mw.process_spider_output(Response("data:,"), spider_output, crawler.spider)
),
list(mw.process_start_requests(spider_output, crawler.spider)),
]:
assert len(processed) == 4
assert isinstance(processed[0], Request)
assert processed[0].url == "data:1,"
assert processed[1] == {"foo": 1}
assert processed[2] == {"foo": 30}
assert isinstance(processed[3], Request)
assert processed[3].url == "data:30,"

View File

@ -1,19 +1,18 @@
from scrapy.http import Request, Response
from scrapy.spidermiddlewares.depth import DepthMiddleware
from scrapy.spiders import Spider
from scrapy.statscollectors import StatsCollector
from scrapy.utils.test import get_crawler
class TestDepthMiddleware:
def setup_method(self):
crawler = get_crawler(Spider)
crawler = get_crawler(Spider, {"DEPTH_LIMIT": 1, "DEPTH_STATS_VERBOSE": True})
self.spider = crawler._create_spider("scrapytest.org")
self.stats = StatsCollector(crawler)
self.stats = crawler.stats
self.stats.open_spider(self.spider)
self.mw = DepthMiddleware(1, self.stats, True)
self.mw = DepthMiddleware.from_crawler(crawler)
def test_process_spider_output(self):
req = Request("http://scrapytest.org")

View File

@ -10,7 +10,7 @@ from scrapy.utils.test import get_crawler
class TestOffsiteMiddleware:
def setup_method(self):
crawler = get_crawler(Spider)
self.spider = crawler._create_spider(**self._get_spiderargs())
self.spider = crawler.spider = crawler._create_spider(**self._get_spiderargs())
self.mw = OffsiteMiddleware.from_crawler(crawler)
self.mw.spider_opened(self.spider)

View File

@ -0,0 +1,44 @@
from twisted.trial.unittest import TestCase
from scrapy.http import Request
from scrapy.spidermiddlewares.start import StartSpiderMiddleware
from scrapy.spiders import Spider
from scrapy.utils.defer import deferred_f_from_coro_f
from scrapy.utils.misc import build_from_crawler
from scrapy.utils.test import get_crawler
class TestMiddleware(TestCase):
@deferred_f_from_coro_f
async def test_async(self):
crawler = get_crawler(Spider)
mw = build_from_crawler(StartSpiderMiddleware, crawler)
async def start():
yield Request("data:,1")
yield Request("data:,2", meta={"is_start_request": True})
yield Request("data:,2", meta={"is_start_request": False})
yield Request("data:,2", meta={"is_start_request": "foo"})
result = [
request.meta["is_start_request"]
async for request in mw.process_start(start())
]
assert result == [True, True, False, "foo"]
@deferred_f_from_coro_f
async def test_sync(self):
crawler = get_crawler(Spider)
mw = build_from_crawler(StartSpiderMiddleware, crawler)
def start():
yield Request("data:,1")
yield Request("data:,2", meta={"is_start_request": True})
yield Request("data:,2", meta={"is_start_request": False})
yield Request("data:,2", meta={"is_start_request": "foo"})
result = [
request.meta["is_start_request"]
for request in mw.process_start_requests(start(), Spider("test"))
]
assert result == [True, True, False, "foo"]

View File

@ -143,7 +143,7 @@ deps =
google-cloud-storage
ipython
robotexclusionrulesparser
uvloop; platform_system != "Windows"
uvloop; platform_system != "Windows" and implementation_name != "pypy"
zstandard; implementation_name != "pypy" # optional for HTTP compress downloader middleware tests
[testenv:extra-deps-pinned]
@ -159,7 +159,7 @@ deps =
google-cloud-storage==1.29.0
ipython==2.0.0
robotexclusionrulesparser==1.6.2
uvloop==0.14.0; platform_system != "Windows"
uvloop==0.14.0; platform_system != "Windows" and implementation_name != "pypy"
zstandard==0.1; implementation_name != "pypy"
install_command = {[pinned]install_command}
setenv =