mirror of https://github.com/scrapy/scrapy.git
Merge 7284a75971 into c9446931a8
This commit is contained in:
commit
172ffbe3c4
|
|
@ -717,6 +717,12 @@ The command line above can generate a directory tree like::
|
|||
Where the first and second files contain exactly 100 items. The last one contains
|
||||
100 items or fewer.
|
||||
|
||||
When :setting:`JOBDIR` is set and the feed URI contains the ``%(batch_id)d``
|
||||
placeholder, the next ``batch_id`` to use is persisted under
|
||||
``<JOBDIR>/feedexport.state`` so that resuming a crawl with the same
|
||||
:setting:`JOBDIR` continues numbering after the last batch produced by the
|
||||
previous run, instead of restarting at ``1`` and overwriting prior files.
|
||||
|
||||
|
||||
.. setting:: FEED_URI_PARAMS
|
||||
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ from __future__ import annotations
|
|||
|
||||
import asyncio
|
||||
import contextlib
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
import sys
|
||||
|
|
@ -31,6 +32,7 @@ from scrapy.utils.asyncio import is_asyncio_available, run_in_thread
|
|||
from scrapy.utils.conf import feed_complete_default_values_from_settings
|
||||
from scrapy.utils.defer import deferred_from_coro, ensure_awaitable
|
||||
from scrapy.utils.ftp import ftp_store_file
|
||||
from scrapy.utils.job import job_dir
|
||||
from scrapy.utils.misc import build_from_crawler, load_object
|
||||
from scrapy.utils.python import without_none_values
|
||||
|
||||
|
|
@ -51,6 +53,9 @@ UriParamsCallableT: TypeAlias = Callable[
|
|||
[dict[str, Any], Spider], dict[str, Any] | None
|
||||
]
|
||||
|
||||
_BATCH_ID_RE = re.compile(r"%\(batch_id\)")
|
||||
_BATCH_STATE_FILENAME = "feedexport.state"
|
||||
|
||||
|
||||
class ItemFilter:
|
||||
"""
|
||||
|
|
@ -459,6 +464,7 @@ class FeedExporter:
|
|||
self.slots: list[FeedSlot] = []
|
||||
self.filters: dict[str, ItemFilter] = {}
|
||||
self._pending_close_coros: list[Coroutine[Any, Any, None]] = []
|
||||
self._persisted_batch_ids: dict[str, int] = self._load_batch_state()
|
||||
|
||||
if not self.settings["FEEDS"] and not self.settings["FEED_URI"]:
|
||||
raise NotConfigured
|
||||
|
|
@ -510,10 +516,13 @@ class FeedExporter:
|
|||
|
||||
def open_spider(self, spider: Spider) -> None:
|
||||
for uri, feed_options in self.feeds.items():
|
||||
uri_params = self._get_uri_params(spider, feed_options["uri_params"])
|
||||
batch_id = self._next_batch_id(uri)
|
||||
uri_params = self._get_uri_params(
|
||||
spider, feed_options["uri_params"], batch_id=batch_id
|
||||
)
|
||||
self.slots.append(
|
||||
self._start_new_batch(
|
||||
batch_id=1,
|
||||
batch_id=batch_id,
|
||||
uri=uri % uri_params,
|
||||
feed_options=feed_options,
|
||||
spider=spider,
|
||||
|
|
@ -560,6 +569,10 @@ class FeedExporter:
|
|||
# In this case, the file is not stored, so no processing is required.
|
||||
return
|
||||
|
||||
# Record that this batch_id was actually written, so a restart with the
|
||||
# same JOBDIR will pick the next unused id and avoid overwriting it.
|
||||
self._record_batch_id(slot.uri_template, slot.batch_id)
|
||||
|
||||
logmsg = f"{slot.format} feed ({slot.itemcount} items) in: {slot.uri}"
|
||||
slot_type = type(slot.storage).__name__
|
||||
assert self.crawler.stats
|
||||
|
|
@ -636,9 +649,10 @@ class FeedExporter:
|
|||
spider, self.feeds[slot.uri_template]["uri_params"], slot
|
||||
)
|
||||
self._pending_close_coros.append(self._close_slot(slot, spider))
|
||||
new_batch_id = slot.batch_id + 1
|
||||
slots.append(
|
||||
self._start_new_batch(
|
||||
batch_id=slot.batch_id + 1,
|
||||
batch_id=new_batch_id,
|
||||
uri=slot.uri_template % uri_params,
|
||||
feed_options=self.feeds[slot.uri_template],
|
||||
spider=spider,
|
||||
|
|
@ -711,12 +725,17 @@ class FeedExporter:
|
|||
spider: Spider,
|
||||
uri_params_function: str | UriParamsCallableT | None,
|
||||
slot: FeedSlot | None = None,
|
||||
*,
|
||||
batch_id: int | None = None,
|
||||
) -> dict[str, Any]:
|
||||
params = {k: getattr(spider, k) for k in dir(spider)}
|
||||
utc_now = datetime.now(tz=timezone.utc)
|
||||
params["time"] = utc_now.replace(microsecond=0).isoformat().replace(":", "-")
|
||||
params["batch_time"] = utc_now.isoformat().replace(":", "-")
|
||||
params["batch_id"] = slot.batch_id + 1 if slot is not None else 1
|
||||
if batch_id is not None:
|
||||
params["batch_id"] = batch_id
|
||||
else:
|
||||
params["batch_id"] = slot.batch_id + 1 if slot is not None else 1
|
||||
uripar_function: UriParamsCallableT = (
|
||||
load_object(uri_params_function)
|
||||
if uri_params_function
|
||||
|
|
@ -731,3 +750,74 @@ class FeedExporter:
|
|||
feed_options.get("item_filter", ItemFilter)
|
||||
)
|
||||
return item_filter_class(feed_options)
|
||||
|
||||
def _batch_state_path(self) -> Path | None:
|
||||
"""Return the path to the persisted batch-id state file when JOBDIR is set."""
|
||||
jobdir = job_dir(self.settings)
|
||||
if not jobdir:
|
||||
return None
|
||||
return Path(jobdir, _BATCH_STATE_FILENAME)
|
||||
|
||||
def _load_batch_state(self) -> dict[str, int]:
|
||||
"""Load the persisted ``{uri_template: last_batch_id_used}`` mapping.
|
||||
|
||||
Returns an empty dict if JOBDIR is not set or the state file does not
|
||||
yet exist. Unreadable or malformed state files are ignored so a
|
||||
corrupted file does not prevent a crawl from starting.
|
||||
"""
|
||||
state_path = self._batch_state_path()
|
||||
if state_path is None or not state_path.exists():
|
||||
return {}
|
||||
try:
|
||||
with state_path.open("r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
except (OSError, ValueError):
|
||||
logger.warning(
|
||||
"Could not read feed exporter batch state at %s; "
|
||||
"starting batch_id counters from 1.",
|
||||
state_path,
|
||||
)
|
||||
return {}
|
||||
if not isinstance(data, dict):
|
||||
return {}
|
||||
return {str(k): int(v) for k, v in data.items()}
|
||||
|
||||
def _save_batch_state(self) -> None:
|
||||
"""Persist the current ``self._persisted_batch_ids`` to JOBDIR."""
|
||||
state_path = self._batch_state_path()
|
||||
if state_path is None:
|
||||
return
|
||||
try:
|
||||
with state_path.open("w", encoding="utf-8") as f:
|
||||
json.dump(self._persisted_batch_ids, f)
|
||||
except OSError:
|
||||
logger.warning(
|
||||
"Could not write feed exporter batch state to %s; "
|
||||
"batch_id may reset on restart.",
|
||||
state_path,
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
def _next_batch_id(self, uri_template: str) -> int:
|
||||
"""Compute the next batch_id to use for *uri_template*.
|
||||
|
||||
When JOBDIR is set and *uri_template* references ``%(batch_id)``, the
|
||||
next id continues from the last value persisted across runs; otherwise
|
||||
it is 1.
|
||||
"""
|
||||
if self._batch_state_path() is None or not _BATCH_ID_RE.search(uri_template):
|
||||
return 1
|
||||
return self._persisted_batch_ids.get(uri_template, 0) + 1
|
||||
|
||||
def _record_batch_id(self, uri_template: str, batch_id: int) -> None:
|
||||
"""Persist that *batch_id* has been used (written to storage) for
|
||||
*uri_template*. The stored value is the maximum batch_id seen, so
|
||||
out-of-order slot close-outs cannot lower it.
|
||||
"""
|
||||
if self._batch_state_path() is None or not _BATCH_ID_RE.search(uri_template):
|
||||
return
|
||||
current = self._persisted_batch_ids.get(uri_template, 0)
|
||||
if batch_id <= current:
|
||||
return
|
||||
self._persisted_batch_ids[uri_template] = batch_id
|
||||
self._save_batch_state()
|
||||
|
|
|
|||
|
|
@ -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 = {
|
||||
|
|
|
|||
Loading…
Reference in New Issue