mirror of https://github.com/scrapy/scrapy.git
Add async def support to pipelines.
This commit is contained in:
parent
8117566974
commit
1f9cef787d
|
|
@ -6,6 +6,8 @@ See documentation in docs/item-pipeline.rst
|
|||
|
||||
from scrapy.middleware import MiddlewareManager
|
||||
from scrapy.utils.conf import build_component_list
|
||||
from scrapy.utils.defer import deferred_f_from_coro_f
|
||||
|
||||
|
||||
|
||||
class ItemPipelineManager(MiddlewareManager):
|
||||
|
|
@ -19,7 +21,7 @@ class ItemPipelineManager(MiddlewareManager):
|
|||
def _add_middleware(self, pipe):
|
||||
super(ItemPipelineManager, self)._add_middleware(pipe)
|
||||
if hasattr(pipe, 'process_item'):
|
||||
self.methods['process_item'].append(pipe.process_item)
|
||||
self.methods['process_item'].append(deferred_f_from_coro_f(pipe.process_item))
|
||||
|
||||
def process_item(self, item, spider):
|
||||
return self._process_chain('process_item', item, spider)
|
||||
|
|
|
|||
|
|
@ -26,6 +26,13 @@ class DeferredPipeline:
|
|||
return d
|
||||
|
||||
|
||||
class AsyncDefPipeline:
|
||||
async def process_item(self, item, spider):
|
||||
await defer.succeed(42)
|
||||
item['pipeline_passed'] = True
|
||||
return item
|
||||
|
||||
|
||||
class ItemSpider(Spider):
|
||||
name = 'itemspider'
|
||||
|
||||
|
|
@ -69,3 +76,9 @@ class PipelineTestCase(unittest.TestCase):
|
|||
crawler = self._create_crawler(DeferredPipeline)
|
||||
yield crawler.crawl(mockserver=self.mockserver)
|
||||
self.assertEqual(len(self.items), 1)
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_asyncdef_pipeline(self):
|
||||
crawler = self._create_crawler(AsyncDefPipeline)
|
||||
yield crawler.crawl(mockserver=self.mockserver)
|
||||
self.assertEqual(len(self.items), 1)
|
||||
|
|
|
|||
Loading…
Reference in New Issue