From 8ddf0811a886bf6a604ae5cf946a7180632b51c4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Daniel=20Gra=C3=B1a?= Date: Tue, 2 Sep 2014 18:07:32 -0300 Subject: [PATCH] Correctly detect when all managed crawlers are done in CrawlerRunner --- scrapy/commands/shell.py | 3 ++- scrapy/crawler.py | 44 ++++++++++++++++++++++++---------------- 2 files changed, 29 insertions(+), 18 deletions(-) diff --git a/scrapy/commands/shell.py b/scrapy/commands/shell.py index ff8c0d156..7c0706482 100644 --- a/scrapy/commands/shell.py +++ b/scrapy/commands/shell.py @@ -53,7 +53,8 @@ class Command(ScrapyCommand): # The crawler is created this way since the Shell manually handles the # crawling engine, so the set up in the crawl method won't work - crawler = self.crawler_process._create_logged_crawler(spidercls) + crawler = self.crawler_process._create_crawler(spidercls) + self.crawler_process._setup_crawler_logging(crawler) # The Shell class needs a persistent engine in the crawler crawler.engine = crawler._create_engine() crawler.engine.start() diff --git a/scrapy/crawler.py b/scrapy/crawler.py index 00de7d0c0..f18760390 100644 --- a/scrapy/crawler.py +++ b/scrapy/crawler.py @@ -76,22 +76,21 @@ class CrawlerRunner(object): smcls = load_object(settings['SPIDER_MANAGER_CLASS']) self.spiders = smcls.from_settings(settings.frozencopy()) self.crawlers = set() - self.crawl_deferreds = set() + self._active = set() def crawl(self, spidercls, *args, **kwargs): - crawler = self._create_logged_crawler(spidercls) - self.crawlers.add(crawler) - - d = crawler.crawl(*args, **kwargs) - self.crawl_deferreds.add(d) - return d - - def _create_logged_crawler(self, spidercls): crawler = self._create_crawler(spidercls) - log_observer = log.start_from_crawler(crawler) - if log_observer: - crawler.signals.connect(log_observer.stop, signals.engine_stopped) - return crawler + self._setup_crawler_logging(crawler) + self.crawlers.add(crawler) + d = crawler.crawl(*args, **kwargs) + self._active.add(d) + + def _done(result): + self.crawlers.discard(crawler) + self._active.discard(d) + return result + + return d.addBoth(_done) def _create_crawler(self, spidercls): if isinstance(spidercls, six.string_types): @@ -100,13 +99,22 @@ class CrawlerRunner(object): crawler_settings = self.settings.copy() spidercls.update_settings(crawler_settings) crawler_settings.freeze() + return Crawler(spidercls, crawler_settings) - crawler = Crawler(spidercls, crawler_settings) - return crawler + def _setup_crawler_logging(self, crawler): + log_observer = log.start_from_crawler(crawler) + if log_observer: + crawler.signals.connect(log_observer.stop, signals.engine_stopped) def stop(self): return defer.DeferredList(c.stop() for c in self.crawlers) + @defer.inlineCallbacks + def join(self): + """Wait for all managed crawlers to complete""" + while self._active: + yield defer.DeferredList(self._active) + class CrawlerProcess(CrawlerRunner): """A class to run multiple scrapy crawlers in a process simultaneously""" @@ -135,13 +143,15 @@ class CrawlerProcess(CrawlerRunner): def start(self, stop_after_crawl=True): if stop_after_crawl: - d = defer.DeferredList(self.crawl_deferreds) + d = self.join() + # Don't start the reactor if the deferreds are already fired if d.called: - # Don't start the reactor if the deferreds are already fired return d.addBoth(lambda _: self._stop_reactor()) + if self.settings.getbool('DNSCACHE_ENABLED'): reactor.installResolver(CachingThreadedResolver(reactor)) + reactor.addSystemEventTrigger('before', 'shutdown', self.stop) reactor.run(installSignalHandlers=False) # blocking call