mirror of https://github.com/scrapy/scrapy.git
Add feed_slot_closed and feed_exporter_closed signals (#5876)
This commit is contained in:
parent
8c8fb67057
commit
b50c032ee9
|
|
@ -307,6 +307,33 @@ spider_error
|
|||
:param spider: the spider which raised the exception
|
||||
:type spider: :class:`~scrapy.Spider` object
|
||||
|
||||
feed_slot_closed
|
||||
~~~~~~~~~~~~~~~~
|
||||
|
||||
.. signal:: feed_slot_closed
|
||||
.. function:: feed_slot_closed(slot)
|
||||
|
||||
Sent when a :ref:`feed exports <topics-feed-exports>` slot is closed.
|
||||
|
||||
This signal supports returning deferreds from its handlers.
|
||||
|
||||
:param slot: the slot closed
|
||||
:type slot: scrapy.extensions.feedexport.FeedSlot
|
||||
|
||||
|
||||
feed_exporter_closed
|
||||
~~~~~~~~~~~~~~~~~~~~
|
||||
|
||||
.. signal:: feed_exporter_closed
|
||||
.. function:: feed_exporter_closed()
|
||||
|
||||
Sent when the :ref:`feed exports <topics-feed-exports>` extension is closed,
|
||||
during the handling of the :signal:`spider_closed` signal by the extension,
|
||||
after all feed exporting has been handled.
|
||||
|
||||
This signal supports returning deferreds from its handlers.
|
||||
|
||||
|
||||
Request signals
|
||||
---------------
|
||||
|
||||
|
|
|
|||
|
|
@ -11,10 +11,11 @@ import warnings
|
|||
from datetime import datetime
|
||||
from pathlib import Path
|
||||
from tempfile import NamedTemporaryFile
|
||||
from typing import IO, Any, Callable, Optional, Tuple, Union
|
||||
from typing import IO, Any, Callable, List, Optional, Tuple, Union
|
||||
from urllib.parse import unquote, urlparse
|
||||
|
||||
from twisted.internet import defer, threads
|
||||
from twisted.internet.defer import DeferredList
|
||||
from w3lib.url import file_uri_to_path
|
||||
from zope.interface import Interface, implementer
|
||||
|
||||
|
|
@ -23,6 +24,8 @@ from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning
|
|||
from scrapy.extensions.postprocessing import PostProcessingManager
|
||||
from scrapy.utils.boto import is_botocore_available
|
||||
from scrapy.utils.conf import feed_complete_default_values_from_settings
|
||||
from scrapy.utils.defer import maybe_deferred_to_future
|
||||
from scrapy.utils.deprecate import create_deprecated_class
|
||||
from scrapy.utils.ftp import ftp_store_file
|
||||
from scrapy.utils.log import failure_to_exc_info
|
||||
from scrapy.utils.misc import create_instance, load_object
|
||||
|
|
@ -271,7 +274,7 @@ class FTPFeedStorage(BlockingFeedStorage):
|
|||
)
|
||||
|
||||
|
||||
class _FeedSlot:
|
||||
class FeedSlot:
|
||||
def __init__(
|
||||
self,
|
||||
file,
|
||||
|
|
@ -309,7 +312,15 @@ class _FeedSlot:
|
|||
self._exporting = False
|
||||
|
||||
|
||||
_FeedSlot = create_deprecated_class(
|
||||
name="_FeedSlot",
|
||||
new_class=FeedSlot,
|
||||
)
|
||||
|
||||
|
||||
class FeedExporter:
|
||||
_pending_deferreds: List[defer.Deferred] = []
|
||||
|
||||
@classmethod
|
||||
def from_crawler(cls, crawler):
|
||||
exporter = cls(crawler)
|
||||
|
|
@ -375,12 +386,18 @@ class FeedExporter:
|
|||
)
|
||||
)
|
||||
|
||||
def close_spider(self, spider):
|
||||
deferred_list = []
|
||||
async def close_spider(self, spider):
|
||||
for slot in self.slots:
|
||||
d = self._close_slot(slot, spider)
|
||||
deferred_list.append(d)
|
||||
return defer.DeferredList(deferred_list) if deferred_list else None
|
||||
self._close_slot(slot, spider)
|
||||
|
||||
# Await all deferreds
|
||||
if self._pending_deferreds:
|
||||
await maybe_deferred_to_future(DeferredList(self._pending_deferreds))
|
||||
|
||||
# Send FEED_EXPORTER_CLOSED signal
|
||||
await maybe_deferred_to_future(
|
||||
self.crawler.signals.send_catch_log_deferred(signals.feed_exporter_closed)
|
||||
)
|
||||
|
||||
def _close_slot(self, slot, spider):
|
||||
def get_file(slot_):
|
||||
|
|
@ -404,6 +421,14 @@ class FeedExporter:
|
|||
d.addErrback(
|
||||
self._handle_store_error, logmsg, spider, type(slot.storage).__name__
|
||||
)
|
||||
self._pending_deferreds.append(d)
|
||||
d.addCallback(
|
||||
lambda _: self.crawler.signals.send_catch_log_deferred(
|
||||
signals.feed_slot_closed, slot=slot
|
||||
)
|
||||
)
|
||||
d.addBoth(lambda _: self._pending_deferreds.remove(d))
|
||||
|
||||
return d
|
||||
|
||||
def _handle_store_error(self, f, logmsg, spider, slot_type):
|
||||
|
|
@ -444,7 +469,7 @@ class FeedExporter:
|
|||
indent=feed_options["indent"],
|
||||
**feed_options["item_export_kwargs"],
|
||||
)
|
||||
slot = _FeedSlot(
|
||||
slot = FeedSlot(
|
||||
file=file,
|
||||
exporter=exporter,
|
||||
storage=storage,
|
||||
|
|
@ -579,7 +604,7 @@ class FeedExporter:
|
|||
self,
|
||||
spider: Spider,
|
||||
uri_params_function: Optional[Union[str, Callable[[dict, Spider], dict]]],
|
||||
slot: Optional[_FeedSlot] = None,
|
||||
slot: Optional[FeedSlot] = None,
|
||||
) -> dict:
|
||||
params = {}
|
||||
for k in dir(spider):
|
||||
|
|
|
|||
|
|
@ -22,6 +22,8 @@ bytes_received = object()
|
|||
item_scraped = object()
|
||||
item_dropped = object()
|
||||
item_error = object()
|
||||
feed_slot_closed = object()
|
||||
feed_exporter_closed = object()
|
||||
|
||||
# for backward compatibility
|
||||
stats_spider_opened = spider_opened
|
||||
|
|
|
|||
|
|
@ -32,18 +32,19 @@ from zope.interface import implementer
|
|||
from zope.interface.verify import verifyObject
|
||||
|
||||
import scrapy
|
||||
from scrapy import signals
|
||||
from scrapy.exceptions import NotConfigured, ScrapyDeprecationWarning
|
||||
from scrapy.exporters import CsvItemExporter, JsonItemExporter
|
||||
from scrapy.extensions.feedexport import (
|
||||
BlockingFeedStorage,
|
||||
FeedExporter,
|
||||
FeedSlot,
|
||||
FileFeedStorage,
|
||||
FTPFeedStorage,
|
||||
GCSFeedStorage,
|
||||
IFeedStorage,
|
||||
S3FeedStorage,
|
||||
StdoutFeedStorage,
|
||||
_FeedSlot,
|
||||
)
|
||||
from scrapy.settings import Settings
|
||||
from scrapy.utils.python import to_unicode
|
||||
|
|
@ -660,8 +661,8 @@ class FeedExportTestBase(ABC, unittest.TestCase):
|
|||
return result
|
||||
|
||||
|
||||
class InstrumentedFeedSlot(_FeedSlot):
|
||||
"""Instrumented _FeedSlot subclass for keeping track of calls to
|
||||
class InstrumentedFeedSlot(FeedSlot):
|
||||
"""Instrumented FeedSlot subclass for keeping track of calls to
|
||||
start_exporting and finish_exporting."""
|
||||
|
||||
def start_exporting(self):
|
||||
|
|
@ -964,7 +965,7 @@ class FeedExportTest(FeedExportTestBase):
|
|||
listener = IsExportingListener()
|
||||
InstrumentedFeedSlot.subscribe__listener(listener)
|
||||
|
||||
with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot):
|
||||
with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot):
|
||||
_ = yield self.exported_data(items, settings)
|
||||
self.assertFalse(listener.start_without_finish)
|
||||
self.assertFalse(listener.finish_without_start)
|
||||
|
|
@ -982,7 +983,7 @@ class FeedExportTest(FeedExportTestBase):
|
|||
listener = IsExportingListener()
|
||||
InstrumentedFeedSlot.subscribe__listener(listener)
|
||||
|
||||
with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot):
|
||||
with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot):
|
||||
_ = yield self.exported_data(items, settings)
|
||||
self.assertFalse(listener.start_without_finish)
|
||||
self.assertFalse(listener.finish_without_start)
|
||||
|
|
@ -1003,7 +1004,7 @@ class FeedExportTest(FeedExportTestBase):
|
|||
listener = IsExportingListener()
|
||||
InstrumentedFeedSlot.subscribe__listener(listener)
|
||||
|
||||
with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot):
|
||||
with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot):
|
||||
_ = yield self.exported_data(items, settings)
|
||||
self.assertFalse(listener.start_without_finish)
|
||||
self.assertFalse(listener.finish_without_start)
|
||||
|
|
@ -1022,7 +1023,7 @@ class FeedExportTest(FeedExportTestBase):
|
|||
listener = IsExportingListener()
|
||||
InstrumentedFeedSlot.subscribe__listener(listener)
|
||||
|
||||
with mock.patch("scrapy.extensions.feedexport._FeedSlot", InstrumentedFeedSlot):
|
||||
with mock.patch("scrapy.extensions.feedexport.FeedSlot", InstrumentedFeedSlot):
|
||||
_ = yield self.exported_data(items, settings)
|
||||
self.assertFalse(listener.start_without_finish)
|
||||
self.assertFalse(listener.finish_without_start)
|
||||
|
|
@ -2651,6 +2652,83 @@ class BatchDeliveriesTest(FeedExportTestBase):
|
|||
stub.assert_no_pending_responses()
|
||||
|
||||
|
||||
# Test that the FeedExporer sends the feed_exporter_closed and feed_slot_closed signals
|
||||
class FeedExporterSignalsTest(unittest.TestCase):
|
||||
items = [
|
||||
{"foo": "bar1", "egg": "spam1"},
|
||||
{"foo": "bar2", "egg": "spam2", "baz": "quux2"},
|
||||
{"foo": "bar3", "baz": "quux3"},
|
||||
]
|
||||
|
||||
with tempfile.NamedTemporaryFile(suffix="json") as tmp:
|
||||
settings = {
|
||||
"FEEDS": {
|
||||
f"file:///{tmp.name}": {
|
||||
"format": "json",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
def feed_exporter_closed_signal_handler(self):
|
||||
self.feed_exporter_closed_received = True
|
||||
|
||||
def feed_slot_closed_signal_handler(self, slot):
|
||||
self.feed_slot_closed_received = True
|
||||
|
||||
def feed_exporter_closed_signal_handler_deferred(self):
|
||||
d = defer.Deferred()
|
||||
d.addCallback(lambda _: setattr(self, "feed_exporter_closed_received", True))
|
||||
d.callback(None)
|
||||
return d
|
||||
|
||||
def feed_slot_closed_signal_handler_deferred(self, slot):
|
||||
d = defer.Deferred()
|
||||
d.addCallback(lambda _: setattr(self, "feed_slot_closed_received", True))
|
||||
d.callback(None)
|
||||
return d
|
||||
|
||||
def run_signaled_feed_exporter(
|
||||
self, feed_exporter_signal_handler, feed_slot_signal_handler
|
||||
):
|
||||
crawler = get_crawler(settings_dict=self.settings)
|
||||
feed_exporter = FeedExporter.from_crawler(crawler)
|
||||
spider = scrapy.Spider("default")
|
||||
spider.crawler = crawler
|
||||
crawler.signals.connect(
|
||||
feed_exporter_signal_handler,
|
||||
signal=signals.feed_exporter_closed,
|
||||
)
|
||||
crawler.signals.connect(
|
||||
feed_slot_signal_handler, signal=signals.feed_slot_closed
|
||||
)
|
||||
feed_exporter.open_spider(spider)
|
||||
for item in self.items:
|
||||
feed_exporter.item_scraped(item, spider)
|
||||
defer.ensureDeferred(feed_exporter.close_spider(spider))
|
||||
|
||||
def test_feed_exporter_signals_sent(self):
|
||||
self.feed_exporter_closed_received = False
|
||||
self.feed_slot_closed_received = False
|
||||
|
||||
self.run_signaled_feed_exporter(
|
||||
self.feed_exporter_closed_signal_handler,
|
||||
self.feed_slot_closed_signal_handler,
|
||||
)
|
||||
self.assertTrue(self.feed_slot_closed_received)
|
||||
self.assertTrue(self.feed_exporter_closed_received)
|
||||
|
||||
def test_feed_exporter_signals_sent_deferred(self):
|
||||
self.feed_exporter_closed_received = False
|
||||
self.feed_slot_closed_received = False
|
||||
|
||||
self.run_signaled_feed_exporter(
|
||||
self.feed_exporter_closed_signal_handler_deferred,
|
||||
self.feed_slot_closed_signal_handler_deferred,
|
||||
)
|
||||
self.assertTrue(self.feed_slot_closed_received)
|
||||
self.assertTrue(self.feed_exporter_closed_received)
|
||||
|
||||
|
||||
class FeedExportInitTest(unittest.TestCase):
|
||||
def test_unsupported_storage(self):
|
||||
settings = {
|
||||
|
|
|
|||
Loading…
Reference in New Issue