From 1f619c7411314d2a333a5b841b5fca4a1447fe6e Mon Sep 17 00:00:00 2001 From: Adrian Chaves Date: Tue, 23 Jun 2026 15:42:56 +0200 Subject: [PATCH] Make batches resumable through new signals --- docs/topics/signals.rst | 53 +++++++++++ scrapy/extensions/feedexport.py | 49 ++++++++++ scrapy/extensions/spiderstate.py | 64 ++++++++++--- scrapy/signals.py | 3 + tests/test_feedexport.py | 158 +++++++++++++++++++++++++++++++ tests/test_spiderstate.py | 102 ++++++++++++++------ 6 files changed, 388 insertions(+), 41 deletions(-) diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index d13733623..72221da06 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -331,6 +331,39 @@ spider_error :param spider: the spider which raised the exception :type spider: :class:`~scrapy.Spider` object +spider_state_loaded +~~~~~~~~~~~~~~~~~~~ + +.. signal:: spider_state_loaded +.. function:: spider_state_loaded(state) + + Sent by the :class:`~scrapy.extensions.spiderstate.SpiderState` extension + after spider state has been loaded at the start of a crawl. The state dict + is either empty (first run) or restored from the JOBDIR file (resumed run). + + Handlers may read or modify the state dict before any items are scraped. + + This signal supports :ref:`asynchronous handlers `. + + :param state: the spider state dict + :type state: dict + +spider_state_saving +~~~~~~~~~~~~~~~~~~~ + +.. signal:: spider_state_saving +.. function:: spider_state_saving(state) + + Sent by the :class:`~scrapy.extensions.spiderstate.SpiderState` extension + just before spider state is written to disk at the end of a crawl. Handlers + may write values into the state dict to have them persisted for the next + run. + + This signal supports :ref:`asynchronous handlers `. + + :param state: the spider state dict + :type state: dict + feed_slot_closed ~~~~~~~~~~~~~~~~ @@ -356,6 +389,26 @@ feed_exporter_closed This signal supports :ref:`asynchronous handlers `. +feed_slots_initialized +~~~~~~~~~~~~~~~~~~~~~~ + +.. signal:: feed_slots_initialized +.. function:: feed_slots_initialized(slots) + + Sent by the :class:`~scrapy.extensions.feedexport.FeedExporter` extension + when its feed slots are finalized and ready for use. + + When the :class:`~scrapy.extensions.spiderstate.SpiderState` extension is + active (i.e. :setting:`JOBDIR` is set), this signal is sent right after + :signal:`spider_state_loaded`, once any saved batch IDs have been applied to + the slots. Otherwise it is sent at :signal:`engine_started`, which is the + earliest point at which it is known that no state will be loaded. + + This signal supports :ref:`asynchronous handlers `. + + :param slots: the list of feed slots + :type slots: list[:class:`scrapy.extensions.feedexport.FeedSlot`] + memusage_warning_reached ~~~~~~~~~~~~~~~~~~~~~~~~ diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index 8029f85c9..462f8a529 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -446,7 +446,10 @@ class FeedExporter: def from_crawler(cls, crawler: Crawler) -> Self: exporter = cls(crawler) crawler.signals.connect(exporter.open_spider, signals.spider_opened) + crawler.signals.connect(exporter._on_state_loaded, signals.spider_state_loaded) + crawler.signals.connect(exporter._on_engine_started, signals.engine_started) crawler.signals.connect(exporter.close_spider, signals.spider_closed) + crawler.signals.connect(exporter._on_state_saving, signals.spider_state_saving) crawler.signals.connect(exporter.item_scraped, signals.item_scraped) return exporter @@ -507,6 +510,7 @@ class FeedExporter: raise NotConfigured def open_spider(self, spider: Spider) -> None: + self._slots_initialized = False for uri, feed_options in self.feeds.items(): uri_params = self._get_uri_params(spider, feed_options["uri_params"]) self.slots.append( @@ -519,6 +523,51 @@ class FeedExporter: ) ) + async def _on_state_loaded(self, state: dict) -> None: + """Update initial batch-1 slots with the correct resumed batch IDs. + + Called via the spider_state_loaded signal after SpiderState has + populated spider.state. Safe to mutate slots here because FeedSlot + opens its file lazily; no I/O has occurred yet. + """ + feed_batch_ids: dict[str, int] = state.get("feed_batch_ids", {}) + if feed_batch_ids: + spider = self.crawler.spider + for slot in self.slots: + saved_id = feed_batch_ids.get(slot.uri_template, 0) + if saved_id == 0: + continue + batch_id = saved_id + 1 + uri_params = self._get_uri_params( + spider, self.feeds[slot.uri_template]["uri_params"] + ) + uri_params["batch_id"] = batch_id + uri = slot.uri_template % uri_params + slot.batch_id = batch_id + slot.uri = uri + slot.storage = self._get_storage(uri, self.feeds[slot.uri_template]) + self._slots_initialized = True + await self.crawler.signals.send_catch_log_async( + signals.feed_slots_initialized, slots=self.slots + ) + + async def _on_engine_started(self) -> None: + if not self._slots_initialized: + self._slots_initialized = True + await self.crawler.signals.send_catch_log_async( + signals.feed_slots_initialized, slots=self.slots + ) + + def _on_state_saving(self, state: dict) -> None: + """Persist the current batch ID for each feed into spider.state. + + Called via the spider_state_saving signal before SpiderState writes + spider.state to disk, so the next run can resume from the right batch. + """ + feed_batch_ids = state.setdefault("feed_batch_ids", {}) + for slot in self.slots: + feed_batch_ids[slot.uri_template] = slot.batch_id + async def close_spider(self, spider: Spider) -> None: self._pending_close_coros.extend( self._close_slot(slot, spider) for slot in self.slots diff --git a/scrapy/extensions/spiderstate.py b/scrapy/extensions/spiderstate.py index 7b8756572..9f93c5eab 100644 --- a/scrapy/extensions/spiderstate.py +++ b/scrapy/extensions/spiderstate.py @@ -1,11 +1,12 @@ from __future__ import annotations import pickle +import warnings from pathlib import Path from typing import TYPE_CHECKING from scrapy import Spider, signals -from scrapy.exceptions import NotConfigured +from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning from scrapy.utils.job import job_dir if TYPE_CHECKING: @@ -20,6 +21,7 @@ class SpiderState: def __init__(self, jobdir: str | None = None): self.jobdir: str | None = jobdir + self.crawler: Crawler | None = None @classmethod def from_crawler(cls, crawler: Crawler) -> Self: @@ -28,23 +30,63 @@ class SpiderState: raise NotConfigured obj = cls(jobdir) - crawler.signals.connect(obj.spider_closed, signal=signals.spider_closed) - crawler.signals.connect(obj.spider_opened, signal=signals.spider_opened) + obj.crawler = crawler + crawler.signals.connect(obj._spider_closed, signal=signals.spider_closed) + crawler.signals.connect(obj._spider_opened, signal=signals.spider_opened) return obj - def spider_closed(self, spider: Spider) -> None: - if self.jobdir: - with Path(self.statefn).open("wb") as f: - assert hasattr(spider, "state") # set in spider_opened - pickle.dump(spider.state, f, protocol=4) - - def spider_opened(self, spider: Spider) -> None: - if self.jobdir and Path(self.statefn).exists(): + def _load_state(self, spider: Spider) -> None: + if Path(self.statefn).exists(): with Path(self.statefn).open("rb") as f: spider.state = pickle.load(f) # type: ignore[attr-defined] # noqa: S301 else: spider.state = {} # type: ignore[attr-defined] + def _persist_state(self, spider: Spider) -> None: + with Path(self.statefn).open("wb") as f: + assert hasattr(spider, "state") + pickle.dump(spider.state, f, protocol=4) + + async def _spider_opened(self, spider: Spider) -> None: + self._load_state(spider) + assert self.crawler is not None + await self.crawler.signals.send_catch_log_async( + signals.spider_state_loaded, state=spider.state + ) + + async def _spider_closed(self, spider: Spider) -> None: + assert self.crawler is not None + await self.crawler.signals.send_catch_log_async( + signals.spider_state_saving, state=spider.state + ) + self._persist_state(spider) + + def spider_opened(self, spider: Spider) -> None: + warnings.warn( + f"{type(self).__qualname__}.spider_opened() is deprecated, " + "use the spider_state_loaded signal instead.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + self._load_state(spider) + assert self.crawler is not None + self.crawler.signals.send_catch_log( + signals.spider_state_loaded, state=spider.state + ) + + def spider_closed(self, spider: Spider) -> None: + warnings.warn( + f"{type(self).__qualname__}.spider_closed() is deprecated, " + "use the spider_state_saving signal instead.", + category=ScrapyDeprecationWarning, + stacklevel=2, + ) + assert self.crawler is not None + self.crawler.signals.send_catch_log( + signals.spider_state_saving, state=spider.state + ) + self._persist_state(spider) + @property def statefn(self) -> str: assert self.jobdir diff --git a/scrapy/signals.py b/scrapy/signals.py index 972f4fd60..bf17d093a 100644 --- a/scrapy/signals.py +++ b/scrapy/signals.py @@ -26,3 +26,6 @@ item_dropped = object() item_error = object() feed_slot_closed = object() feed_exporter_closed = object() +feed_slots_initialized = object() +spider_state_loaded = object() +spider_state_saving = object() diff --git a/tests/test_feedexport.py b/tests/test_feedexport.py index c1d6f04eb..36f05315b 100644 --- a/tests/test_feedexport.py +++ b/tests/test_feedexport.py @@ -1402,3 +1402,161 @@ class TestFeedExportInit: crawler = get_crawler(settings_dict=settings) exporter = FeedExporter.from_crawler(crawler) assert isinstance(exporter, FeedExporter) + + +class TestFeedExporterBatchIdState: + """Tests that batch_id is persisted across JOBDIR-based resumption.""" + + items = [{"foo": "bar"}, {"foo": "baz"}] + + def _make_exporter(self, uri_template, batch_item_count=1): + settings = { + "FEEDS": { + uri_template: { + "format": "jl", + "batch_item_count": batch_item_count, + }, + }, + } + return FeedExporter.from_crawler(get_crawler(settings_dict=settings)) + + @coroutine_test + async def test_fresh_crawl_starts_at_one(self): + """Without saved state, batch_id starts at 1.""" + with tempfile.TemporaryDirectory() as tmpdir: + uri = f"file://{tmpdir}/feed-%(batch_id)d.jl" + exporter = self._make_exporter(uri) + spider = scrapy.Spider("testspider") + + exporter.open_spider(spider) + assert exporter.slots[0].batch_id == 1 + + @coroutine_test + async def test_multiple_feeds_tracked_independently(self): + """Each feed URI is tracked independently in spider.state.""" + with tempfile.TemporaryDirectory() as tmpdir: + uri_a = f"file://{tmpdir}/a-%(batch_id)d.jl" + uri_b = f"file://{tmpdir}/b-%(batch_id)d.jl" + settings = { + "FEEDS": { + uri_a: {"format": "jl", "batch_item_count": 1}, + uri_b: {"format": "jl", "batch_item_count": 1}, + }, + } + exporter = FeedExporter.from_crawler(get_crawler(settings_dict=settings)) + spider = scrapy.Spider("testspider") + spider.state = {"feed_batch_ids": {uri_a: 3, uri_b: 7}} + + exporter.open_spider(spider) + await exporter._on_state_loaded(spider.state) + batch_ids = {slot.uri_template: slot.batch_id for slot in exporter.slots} + assert batch_ids[uri_a] == 4 + assert batch_ids[uri_b] == 8 + + @coroutine_test + async def test_no_jobdir_no_error(self): + """open_spider and close_spider work when JOBDIR is not configured. + + Without JOBDIR, SpiderState is not active, so spider_state_loaded and + spider_state_saving never fire. FeedExporter falls back to batch_id=1 + with no state persisted. + """ + with tempfile.TemporaryDirectory() as tmpdir: + uri = f"file://{tmpdir}/feed-%(batch_id)d.jl" + exporter = self._make_exporter(uri) + spider = scrapy.Spider("testspider") + + exporter.open_spider(spider) + assert exporter.slots[0].batch_id == 1 + + for item in self.items: + exporter.item_scraped(item, spider) + await exporter.close_spider(spider) + + @coroutine_test + async def test_batch_id_persists_across_jobdir_runs(self): + """Resuming with JOBDIR continues batch numbering instead of overwriting files. + + Uses the full Scrapy machinery (real crawl_async()) so that any internal + change to signal ordering or extension wiring is caught. SpiderState + fires spider_state_loaded / spider_state_saving; FeedExporter uses those + to resume from the right batch ID, so the second run's files (feed-4.jl, + feed-5.jl) never collide with the first run's files (feed-1.jl, feed-2.jl). + """ + with tempfile.TemporaryDirectory() as tmpdir: + jobdir = Path(tmpdir) / "jobdir" + jobdir.mkdir() + feed_uri_template = f"file://{tmpdir}/feed-%(batch_id)d.jl" + settings = { + "FEEDS": {feed_uri_template: {"format": "jl", "batch_item_count": 1}}, + "JOBDIR": str(jobdir), + } + + class BatchSpider(scrapy.Spider): + name = "batchspider" + + async def start(self): + for item in [{"foo": "bar"}, {"foo": "baz"}]: + yield item + + # First run: 2 items → feed-1.jl, feed-2.jl + await get_crawler(BatchSpider, settings).crawl_async() + content_run1 = { + 1: Path(f"{tmpdir}/feed-1.jl").read_bytes(), + 2: Path(f"{tmpdir}/feed-2.jl").read_bytes(), + } + + # Second run (resume): same JOBDIR, 2 more items + await get_crawler(BatchSpider, settings).crawl_async() + + # First-run files must be untouched + assert Path(f"{tmpdir}/feed-1.jl").read_bytes() == content_run1[1] + assert Path(f"{tmpdir}/feed-2.jl").read_bytes() == content_run1[2] + # Second run created new files with higher batch IDs + assert Path(f"{tmpdir}/feed-4.jl").exists() + assert Path(f"{tmpdir}/feed-5.jl").exists() + + @coroutine_test + async def test_feed_slots_initialized_fires_after_state_loaded(self): + """feed_slots_initialized fires with correct slots when state is restored.""" + with tempfile.TemporaryDirectory() as tmpdir: + uri = f"file://{tmpdir}/feed-%(batch_id)d.jl" + exporter = self._make_exporter(uri) + spider = scrapy.Spider("testspider") + spider.state = {"feed_batch_ids": {uri: 5}} + + received: list[list] = [] + + def on_initialized(slots): + received.append(list(slots)) + + exporter.crawler.signals.connect( + on_initialized, signal=signals.feed_slots_initialized, weak=False + ) + exporter.open_spider(spider) + await exporter._on_state_loaded(spider.state) + + assert len(received) == 1 + assert received[0][0].batch_id == 6 + + @coroutine_test + async def test_feed_slots_initialized_fires_from_engine_started_without_state(self): + """feed_slots_initialized fires via engine_started when no state is loaded.""" + with tempfile.TemporaryDirectory() as tmpdir: + uri = f"file://{tmpdir}/feed-%(batch_id)d.jl" + exporter = self._make_exporter(uri) + spider = scrapy.Spider("testspider") + + received: list[list] = [] + + def on_initialized(slots): + received.append(list(slots)) + + exporter.crawler.signals.connect( + on_initialized, signal=signals.feed_slots_initialized, weak=False + ) + exporter.open_spider(spider) + await exporter._on_engine_started() + + assert len(received) == 1 + assert received[0][0].batch_id == 1 diff --git a/tests/test_spiderstate.py b/tests/test_spiderstate.py index e44cfca90..255ff21f1 100644 --- a/tests/test_spiderstate.py +++ b/tests/test_spiderstate.py @@ -5,48 +5,90 @@ from typing import TYPE_CHECKING import pytest -from scrapy.exceptions import NotConfigured +from scrapy import signals +from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning from scrapy.extensions.spiderstate import SpiderState from scrapy.spiders import Spider from scrapy.utils.test import get_crawler +from tests.utils.decorators import coroutine_test if TYPE_CHECKING: from pathlib import Path -def test_store_load(tmp_path: Path) -> None: - jobdir = str(tmp_path) - +def test_deprecated_methods(tmp_path: Path) -> None: + crawler = get_crawler(Spider, {"JOBDIR": str(tmp_path)}) + ss = SpiderState.from_crawler(crawler) spider = Spider(name="default") - dt = datetime.now(tz=timezone.utc) - - ss = SpiderState(jobdir) - ss.spider_opened(spider) - assert hasattr(spider, "state") - spider.state["one"] = 1 - spider.state["dt"] = dt - ss.spider_closed(spider) - - spider2 = Spider(name="default") - ss2 = SpiderState(jobdir) - ss2.spider_opened(spider2) - assert hasattr(spider2, "state") - assert spider2.state == {"one": 1, "dt": dt} - ss2.spider_closed(spider2) - - -def test_state_attribute() -> None: - # state attribute must be present if jobdir is not set, to provide a - # consistent interface - spider = Spider(name="default") - ss = SpiderState() - ss.spider_opened(spider) - assert hasattr(spider, "state") - assert spider.state == {} - ss.spider_closed(spider) + with pytest.warns(ScrapyDeprecationWarning, match="spider_opened"): + ss.spider_opened(spider) + with pytest.warns(ScrapyDeprecationWarning, match="spider_closed"): + ss.spider_closed(spider) def test_not_configured() -> None: crawler = get_crawler(Spider) with pytest.raises(NotConfigured): SpiderState.from_crawler(crawler) + + +@coroutine_test +async def test_store_load(tmp_path: Path) -> None: + jobdir = str(tmp_path) + spider = Spider(name="default") + dt = datetime.now(tz=timezone.utc) + + crawler = get_crawler(Spider, {"JOBDIR": jobdir}) + ss = SpiderState.from_crawler(crawler) + await ss._spider_opened(spider) + assert hasattr(spider, "state") + spider.state["one"] = 1 + spider.state["dt"] = dt + await ss._spider_closed(spider) + + spider2 = Spider(name="default") + crawler2 = get_crawler(Spider, {"JOBDIR": jobdir}) + ss2 = SpiderState.from_crawler(crawler2) + await ss2._spider_opened(spider2) + assert hasattr(spider2, "state") + assert spider2.state == {"one": 1, "dt": dt} + await ss2._spider_closed(spider2) + + +@coroutine_test +async def test_spider_state_loaded_signal_fires(tmp_path: Path) -> None: + crawler = get_crawler(Spider, {"JOBDIR": str(tmp_path)}) + ss = SpiderState.from_crawler(crawler) + + received: list[dict] = [] + + def on_loaded(state: dict) -> None: + received.append(dict(state)) + + crawler.signals.connect(on_loaded, signal=signals.spider_state_loaded, weak=False) + + spider = Spider(name="default") + await ss._spider_opened(spider) + + assert received == [{}] + assert spider.state == {} + + +@coroutine_test +async def test_spider_state_saving_signal_fires(tmp_path: Path) -> None: + crawler = get_crawler(Spider, {"JOBDIR": str(tmp_path)}) + ss = SpiderState.from_crawler(crawler) + + spider = Spider(name="default") + await ss._spider_opened(spider) + spider.state["key"] = "value" + + saving_calls: list[dict] = [] + + def on_saving(state: dict) -> None: + saving_calls.append(dict(state)) + + crawler.signals.connect(on_saving, signal=signals.spider_state_saving, weak=False) + await ss._spider_closed(spider) + + assert saving_calls == [{"key": "value"}]