Change default priority queue to DownloaderAwarePriorityQueue (#6940)

* Change default priority queue to DownloaderAwarePriorityQueue

* Fix documentation building

* Simplify test_start_already_running_exception changes.

* Modernize the test.

* Fix TestEngineCloseSpider.

* Fix typing.

* Remove special slot=None handling.

---------

Co-authored-by: Andrey Rakhmatullin <wrar@wrar.name>
This commit is contained in:
Thalison Fernandes 2025-12-11 07:25:17 -03:00 committed by GitHub
parent 1a3e343dc4
commit d8583a89c7
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
5 changed files with 46 additions and 16 deletions

View File

@ -41,19 +41,6 @@ efficient broad crawl.
.. _broad-crawls-scheduler-priority-queue: .. _broad-crawls-scheduler-priority-queue:
Use the right :setting:`SCHEDULER_PRIORITY_QUEUE`
=================================================
Scrapys default scheduler priority queue is ``'scrapy.pqueues.ScrapyPriorityQueue'``.
It works best during single-domain crawl. It does not work well with crawling
many different domains in parallel
To apply the recommended priority queue use:
.. code-block:: python
SCHEDULER_PRIORITY_QUEUE = "scrapy.pqueues.DownloaderAwarePriorityQueue"
.. _broad-crawls-concurrency: .. _broad-crawls-concurrency:
Increase concurrency Increase concurrency

View File

@ -1739,10 +1739,10 @@ Type of in-memory queue used by the scheduler. Other available type is:
SCHEDULER_PRIORITY_QUEUE SCHEDULER_PRIORITY_QUEUE
------------------------ ------------------------
Default: ``'scrapy.pqueues.ScrapyPriorityQueue'`` Default: ``'scrapy.pqueues.DownloaderAwarePriorityQueue'``
Type of priority queue used by the scheduler. Another available type is Type of priority queue used by the scheduler. Another available type is
``scrapy.pqueues.DownloaderAwarePriorityQueue``. ``scrapy.pqueues.ScrapyPriorityQueue``.
``scrapy.pqueues.DownloaderAwarePriorityQueue`` works better than ``scrapy.pqueues.DownloaderAwarePriorityQueue`` works better than
``scrapy.pqueues.ScrapyPriorityQueue`` when you crawl many different ``scrapy.pqueues.ScrapyPriorityQueue`` when you crawl many different
domains in parallel. domains in parallel.

View File

@ -479,7 +479,7 @@ SCHEDULER = "scrapy.core.scheduler.Scheduler"
SCHEDULER_DEBUG = False SCHEDULER_DEBUG = False
SCHEDULER_DISK_QUEUE = "scrapy.squeues.PickleLifoDiskQueue" SCHEDULER_DISK_QUEUE = "scrapy.squeues.PickleLifoDiskQueue"
SCHEDULER_MEMORY_QUEUE = "scrapy.squeues.LifoMemoryQueue" SCHEDULER_MEMORY_QUEUE = "scrapy.squeues.LifoMemoryQueue"
SCHEDULER_PRIORITY_QUEUE = "scrapy.pqueues.ScrapyPriorityQueue" SCHEDULER_PRIORITY_QUEUE = "scrapy.pqueues.DownloaderAwarePriorityQueue"
SCHEDULER_START_DISK_QUEUE = "scrapy.squeues.PickleFifoDiskQueue" SCHEDULER_START_DISK_QUEUE = "scrapy.squeues.PickleFifoDiskQueue"
SCHEDULER_START_MEMORY_QUEUE = "scrapy.squeues.FifoMemoryQueue" SCHEDULER_START_MEMORY_QUEUE = "scrapy.squeues.FifoMemoryQueue"

View File

@ -1,11 +1,13 @@
import time import time
from typing import Any from typing import Any
import pytest
from twisted.internet.defer import inlineCallbacks from twisted.internet.defer import inlineCallbacks
from scrapy import Request from scrapy import Request
from scrapy.core.downloader import Downloader, Slot from scrapy.core.downloader import Downloader, Slot
from scrapy.crawler import CrawlerRunner from scrapy.crawler import CrawlerRunner
from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future
from scrapy.utils.spider import DefaultSpider from scrapy.utils.spider import DefaultSpider
from scrapy.utils.test import get_crawler from scrapy.utils.test import get_crawler
from tests.mockserver.http import MockServer from tests.mockserver.http import MockServer
@ -104,3 +106,33 @@ def test_params():
assert getattr(expected, param) == getattr(actual, param), ( assert getattr(expected, param) == getattr(actual, param), (
f"Slot.{param}: {getattr(expected, param)!r} != {getattr(actual, param)!r}" f"Slot.{param}: {getattr(expected, param)!r} != {getattr(actual, param)!r}"
) )
@pytest.mark.parametrize(
"priority_queue_class",
[
"scrapy.pqueues.ScrapyPriorityQueue",
"scrapy.pqueues.DownloaderAwarePriorityQueue",
],
)
@deferred_f_from_coro_f
async def test_none_slot_with_priority_queue(
mockserver: MockServer, priority_queue_class: str
) -> None:
"""Test specific cases for None slot handling with different priority queues."""
crawler = get_crawler(
DownloaderSlotsSettingsTestSpider,
settings_dict={"SCHEDULER_PRIORITY_QUEUE": priority_queue_class},
)
await maybe_deferred_to_future(crawler.crawl(mockserver=mockserver))
assert isinstance(crawler.spider, DownloaderSlotsSettingsTestSpider)
assert hasattr(crawler.spider, "times")
assert None not in crawler.spider.times
assert crawler.spider.default_slot in crawler.spider.times
assert len(crawler.spider.times[crawler.spider.default_slot]) == 2
assert crawler.stats
stats = crawler.stats
assert stats.get_value("spider_exceptions", 0) == 0
assert stats.get_value("downloader/exception_count", 0) == 0

View File

@ -438,6 +438,7 @@ class TestEngine(TestEngineBase):
crawler = get_crawler(DefaultSpider) crawler = get_crawler(DefaultSpider)
crawler.spider = crawler._create_spider() crawler.spider = crawler._create_spider()
e = ExecutionEngine(crawler, lambda _: None) e = ExecutionEngine(crawler, lambda _: None)
crawler.engine = e
yield deferred_from_coro(e.open_spider_async()) yield deferred_from_coro(e.open_spider_async())
_schedule_coro(e.start_async()) _schedule_coro(e.start_async())
with pytest.raises(RuntimeError, match="Engine already running"): with pytest.raises(RuntimeError, match="Engine already running"):
@ -450,6 +451,7 @@ class TestEngine(TestEngineBase):
crawler = get_crawler(DefaultSpider) crawler = get_crawler(DefaultSpider)
crawler.spider = crawler._create_spider() crawler.spider = crawler._create_spider()
e = ExecutionEngine(crawler, lambda _: None) e = ExecutionEngine(crawler, lambda _: None)
crawler.engine = e
await e.open_spider_async() await e.open_spider_async()
with pytest.raises(RuntimeError, match="Engine already running"): with pytest.raises(RuntimeError, match="Engine already running"):
await asyncio.gather(e.start_async(), e.start_async()) await asyncio.gather(e.start_async(), e.start_async())
@ -645,6 +647,7 @@ class TestEngineCloseSpider:
@deferred_f_from_coro_f @deferred_f_from_coro_f
async def test_no_slot(self, crawler: Crawler) -> None: async def test_no_slot(self, crawler: Crawler) -> None:
engine = ExecutionEngine(crawler, lambda _: None) engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
slot = engine._slot slot = engine._slot
engine._slot = None engine._slot = None
@ -666,6 +669,7 @@ class TestEngineCloseSpider:
self, crawler: Crawler, caplog: pytest.LogCaptureFixture self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None: ) -> None:
engine = ExecutionEngine(crawler, lambda _: None) engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
assert engine._slot assert engine._slot
del engine._slot.heartbeat del engine._slot.heartbeat
@ -677,6 +681,7 @@ class TestEngineCloseSpider:
self, crawler: Crawler, caplog: pytest.LogCaptureFixture self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None: ) -> None:
engine = ExecutionEngine(crawler, lambda _: None) engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
del engine.downloader.slots del engine.downloader.slots
await engine.close_spider_async() await engine.close_spider_async()
@ -687,6 +692,7 @@ class TestEngineCloseSpider:
self, crawler: Crawler, caplog: pytest.LogCaptureFixture self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None: ) -> None:
engine = ExecutionEngine(crawler, lambda _: None) engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
engine.scraper.slot = None engine.scraper.slot = None
await engine.close_spider_async() await engine.close_spider_async()
@ -697,6 +703,7 @@ class TestEngineCloseSpider:
self, crawler: Crawler, caplog: pytest.LogCaptureFixture self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None: ) -> None:
engine = ExecutionEngine(crawler, lambda _: None) engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
assert engine._slot assert engine._slot
del cast("Scheduler", engine._slot.scheduler).dqs del cast("Scheduler", engine._slot.scheduler).dqs
@ -708,6 +715,7 @@ class TestEngineCloseSpider:
self, crawler: Crawler, caplog: pytest.LogCaptureFixture self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None: ) -> None:
engine = ExecutionEngine(crawler, lambda _: None) engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
signal_manager = engine.signals signal_manager = engine.signals
del engine.signals del engine.signals
@ -725,6 +733,7 @@ class TestEngineCloseSpider:
self, crawler: Crawler, caplog: pytest.LogCaptureFixture self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None: ) -> None:
engine = ExecutionEngine(crawler, lambda _: None) engine = ExecutionEngine(crawler, lambda _: None)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
del cast("MemoryStatsCollector", crawler.stats).spider_stats del cast("MemoryStatsCollector", crawler.stats).spider_stats
await engine.close_spider_async() await engine.close_spider_async()
@ -735,6 +744,7 @@ class TestEngineCloseSpider:
self, crawler: Crawler, caplog: pytest.LogCaptureFixture self, crawler: Crawler, caplog: pytest.LogCaptureFixture
) -> None: ) -> None:
engine = ExecutionEngine(crawler, lambda _: defer.fail(ValueError())) engine = ExecutionEngine(crawler, lambda _: defer.fail(ValueError()))
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
await engine.close_spider_async() await engine.close_spider_async()
assert "Error running spider_closed_callback" in caplog.text assert "Error running spider_closed_callback" in caplog.text
@ -747,6 +757,7 @@ class TestEngineCloseSpider:
raise ValueError raise ValueError
engine = ExecutionEngine(crawler, cb) engine = ExecutionEngine(crawler, cb)
crawler.engine = engine
await engine.open_spider_async() await engine.open_spider_async()
await engine.close_spider_async() await engine.close_spider_async()
assert "Error running spider_closed_callback" in caplog.text assert "Error running spider_closed_callback" in caplog.text