From e5b23f4b00962df76d8302ebf869ef4a4319e142 Mon Sep 17 00:00:00 2001 From: BroodingKangaroo Date: Wed, 18 Mar 2020 11:26:59 +0300 Subject: [PATCH 1/7] fix #4250: add batch deliveries --- scrapy/extensions/feedexport.py | 54 +++++++++++++++++++++-------- scrapy/settings/default_settings.py | 1 + 2 files changed, 41 insertions(+), 14 deletions(-) diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index 998d2a5d1..906f99fee 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -241,6 +241,7 @@ class FeedExporter: self.storages = self._load_components('FEED_STORAGES') self.exporters = self._load_components('FEED_EXPORTERS') + self.storage_batch = self.settings.getint('FEED_STORAGE_BATCH') for uri, feed in self.feeds.items(): if not self._storage_supported(uri): raise NotConfigured @@ -250,19 +251,7 @@ class FeedExporter: def open_spider(self, spider): for uri, feed in self.feeds.items(): uri = uri % self._get_uri_params(spider, feed['uri_params']) - storage = self._get_storage(uri) - file = storage.open(spider) - exporter = self._get_exporter( - file=file, - format=feed['format'], - fields_to_export=feed['fields'], - encoding=feed['encoding'], - indent=feed['indent'], - ) - slot = _FeedSlot(file, exporter, storage, uri, feed['format'], feed['store_empty']) - self.slots.append(slot) - if slot.store_empty: - slot.start_exporting() + self.slots.append(self._start_new_batch(None, uri, feed, spider)) def close_spider(self, spider): deferred_list = [] @@ -285,11 +274,48 @@ class FeedExporter: deferred_list.append(d) return defer.DeferredList(deferred_list) if deferred_list else None + def _start_new_batch(self, previous_batch_slot, uri, feed, spider): + """ + Redirect the output data stream to a new file. + Execute multiple times if 'FEED_STORAGE_BATCH' setting is greater than zero. + """ + if previous_batch_slot is not None: + previous_batch_slot.exporter.finish_exporting() + previous_batch_slot.storage.store(previous_batch_slot.file) + storage = self._get_storage(uri) + file = storage.open(spider) + exporter = self._get_exporter( + file=file, + format=feed['format'], + fields_to_export=feed['fields'], + encoding=feed['encoding'], + indent=feed['indent'] + ) + slot = _FeedSlot(file, exporter, storage, uri, feed['format'], feed['store_empty']) + if slot.store_empty: + slot.start_exporting() + return slot + + def _get_uri_of_partial(self, slot, feed, spider): + """Get uri for each partial using datetime.now().isoformat()""" + uri = (slot.uri % self._get_uri_params(spider, feed['uri_params'])).split('.')[0] + '.' + uri = uri + datetime.now().isoformat() + '.' + feed['format'] + return uri + def item_scraped(self, item, spider): - for slot in self.slots: + slots = [] + for idx, slot in enumerate(self.slots): slot.start_exporting() slot.exporter.export_item(item) slot.itemcount += 1 + if self.storage_batch and slot.itemcount % self.storage_batch == 0: + uri = self._get_uri_of_partial(slot, self.feeds[slot.uri], spider) + slots.append(self._start_new_batch(slot, uri, self.feeds[slot.uri], spider)) + self.feeds[uri] = self.feeds[slot.uri] + self.feeds.pop(slot.uri) + self.slots[idx] = None + self.slots = [slot for slot in self.slots if slot is not None] + self.slots.extend(slots) def _load_components(self, setting_prefix): conf = without_none_values(self.settings.getwithbase(setting_prefix)) diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 077317c81..690e044c5 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -146,6 +146,7 @@ FEED_STORAGES_BASE = { 's3': 'scrapy.extensions.feedexport.S3FeedStorage', 'ftp': 'scrapy.extensions.feedexport.FTPFeedStorage', } +FEED_STORAGE_BATCH = 0 FEED_EXPORTERS = {} FEED_EXPORTERS_BASE = { 'json': 'scrapy.exporters.JsonItemExporter', From 8b4566ff93843cdf17ada069dc09261a99971d26 Mon Sep 17 00:00:00 2001 From: BroodingKangaroo Date: Wed, 18 Mar 2020 14:21:21 +0300 Subject: [PATCH 2/7] fix wrong name of first file in partial deliveries --- scrapy/extensions/feedexport.py | 13 ++++++++----- 1 file changed, 8 insertions(+), 5 deletions(-) diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index 906f99fee..4f7c6bf07 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -249,6 +249,8 @@ class FeedExporter: raise NotConfigured def open_spider(self, spider): + if self.storage_batch: + self.feeds = {self._get_uri_of_partial(uri, feed, spider): feed for uri, feed in self.feeds.items()} for uri, feed in self.feeds.items(): uri = uri % self._get_uri_params(spider, feed['uri_params']) self.slots.append(self._start_new_batch(None, uri, feed, spider)) @@ -296,11 +298,11 @@ class FeedExporter: slot.start_exporting() return slot - def _get_uri_of_partial(self, slot, feed, spider): + def _get_uri_of_partial(self, template_uri, feed, spider): """Get uri for each partial using datetime.now().isoformat()""" - uri = (slot.uri % self._get_uri_params(spider, feed['uri_params'])).split('.')[0] + '.' - uri = uri + datetime.now().isoformat() + '.' + feed['format'] - return uri + template_uri = (template_uri % self._get_uri_params(spider, feed['uri_params'])) + uri_name = template_uri.split('.')[0] + return '{}.{}.{}'.format(uri_name, datetime.now().isoformat(), feed["format"]) def item_scraped(self, item, spider): slots = [] @@ -309,11 +311,12 @@ class FeedExporter: slot.exporter.export_item(item) slot.itemcount += 1 if self.storage_batch and slot.itemcount % self.storage_batch == 0: - uri = self._get_uri_of_partial(slot, self.feeds[slot.uri], spider) + uri = self._get_uri_of_partial(slot.uri, self.feeds[slot.uri], spider) slots.append(self._start_new_batch(slot, uri, self.feeds[slot.uri], spider)) self.feeds[uri] = self.feeds[slot.uri] self.feeds.pop(slot.uri) self.slots[idx] = None + self.slots = [slot for slot in self.slots if slot is not None] self.slots.extend(slots) From 0723e3f4f9777a87d0df3b2e2fddfeac9099dd3b Mon Sep 17 00:00:00 2001 From: BroodingKangaroo Date: Thu, 19 Mar 2020 21:17:02 +0300 Subject: [PATCH 3/7] add batch_id, add error if uri is specified incorrectly --- scrapy/extensions/feedexport.py | 73 ++++++++++++++++++++--------- scrapy/settings/default_settings.py | 2 +- 2 files changed, 52 insertions(+), 23 deletions(-) diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index 4f7c6bf07..38b25bf4a 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -180,14 +180,16 @@ class FTPFeedStorage(BlockingFeedStorage): class _FeedSlot: - def __init__(self, file, exporter, storage, uri, format, store_empty): + def __init__(self, file, exporter, storage, uri, format, store_empty, batch_id, template_uri): self.file = file self.exporter = exporter self.storage = storage # feed params - self.uri = uri + self.batch_id = batch_id self.format = format self.store_empty = store_empty + self.template_uri = template_uri + self.uri = uri # flags self.itemcount = 0 self._exporting = False @@ -241,19 +243,28 @@ class FeedExporter: self.storages = self._load_components('FEED_STORAGES') self.exporters = self._load_components('FEED_EXPORTERS') - self.storage_batch = self.settings.getint('FEED_STORAGE_BATCH') + self.storage_batch_size = self.settings.getint('FEED_STORAGE_BATCH_SIZE') for uri, feed in self.feeds.items(): if not self._storage_supported(uri): raise NotConfigured + if not self._batch_deliveries_supported(uri): + raise NotConfigured if not self._exporter_supported(feed['format']): raise NotConfigured def open_spider(self, spider): - if self.storage_batch: - self.feeds = {self._get_uri_of_partial(uri, feed, spider): feed for uri, feed in self.feeds.items()} for uri, feed in self.feeds.items(): - uri = uri % self._get_uri_params(spider, feed['uri_params']) - self.slots.append(self._start_new_batch(None, uri, feed, spider)) + batch_id = 1 + uri_params = self._get_uri_params(spider, feed['uri_params']) + uri_params['batch_id'] = batch_id + self.slots.append(self._start_new_batch( + previous_batch_slot=None, + uri=uri % uri_params, + feed=feed, + spider=spider, + batch_id=batch_id, + template_uri=uri + )) def close_spider(self, spider): deferred_list = [] @@ -276,10 +287,17 @@ class FeedExporter: deferred_list.append(d) return defer.DeferredList(deferred_list) if deferred_list else None - def _start_new_batch(self, previous_batch_slot, uri, feed, spider): + def _start_new_batch(self, previous_batch_slot, uri, feed, spider, batch_id, template_uri): """ Redirect the output data stream to a new file. Execute multiple times if 'FEED_STORAGE_BATCH' setting is greater than zero. + :param previous_batch_slot: slot of previous batch. We need to call slot.storage.store + to get the file properly closed. + :param uri: uri of the new batch to start + :param feed: dict with parameters of feed + :param spider: user spider + :param batch_id: sequential batch id starting at 1 + :param template_uri: template uri which contains %(time)s or %(batch_id)s to create new uri """ if previous_batch_slot is not None: previous_batch_slot.exporter.finish_exporting() @@ -293,30 +311,30 @@ class FeedExporter: encoding=feed['encoding'], indent=feed['indent'] ) - slot = _FeedSlot(file, exporter, storage, uri, feed['format'], feed['store_empty']) + slot = _FeedSlot(file, exporter, storage, uri, feed['format'], feed['store_empty'], batch_id, template_uri) if slot.store_empty: slot.start_exporting() return slot - def _get_uri_of_partial(self, template_uri, feed, spider): - """Get uri for each partial using datetime.now().isoformat()""" - template_uri = (template_uri % self._get_uri_params(spider, feed['uri_params'])) - uri_name = template_uri.split('.')[0] - return '{}.{}.{}'.format(uri_name, datetime.now().isoformat(), feed["format"]) - def item_scraped(self, item, spider): slots = [] for idx, slot in enumerate(self.slots): slot.start_exporting() slot.exporter.export_item(item) slot.itemcount += 1 - if self.storage_batch and slot.itemcount % self.storage_batch == 0: - uri = self._get_uri_of_partial(slot.uri, self.feeds[slot.uri], spider) - slots.append(self._start_new_batch(slot, uri, self.feeds[slot.uri], spider)) - self.feeds[uri] = self.feeds[slot.uri] - self.feeds.pop(slot.uri) + if self.storage_batch_size and slot.itemcount % self.storage_batch_size == 0: + batch_id = slot.batch_id + 1 + uri_params = self._get_uri_params(spider, self.feeds[slot.template_uri]['uri_params']) + uri_params['batch_id'] = batch_id + self.slots.append(self._start_new_batch( + previous_batch_slot=slot, + uri=slot.template_uri % uri_params, + feed=self.feeds[slot.template_uri], + spider=spider, + batch_id=batch_id, + template_uri=slot.template_uri + )) self.slots[idx] = None - self.slots = [slot for slot in self.slots if slot is not None] self.slots.extend(slots) @@ -335,6 +353,17 @@ class FeedExporter: return True logger.error("Unknown feed format: %(format)s", {'format': format}) + def _batch_deliveries_supported(self, uri): + """ + If FEED_STORAGE_BATCH_SIZE setting is specified uri has to contain %(time)s or %(batch_id)s + to distinguish different files of partial output + """ + if not self.storage_batch_size: + return True + if '%(time)s' in uri or '%(batch_id)s' in uri: + return True + logger.error('%(time)s or %(batch_id)s must be in uri if FEED_STORAGE_BATCH_SIZE setting is specified') + def _storage_supported(self, uri): scheme = urlparse(uri).scheme if scheme in self.storages: @@ -364,7 +393,7 @@ class FeedExporter: params = {} for k in dir(spider): params[k] = getattr(spider, k) - ts = datetime.utcnow().replace(microsecond=0).isoformat().replace(':', '-') + ts = datetime.utcnow().isoformat().replace(':', '-') params['time'] = ts uripar_function = load_object(uri_params) if uri_params else lambda x, y: None uripar_function(params, spider) diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 690e044c5..7f90a2280 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -146,7 +146,7 @@ FEED_STORAGES_BASE = { 's3': 'scrapy.extensions.feedexport.S3FeedStorage', 'ftp': 'scrapy.extensions.feedexport.FTPFeedStorage', } -FEED_STORAGE_BATCH = 0 +FEED_STORAGE_BATCH_SIZE = 0 FEED_EXPORTERS = {} FEED_EXPORTERS_BASE = { 'json': 'scrapy.exporters.JsonItemExporter', From d11411b402ae68874c6ccc2883836be0b9cf8326 Mon Sep 17 00:00:00 2001 From: BroodingKangaroo Date: Sat, 21 Mar 2020 10:48:13 +0300 Subject: [PATCH 4/7] fix comments --- scrapy/extensions/feedexport.py | 31 ++++++++++++++++++----------- scrapy/settings/default_settings.py | 2 +- 2 files changed, 20 insertions(+), 13 deletions(-) diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index 38b25bf4a..ab0a0de37 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -25,7 +25,6 @@ from scrapy.utils.log import failure_to_exc_info from scrapy.utils.misc import create_instance, load_object from scrapy.utils.python import without_none_values - logger = logging.getLogger(__name__) @@ -243,7 +242,7 @@ class FeedExporter: self.storages = self._load_components('FEED_STORAGES') self.exporters = self._load_components('FEED_EXPORTERS') - self.storage_batch_size = self.settings.getint('FEED_STORAGE_BATCH_SIZE') + self.storage_batch_size = self.settings.get('FEED_STORAGE_BATCH_SIZE', None) for uri, feed in self.feeds.items(): if not self._storage_supported(uri): raise NotConfigured @@ -263,7 +262,7 @@ class FeedExporter: feed=feed, spider=spider, batch_id=batch_id, - template_uri=uri + template_uri=uri, )) def close_spider(self, spider): @@ -290,7 +289,7 @@ class FeedExporter: def _start_new_batch(self, previous_batch_slot, uri, feed, spider, batch_id, template_uri): """ Redirect the output data stream to a new file. - Execute multiple times if 'FEED_STORAGE_BATCH' setting is greater than zero. + Execute multiple times if 'FEED_STORAGE_BATCH' setting is specified. :param previous_batch_slot: slot of previous batch. We need to call slot.storage.store to get the file properly closed. :param uri: uri of the new batch to start @@ -309,9 +308,18 @@ class FeedExporter: format=feed['format'], fields_to_export=feed['fields'], encoding=feed['encoding'], - indent=feed['indent'] + indent=feed['indent'], + ) + slot = _FeedSlot( + file=file, + exporter=exporter, + storage=storage, + uri=uri, + format=feed['format'], + store_empty=feed['store_empty'], + batch_id=batch_id, + template_uri=template_uri, ) - slot = _FeedSlot(file, exporter, storage, uri, feed['format'], feed['store_empty'], batch_id, template_uri) if slot.store_empty: slot.start_exporting() return slot @@ -326,13 +334,13 @@ class FeedExporter: batch_id = slot.batch_id + 1 uri_params = self._get_uri_params(spider, self.feeds[slot.template_uri]['uri_params']) uri_params['batch_id'] = batch_id - self.slots.append(self._start_new_batch( + slots.append(self._start_new_batch( previous_batch_slot=slot, uri=slot.template_uri % uri_params, feed=self.feeds[slot.template_uri], spider=spider, batch_id=batch_id, - template_uri=slot.template_uri + template_uri=slot.template_uri, )) self.slots[idx] = None self.slots = [slot for slot in self.slots if slot is not None] @@ -358,11 +366,10 @@ class FeedExporter: If FEED_STORAGE_BATCH_SIZE setting is specified uri has to contain %(time)s or %(batch_id)s to distinguish different files of partial output """ - if not self.storage_batch_size: + if self.storage_batch_size is None or '%(time)s' in uri or '%(batch_id)s' in uri: return True - if '%(time)s' in uri or '%(batch_id)s' in uri: - return True - logger.error('%(time)s or %(batch_id)s must be in uri if FEED_STORAGE_BATCH_SIZE setting is specified') + logger.warning('%(time)s or %(batch_id)s must be in uri if FEED_STORAGE_BATCH_SIZE setting is specified') + return False def _storage_supported(self, uri): scheme = urlparse(uri).scheme diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 7f90a2280..c3463a505 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -146,7 +146,7 @@ FEED_STORAGES_BASE = { 's3': 'scrapy.extensions.feedexport.S3FeedStorage', 'ftp': 'scrapy.extensions.feedexport.FTPFeedStorage', } -FEED_STORAGE_BATCH_SIZE = 0 +FEED_STORAGE_BATCH_SIZE = None FEED_EXPORTERS = {} FEED_EXPORTERS_BASE = { 'json': 'scrapy.exporters.JsonItemExporter', From 39d0d13d3f7bd671d5b29646b209c62e23373fab Mon Sep 17 00:00:00 2001 From: BroodingKangaroo Date: Thu, 26 Mar 2020 14:18:35 +0300 Subject: [PATCH 5/7] Add partial deliveries tests --- tests/test_feedexport.py | 195 +++++++++++++++++++++++++++++++-------- 1 file changed, 159 insertions(+), 36 deletions(-) diff --git a/tests/test_feedexport.py b/tests/test_feedexport.py index c5589e52f..1ebe44e12 100644 --- a/tests/test_feedexport.py +++ b/tests/test_feedexport.py @@ -6,6 +6,7 @@ import shutil import string import tempfile import warnings +from abc import ABC, abstractmethod from io import BytesIO from pathlib import Path from string import ascii_letters, digits @@ -21,8 +22,9 @@ from zope.interface.verify import verifyObject import scrapy from scrapy.crawler import CrawlerRunner +from scrapy.exceptions import NotConfigured from scrapy.exporters import CsvItemExporter -from scrapy.extensions.feedexport import (BlockingFeedStorage, FileFeedStorage, FTPFeedStorage, +from scrapy.extensions.feedexport import (BlockingFeedStorage, FeedExporter, FileFeedStorage, FTPFeedStorage, IFeedStorage, S3FeedStorage, StdoutFeedStorage) from scrapy.settings import Settings from scrapy.utils.python import to_unicode @@ -76,6 +78,7 @@ class FTPFeedStorageTest(unittest.TestCase): def get_test_spider(self, settings=None): class TestSpider(scrapy.Spider): name = 'test_spider' + crawler = get_crawler(settings_dict=settings) spider = TestSpider.from_crawler(crawler) return spider @@ -129,6 +132,7 @@ class BlockingFeedStorageTest(unittest.TestCase): def get_test_spider(self, settings=None): class TestSpider(scrapy.Spider): name = 'test_spider' + crawler = get_crawler(settings_dict=settings) spider = TestSpider.from_crawler(crawler) return spider @@ -390,23 +394,63 @@ class FromCrawlerFileFeedStorage(FileFeedStorage, FromCrawlerMixin): pass -class FeedExportTest(unittest.TestCase): +class FeedExportTestBase(ABC, unittest.TestCase): + __test__ = False class MyItem(scrapy.Item): foo = scrapy.Field() egg = scrapy.Field() baz = scrapy.Field() + def _random_temp_filename(self, inter_dir=''): + chars = [random.choice(ascii_letters + digits) for _ in range(15)] + filename = ''.join(chars) + return os.path.join(self.temp_dir, inter_dir, filename) + def setUp(self): self.temp_dir = tempfile.mkdtemp() def tearDown(self): shutil.rmtree(self.temp_dir, ignore_errors=True) - def _random_temp_filename(self): - chars = [random.choice(ascii_letters + digits) for _ in range(15)] - filename = ''.join(chars) - return os.path.join(self.temp_dir, filename) + @defer.inlineCallbacks + def exported_data(self, items, settings): + """ + Return exported data which a spider yielding ``items`` would return. + """ + + class TestSpider(scrapy.Spider): + name = 'testspider' + + def parse(self, response): + for item in items: + yield item + + data = yield self.run_and_export(TestSpider, settings) + defer.returnValue(data) + + @defer.inlineCallbacks + def exported_no_data(self, settings): + """ + Return exported data which a spider yielding no ``items`` would return. + """ + + class TestSpider(scrapy.Spider): + name = 'testspider' + + def parse(self, response): + pass + + data = yield self.run_and_export(TestSpider, settings) + defer.returnValue(data) + + @abstractmethod + def run_and_export(self, spider_cls, settings): + pass + + +class FeedExportTest(FeedExportTestBase): + __test__ = True @defer.inlineCallbacks def run_and_export(self, spider_cls, settings): @@ -417,7 +461,6 @@ class FeedExportTest(unittest.TestCase): urljoin('file:', pathname2url(str(file_path))): feed for file_path, feed in FEEDS.items() } - content = {} try: with MockServer() as s: @@ -435,35 +478,6 @@ class FeedExportTest(unittest.TestCase): defer.returnValue(content) - @defer.inlineCallbacks - def exported_data(self, items, settings): - """ - Return exported data which a spider yielding ``items`` would return. - """ - class TestSpider(scrapy.Spider): - name = 'testspider' - - def parse(self, response): - for item in items: - yield item - - data = yield self.run_and_export(TestSpider, settings) - defer.returnValue(data) - - @defer.inlineCallbacks - def exported_no_data(self, settings): - """ - Return exported data which a spider yielding no ``items`` would return. - """ - class TestSpider(scrapy.Spider): - name = 'testspider' - - def parse(self, response): - pass - - data = yield self.run_and_export(TestSpider, settings) - defer.returnValue(data) - @defer.inlineCallbacks def assertExportedCsv(self, items, header, rows, settings=None, ordered=True): settings = settings or {} @@ -970,3 +984,112 @@ class FeedExportTest(unittest.TestCase): } data = yield self.exported_no_data(settings) self.assertEqual(data['csv'], b'') + + +class PartialDeliveriesTest(FeedExportTestBase): + __test__ = True + _file_mark = '_%(time)s_#%(batch_id)s' + + @defer.inlineCallbacks + def run_and_export(self, spider_cls, settings): + """ Run spider with specified settings; return exported data. """ + + FEEDS = settings.get('FEEDS') or {} + settings['FEEDS'] = { + urljoin('file:', file_path): feed + for file_path, feed in FEEDS.items() + } + from collections import defaultdict + content = defaultdict(list) + try: + with MockServer() as s: + runner = CrawlerRunner(Settings(settings)) + spider_cls.start_urls = [s.url('/')] + yield runner.crawl(spider_cls) + + for path, feed in FEEDS.items(): + dir_name = os.path.dirname(path) + for file in sorted(os.listdir(dir_name)): + with open(os.path.join(dir_name, file), 'rb') as f: + data = f.read() + content[feed['format']].append(data) + finally: + pass + defer.returnValue(content) + + @defer.inlineCallbacks + def assertPartialExported(self, items, rows, settings=None): + settings = settings or {} + settings.update({ + 'FEEDS': { + os.path.join(self._random_temp_filename(), 'jl', self._file_mark): {'format': 'jl'}, + }, + }) + data = yield self.exported_data(items, settings) + data['jl'] = b''.join(data['jl']) + parsed = [json.loads(to_unicode(line)) for line in data['jl'].splitlines()] + + rows = [{k: v for k, v in row.items() if v} for row in rows] + self.assertEqual(rows, parsed) + + @defer.inlineCallbacks + def test_partial_deliveries(self): + items = [ + self.MyItem({'foo': 'bar1', 'egg': 'spam1'}), + self.MyItem({'foo': 'bar2', 'egg': 'spam2', 'baz': 'quux2'}), + self.MyItem({'foo': 'bar3', 'baz': 'quux3'}), + ] + rows = [ + {'egg': 'spam1', 'foo': 'bar1', 'baz': ''}, + {'egg': 'spam2', 'foo': 'bar2', 'baz': 'quux2'}, + {'foo': 'bar3', 'baz': 'quux3'} + ] + settings = { + 'FEED_STORAGE_BATCH_SIZE': 1 + } + yield self.assertPartialExported(items, rows, settings=settings) + + def test_wrong_path(self): + settings = { + 'FEEDS': { + self._random_temp_filename(): {'format': 'xml'}, + }, + 'FEED_STORAGE_BATCH_SIZE': 1 + } + crawler = get_crawler(settings_dict=settings) + self.assertRaises(NotConfigured, FeedExporter, crawler) + + @defer.inlineCallbacks + def test_export_no_items_not_store_empty(self): + for fmt in ('json', 'jsonlines', 'xml', 'csv'): + settings = { + 'FEEDS': { + os.path.join(self._random_temp_filename(), fmt, self._file_mark): {'format': fmt}, + }, + 'FEED_STORAGE_BATCH_SIZE': 1 + } + data = yield self.exported_no_data(settings) + data[fmt] = b''.join(data[fmt]) + self.assertEqual(data[fmt], b'') + + @defer.inlineCallbacks + def test_export_no_items_store_empty(self): + formats = ( + ('json', b'[]'), + ('jsonlines', b''), + ('xml', b'\n'), + ('csv', b''), + ) + + for fmt, expctd in formats: + settings = { + 'FEEDS': { + os.path.join(self._random_temp_filename(), fmt, self._file_mark): {'format': fmt}, + }, + 'FEED_STORE_EMPTY': True, + 'FEED_EXPORT_INDENT': None, + 'FEED_STORAGE_BATCH_SIZE': 1 + } + data = yield self.exported_no_data(settings) + data[fmt] = b''.join(data[fmt]) + self.assertEqual(data[fmt], expctd) From ffa8a533e74478a5c81fbf453f2c65601bb1d244 Mon Sep 17 00:00:00 2001 From: BroodingKangaroo Date: Sat, 28 Mar 2020 11:40:16 +0300 Subject: [PATCH 6/7] Set batch_id in _get_uri_params --- scrapy/extensions/feedexport.py | 25 +++++++++++-------------- 1 file changed, 11 insertions(+), 14 deletions(-) diff --git a/scrapy/extensions/feedexport.py b/scrapy/extensions/feedexport.py index ab0a0de37..06ea6c5b2 100644 --- a/scrapy/extensions/feedexport.py +++ b/scrapy/extensions/feedexport.py @@ -253,15 +253,12 @@ class FeedExporter: def open_spider(self, spider): for uri, feed in self.feeds.items(): - batch_id = 1 - uri_params = self._get_uri_params(spider, feed['uri_params']) - uri_params['batch_id'] = batch_id + uri_params = self._get_uri_params(spider, feed['uri_params'], None) self.slots.append(self._start_new_batch( previous_batch_slot=None, uri=uri % uri_params, feed=feed, spider=spider, - batch_id=batch_id, template_uri=uri, )) @@ -286,7 +283,7 @@ class FeedExporter: deferred_list.append(d) return defer.DeferredList(deferred_list) if deferred_list else None - def _start_new_batch(self, previous_batch_slot, uri, feed, spider, batch_id, template_uri): + def _start_new_batch(self, previous_batch_slot, uri, feed, spider, template_uri): """ Redirect the output data stream to a new file. Execute multiple times if 'FEED_STORAGE_BATCH' setting is specified. @@ -295,12 +292,15 @@ class FeedExporter: :param uri: uri of the new batch to start :param feed: dict with parameters of feed :param spider: user spider - :param batch_id: sequential batch id starting at 1 :param template_uri: template uri which contains %(time)s or %(batch_id)s to create new uri """ if previous_batch_slot is not None: + previous_batch_id = previous_batch_slot.batch_id previous_batch_slot.exporter.finish_exporting() previous_batch_slot.storage.store(previous_batch_slot.file) + else: + previous_batch_id = 0 + storage = self._get_storage(uri) file = storage.open(spider) exporter = self._get_exporter( @@ -317,7 +317,7 @@ class FeedExporter: uri=uri, format=feed['format'], store_empty=feed['store_empty'], - batch_id=batch_id, + batch_id=previous_batch_id + 1, template_uri=template_uri, ) if slot.store_empty: @@ -331,15 +331,12 @@ class FeedExporter: slot.exporter.export_item(item) slot.itemcount += 1 if self.storage_batch_size and slot.itemcount % self.storage_batch_size == 0: - batch_id = slot.batch_id + 1 - uri_params = self._get_uri_params(spider, self.feeds[slot.template_uri]['uri_params']) - uri_params['batch_id'] = batch_id + uri_params = self._get_uri_params(spider, self.feeds[slot.template_uri]['uri_params'], slot) slots.append(self._start_new_batch( previous_batch_slot=slot, uri=slot.template_uri % uri_params, feed=self.feeds[slot.template_uri], spider=spider, - batch_id=batch_id, template_uri=slot.template_uri, )) self.slots[idx] = None @@ -396,12 +393,12 @@ class FeedExporter: def _get_storage(self, uri): return self._get_instance(self.storages[urlparse(uri).scheme], uri) - def _get_uri_params(self, spider, uri_params): + def _get_uri_params(self, spider, uri_params, slot): params = {} for k in dir(spider): params[k] = getattr(spider, k) - ts = datetime.utcnow().isoformat().replace(':', '-') - params['time'] = ts + params['batch_id'] = slot.batch_id + 1 if slot is not None else 1 + params['time'] = datetime.utcnow().isoformat().replace(':', '-') uripar_function = load_object(uri_params) if uri_params else lambda x, y: None uripar_function(params, spider) return params From 963580463b96315eb58319e6d35b4cd52672371a Mon Sep 17 00:00:00 2001 From: BroodingKangaroo Date: Wed, 15 Apr 2020 20:14:33 +0300 Subject: [PATCH 7/7] Update tests --- tests/test_feedexport.py | 203 ++++++++++++++++++++++++++++++--------- 1 file changed, 159 insertions(+), 44 deletions(-) diff --git a/tests/test_feedexport.py b/tests/test_feedexport.py index 1ebe44e12..c6cd867b1 100644 --- a/tests/test_feedexport.py +++ b/tests/test_feedexport.py @@ -7,6 +7,7 @@ import string import tempfile import warnings from abc import ABC, abstractmethod +from collections import defaultdict from io import BytesIO from pathlib import Path from string import ascii_letters, digits @@ -444,10 +445,31 @@ class FeedExportTestBase(ABC, unittest.TestCase): data = yield self.run_and_export(TestSpider, settings) defer.returnValue(data) + @defer.inlineCallbacks + def assertExported(self, items, header, rows, settings=None, ordered=True): + yield self.assertExportedCsv(items, header, rows, settings, ordered) + yield self.assertExportedJsonLines(items, rows, settings) + yield self.assertExportedXml(items, rows, settings) + yield self.assertExportedPickle(items, rows, settings) + yield self.assertExportedMarshal(items, rows, settings) + yield self.assertExportedMultiple(items, rows, settings) + @abstractmethod def run_and_export(self, spider_cls, settings): pass + def _load_until_eof(self, data, load_func): + result = [] + with tempfile.TemporaryFile() as temp: + temp.write(data) + temp.seek(0) + while True: + try: + result.append(load_func(temp)) + except EOFError: + break + return result + class FeedExportTest(FeedExportTestBase): __test__ = True @@ -478,6 +500,22 @@ class FeedExportTest(FeedExportTestBase): defer.returnValue(content) + @defer.inlineCallbacks + def exported_data(self, items, settings): + """ + Return exported data which a spider yielding ``items`` would return. + """ + + class TestSpider(scrapy.Spider): + name = 'testspider' + + def parse(self, response): + for item in items: + yield item + + data = yield self.run_and_export(TestSpider, settings) + defer.returnValue(data) + @defer.inlineCallbacks def assertExportedCsv(self, items, header, rows, settings=None, ordered=True): settings = settings or {} @@ -543,18 +581,6 @@ class FeedExportTest(FeedExportTestBase): json_rows = json.loads(to_unicode(data['json'])) self.assertEqual(rows, json_rows) - def _load_until_eof(self, data, load_func): - result = [] - with tempfile.TemporaryFile() as temp: - temp.write(data) - temp.seek(0) - while True: - try: - result.append(load_func(temp)) - except EOFError: - break - return result - @defer.inlineCallbacks def assertExportedPickle(self, items, rows, settings=None): settings = settings or {} @@ -583,15 +609,6 @@ class FeedExportTest(FeedExportTestBase): result = self._load_until_eof(data['marshal'], load_func=marshal.load) self.assertEqual(expected, result) - @defer.inlineCallbacks - def assertExported(self, items, header, rows, settings=None, ordered=True): - yield self.assertExportedCsv(items, header, rows, settings, ordered) - yield self.assertExportedJsonLines(items, rows, settings) - yield self.assertExportedXml(items, rows, settings) - yield self.assertExportedPickle(items, rows, settings) - yield self.assertExportedMarshal(items, rows, settings) - yield self.assertExportedMultiple(items, rows, settings) - @defer.inlineCallbacks def test_export_items(self): # feed exporters use field names from Item @@ -615,7 +632,7 @@ class FeedExportTest(FeedExportTestBase): }, } data = yield self.exported_no_data(settings) - self.assertEqual(data[fmt], b'') + self.assertEqual(b'', data[fmt]) @defer.inlineCallbacks def test_export_no_items_store_empty(self): @@ -635,7 +652,7 @@ class FeedExportTest(FeedExportTestBase): 'FEED_EXPORT_INDENT': None, } data = yield self.exported_no_data(settings) - self.assertEqual(data[fmt], expctd) + self.assertEqual(expctd, data[fmt]) @defer.inlineCallbacks def test_export_multiple_item_classes(self): @@ -734,7 +751,8 @@ class FeedExportTest(FeedExportTestBase): formats = { 'json': u'[{"foo": "Test\\u00d6"}]'.encode('utf-8'), 'jsonlines': u'{"foo": "Test\\u00d6"}\n'.encode('utf-8'), - 'xml': u'\nTest\xd6'.encode('utf-8'), + 'xml': u'\nTest\xd6'.encode( + 'utf-8'), 'csv': u'foo\r\nTest\xd6\r\n'.encode('utf-8'), } @@ -751,7 +769,8 @@ class FeedExportTest(FeedExportTestBase): formats = { 'json': u'[{"foo": "Test\xd6"}]'.encode('latin-1'), 'jsonlines': u'{"foo": "Test\xd6"}\n'.encode('latin-1'), - 'xml': u'\nTest\xd6'.encode('latin-1'), + 'xml': u'\nTest\xd6'.encode( + 'latin-1'), 'csv': u'foo\r\nTest\xd6\r\n'.encode('latin-1'), } @@ -772,7 +791,8 @@ class FeedExportTest(FeedExportTestBase): formats = { 'json': u'[\n{"bar": "BAR"}\n]'.encode('utf-8'), - 'xml': u'\n\n \n FOO\n \n'.encode('latin-1'), + 'xml': u'\n\n \n FOO\n \n'.encode( + 'latin-1'), 'csv': u'bar,foo\r\nBAR,FOO\r\n'.encode('utf-8'), } @@ -988,7 +1008,7 @@ class FeedExportTest(FeedExportTestBase): class PartialDeliveriesTest(FeedExportTestBase): __test__ = True - _file_mark = '_%(time)s_#%(batch_id)s' + _file_mark = '_%(time)s_#%(batch_id)s_' @defer.inlineCallbacks def run_and_export(self, spider_cls, settings): @@ -999,7 +1019,6 @@ class PartialDeliveriesTest(FeedExportTestBase): urljoin('file:', file_path): feed for file_path, feed in FEEDS.items() } - from collections import defaultdict content = defaultdict(list) try: with MockServer() as s: @@ -1014,26 +1033,120 @@ class PartialDeliveriesTest(FeedExportTestBase): data = f.read() content[feed['format']].append(data) finally: - pass + self.tearDown() defer.returnValue(content) @defer.inlineCallbacks - def assertPartialExported(self, items, rows, settings=None): + def assertExportedJsonLines(self, items, rows, settings=None): settings = settings or {} settings.update({ 'FEEDS': { os.path.join(self._random_temp_filename(), 'jl', self._file_mark): {'format': 'jl'}, }, }) - data = yield self.exported_data(items, settings) - data['jl'] = b''.join(data['jl']) - parsed = [json.loads(to_unicode(line)) for line in data['jl'].splitlines()] - + batch_size = settings['FEED_STORAGE_BATCH_SIZE'] rows = [{k: v for k, v in row.items() if v} for row in rows] - self.assertEqual(rows, parsed) + data = yield self.exported_data(items, settings) + for batch in data['jl']: + got_batch = [json.loads(to_unicode(batch_item)) for batch_item in batch.splitlines()] + expected_batch, rows = rows[:batch_size], rows[batch_size:] + self.assertEqual(expected_batch, got_batch) @defer.inlineCallbacks - def test_partial_deliveries(self): + def assertExportedCsv(self, items, header, rows, settings=None, ordered=True): + settings = settings or {} + settings.update({ + 'FEEDS': { + os.path.join(self._random_temp_filename(), 'csv', self._file_mark): {'format': 'csv'}, + }, + }) + batch_size = settings['FEED_STORAGE_BATCH_SIZE'] + data = yield self.exported_data(items, settings) + for batch in data['csv']: + got_batch = csv.DictReader(to_unicode(batch).splitlines()) + self.assertEqual(list(header), got_batch.fieldnames) + expected_batch, rows = rows[:batch_size], rows[batch_size:] + self.assertEqual(expected_batch, list(got_batch)) + + @defer.inlineCallbacks + def assertExportedXml(self, items, rows, settings=None): + settings = settings or {} + settings.update({ + 'FEEDS': { + os.path.join(self._random_temp_filename(), 'xml', self._file_mark): {'format': 'xml'}, + }, + }) + batch_size = settings['FEED_STORAGE_BATCH_SIZE'] + rows = [{k: v for k, v in row.items() if v} for row in rows] + data = yield self.exported_data(items, settings) + for batch in data['xml']: + root = lxml.etree.fromstring(batch) + got_batch = [{e.tag: e.text for e in it} for it in root.findall('item')] + expected_batch, rows = rows[:batch_size], rows[batch_size:] + self.assertEqual(expected_batch, got_batch) + + @defer.inlineCallbacks + def assertExportedMultiple(self, items, rows, settings=None): + settings = settings or {} + settings.update({ + 'FEEDS': { + os.path.join(self._random_temp_filename(), 'xml', self._file_mark): {'format': 'xml'}, + os.path.join(self._random_temp_filename(), 'json', self._file_mark): {'format': 'json'}, + }, + }) + batch_size = settings['FEED_STORAGE_BATCH_SIZE'] + rows = [{k: v for k, v in row.items() if v} for row in rows] + data = yield self.exported_data(items, settings) + # XML + xml_rows = rows.copy() + for batch in data['xml']: + root = lxml.etree.fromstring(batch) + got_batch = [{e.tag: e.text for e in it} for it in root.findall('item')] + expected_batch, xml_rows = xml_rows[:batch_size], xml_rows[batch_size:] + self.assertEqual(expected_batch, got_batch) + # JSON + json_rows = rows.copy() + for batch in data['json']: + got_batch = json.loads(batch) + expected_batch, json_rows = json_rows[:batch_size], json_rows[batch_size:] + self.assertEqual(expected_batch, got_batch) + + @defer.inlineCallbacks + def assertExportedPickle(self, items, rows, settings=None): + settings = settings or {} + settings.update({ + 'FEEDS': { + os.path.join(self._random_temp_filename(), 'pickle', self._file_mark): {'format': 'pickle'}, + }, + }) + batch_size = settings['FEED_STORAGE_BATCH_SIZE'] + rows = [{k: v for k, v in row.items() if v} for row in rows] + data = yield self.exported_data(items, settings) + import pickle + for batch in data['pickle']: + got_batch = self._load_until_eof(batch, load_func=pickle.load) + expected_batch, rows = rows[:batch_size], rows[batch_size:] + self.assertEqual(expected_batch, got_batch) + + @defer.inlineCallbacks + def assertExportedMarshal(self, items, rows, settings=None): + settings = settings or {} + settings.update({ + 'FEEDS': { + os.path.join(self._random_temp_filename(), 'marshal', self._file_mark): {'format': 'marshal'}, + }, + }) + batch_size = settings['FEED_STORAGE_BATCH_SIZE'] + rows = [{k: v for k, v in row.items() if v} for row in rows] + data = yield self.exported_data(items, settings) + import marshal + for batch in data['marshal']: + got_batch = self._load_until_eof(batch, load_func=marshal.load) + expected_batch, rows = rows[:batch_size], rows[batch_size:] + self.assertEqual(expected_batch, got_batch) + + @defer.inlineCallbacks + def test_export_items(self): items = [ self.MyItem({'foo': 'bar1', 'egg': 'spam1'}), self.MyItem({'foo': 'bar2', 'egg': 'spam2', 'baz': 'quux2'}), @@ -1042,14 +1155,16 @@ class PartialDeliveriesTest(FeedExportTestBase): rows = [ {'egg': 'spam1', 'foo': 'bar1', 'baz': ''}, {'egg': 'spam2', 'foo': 'bar2', 'baz': 'quux2'}, - {'foo': 'bar3', 'baz': 'quux3'} + {'foo': 'bar3', 'baz': 'quux3', 'egg': ''} ] settings = { - 'FEED_STORAGE_BATCH_SIZE': 1 + 'FEED_STORAGE_BATCH_SIZE': 2 } - yield self.assertPartialExported(items, rows, settings=settings) + header = self.MyItem.fields.keys() + yield self.assertExported(items, header, rows, settings=settings) def test_wrong_path(self): + """If path without %(time)s or %(batch_id)s an exception must be raised""" settings = { 'FEEDS': { self._random_temp_filename(): {'format': 'xml'}, @@ -1069,8 +1184,8 @@ class PartialDeliveriesTest(FeedExportTestBase): 'FEED_STORAGE_BATCH_SIZE': 1 } data = yield self.exported_no_data(settings) - data[fmt] = b''.join(data[fmt]) - self.assertEqual(data[fmt], b'') + data = dict(data) + self.assertEqual(b'', data[fmt][0]) @defer.inlineCallbacks def test_export_no_items_store_empty(self): @@ -1088,8 +1203,8 @@ class PartialDeliveriesTest(FeedExportTestBase): }, 'FEED_STORE_EMPTY': True, 'FEED_EXPORT_INDENT': None, - 'FEED_STORAGE_BATCH_SIZE': 1 + 'FEED_STORAGE_BATCH_SIZE': 1, } data = yield self.exported_no_data(settings) - data[fmt] = b''.join(data[fmt]) - self.assertEqual(data[fmt], expctd) + data = dict(data) + self.assertEqual(expctd, data[fmt][0])