diff --git a/docs/topics/signals.rst b/docs/topics/signals.rst index bd14f5993..8cdb1a5c4 100644 --- a/docs/topics/signals.rst +++ b/docs/topics/signals.rst @@ -333,6 +333,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 ~~~~~~~~~~~~~~~~ @@ -358,6 +391,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 fa708a2e1..e8a13a806 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -448,7 +448,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 @@ -509,6 +512,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( @@ -521,6 +525,52 @@ class FeedExporter: ) ) + async def _on_state_loaded(self, state: dict[str, Any]) -> 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 + assert spider is not None + for slot in self.slots: + saved_id = feed_batch_ids.get(slot.uri_template, 0) + if saved_id == 0: + continue + batch_id = saved_id + 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[str, Any]) -> 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..e457d14b2 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,67 @@ 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, # type: ignore[attr-defined] + ) + + 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, # type: ignore[attr-defined] + ) + 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, # type: ignore[attr-defined] + ) + + 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, # type: ignore[attr-defined] + ) + 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 7d751f188..d7ce9a0dd 100644 --- a/tests/test_feedexport.py +++ b/tests/test_feedexport.py @@ -1408,3 +1408,163 @@ 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"{path_to_url(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"{path_to_url(tmpdir)}/a-%(batch_id)d.jl" + uri_b = f"{path_to_url(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) + exporter.crawler.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] == 3 + assert batch_ids[uri_b] == 7 + + @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"{path_to_url(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-3.jl, + feed-4.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"{path_to_url(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-3.jl").exists() + assert Path(f"{tmpdir}/feed-4.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"{path_to_url(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) + exporter.crawler.spider = spider + await exporter._on_state_loaded(spider.state) + + assert len(received) == 1 + assert received[0][0].batch_id == 5 + + @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"{path_to_url(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_feedexport_batch.py b/tests/test_feedexport_batch.py index 469985288..5a1d896a7 100644 --- a/tests/test_feedexport_batch.py +++ b/tests/test_feedexport_batch.py @@ -364,6 +364,61 @@ class TestBatchDeliveries(TestFeedExportBase): data = await self.exported_data(items, settings) assert len(items) == len(data["json"]) + @coroutine_test + async def test_jobdir_batch_id_continues_after_restart(self): + """Regression test for #5153. + + When JOBDIR is set and the feed URI references ``%(batch_id)``, the + batch_id counter must persist across restarts so that re-running the + crawl with the same JOBDIR does not overwrite previously-written + batch files. + """ + items = [ + self.MyItem({"foo": "bar1", "egg": "spam1"}), + self.MyItem({"foo": "bar2", "egg": "spam2"}), + ] + feed_dir = self._random_temp_filename() + jobdir = self._random_temp_filename() + uri_template = feed_dir / "%(batch_id)d.jl" + + def make_settings(): + return { + "FEEDS": { + uri_template: {"format": "jl"}, + }, + "FEED_EXPORT_BATCH_ITEM_COUNT": 1, + "JOBDIR": str(jobdir), + } + + # First run: should produce 1.jl and 2.jl (one per item). + await self.exported_data(items, make_settings()) + first_run_files = sorted(p.name for p in feed_dir.iterdir()) + assert first_run_files == ["1.jl", "2.jl"] + # Mark each file so that we can detect overwrites. + sentinel = b"SENTINEL-FROM-FIRST-RUN\n" + for name in first_run_files: + with (feed_dir / name).open("ab") as f: + f.write(sentinel) + first_run_contents = { + name: (feed_dir / name).read_bytes() for name in first_run_files + } + + # Second run with the same JOBDIR: must NOT overwrite the prior files + # and must start the batch_id counter at 3. + more_items = [ + self.MyItem({"foo": "bar3", "egg": "spam3"}), + self.MyItem({"foo": "bar4", "egg": "spam4"}), + ] + await self.exported_data(more_items, make_settings()) + + for name, prior in first_run_contents.items(): + assert (feed_dir / name).read_bytes() == prior, ( + f"{name} was overwritten by the second run" + ) + + all_files = sorted(p.name for p in feed_dir.iterdir()) + assert all_files == ["1.jl", "2.jl", "3.jl", "4.jl"] + @inline_callbacks_test def test_stats_batch_file_success(self): settings = { diff --git a/tests/test_spiderstate.py b/tests/test_spiderstate.py index e44cfca90..968da3da6 100644 --- a/tests/test_spiderstate.py +++ b/tests/test_spiderstate.py @@ -1,52 +1,94 @@ from __future__ import annotations from datetime import datetime, timezone -from typing import TYPE_CHECKING +from typing import TYPE_CHECKING, Any 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[str, Any]] = [] + + def on_loaded(state: dict[str, Any]) -> 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 == {} # type: ignore[attr-defined] + + +@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" # type: ignore[attr-defined] + + saving_calls: list[dict[str, Any]] = [] + + def on_saving(state: dict[str, Any]) -> 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"}]