diff --git a/docs/topics/settings.rst b/docs/topics/settings.rst index 35cce502a..79b679fad 100644 --- a/docs/topics/settings.rst +++ b/docs/topics/settings.rst @@ -360,13 +360,13 @@ Whether to collect depth stats. .. setting:: DOMAIN_SCHEDULER -DOMAIN_SCHEDULER +SPIDER_SCHEDULER ---------------- -Default: ``'scrapy.contrib.domainsch.FifoDomainScheduler'`` +Default: ``'scrapy.contrib.spiderscheduler.FifoSpiderScheduler'`` -The Domain Scheduler to use. The domain scheduler returns the next domain -(spider) to scrape. +The Spider Scheduler to use. The spider scheduler returns the next spider to +scrape. .. setting:: DOWNLOADER_DEBUG diff --git a/scrapy/command/commands/fetch.py b/scrapy/command/commands/fetch.py index 6509bf317..bbc2efa9e 100644 --- a/scrapy/command/commands/fetch.py +++ b/scrapy/command/commands/fetch.py @@ -1,7 +1,7 @@ import pprint from scrapy.command import ScrapyCommand -from scrapy.fetcher import fetch +from scrapy.utils.fetch import fetch class Command(ScrapyCommand): @@ -20,11 +20,11 @@ class Command(ScrapyCommand): def add_options(self, parser): ScrapyCommand.add_options(self, parser) parser.add_option("--headers", dest="headers", action="store_true", \ - help="print HTTP headers instead of body") + help="print response HTTP headers instead of body") def run(self, args, opts): - if not args: - print "A URL is required" + if len(args) != 1: + print "One URL is required" return responses = fetch(args) diff --git a/scrapy/command/commands/parse.py b/scrapy/command/commands/parse.py index 1cf1f0f4f..5d5143f08 100644 --- a/scrapy/command/commands/parse.py +++ b/scrapy/command/commands/parse.py @@ -1,5 +1,5 @@ from scrapy.command import ScrapyCommand -from scrapy.fetcher import fetch +from scrapy.utils.fetch import fetch from scrapy.http import Request from scrapy.item import BaseItem from scrapy.spider import spiders diff --git a/scrapy/conf/default_settings.py b/scrapy/conf/default_settings.py index 02d2e712b..17b4f7db0 100644 --- a/scrapy/conf/default_settings.py +++ b/scrapy/conf/default_settings.py @@ -40,7 +40,7 @@ DEFAULT_REQUEST_HEADERS = { DEPTH_LIMIT = 0 DEPTH_STATS = True -DOMAIN_SCHEDULER = 'scrapy.contrib.domainsch.FifoDomainScheduler' +SPIDER_SCHEDULER = 'scrapy.contrib.spiderscheduler.FifoSpiderScheduler' DOWNLOAD_DELAY = 0 DOWNLOAD_TIMEOUT = 180 # 3mins diff --git a/scrapy/contrib/closedomain.py b/scrapy/contrib/closedomain.py index 2bd9b02d8..48f04bf15 100644 --- a/scrapy/contrib/closedomain.py +++ b/scrapy/contrib/closedomain.py @@ -28,17 +28,17 @@ class CloseDomain(object): dispatcher.connect(self.item_passed, signal=signals.item_passed) dispatcher.connect(self.domain_closed, signal=signals.domain_closed) - def domain_opened(self, domain): - self.tasks[domain] = reactor.callLater(self.timeout, scrapyengine.close_domain, \ - domain=domain, reason='closedomain_timeout') + def domain_opened(self, spider): + self.tasks[spider] = reactor.callLater(self.timeout, scrapyengine.close_spider, \ + spider=spider, reason='closedomain_timeout') def item_passed(self, item, spider): - self.counts[spider.domain_name] += 1 - if self.counts[spider.domain_name] == self.itempassed: - scrapyengine.close_domain(spider.domain_name, 'closedomain_itempassed') + self.counts[spider] += 1 + if self.counts[spider] == self.itempassed: + scrapyengine.close_spider(spider, 'closedomain_itempassed') - def domain_closed(self, domain): - self.counts.pop(domain, None) - tsk = self.tasks.pop(domain, None) + def domain_closed(self, spider): + self.counts.pop(spider, None) + tsk = self.tasks.pop(spider, None) if tsk and not tsk.called: tsk.cancel() diff --git a/scrapy/contrib/domainsch.py b/scrapy/contrib/domainsch.py deleted file mode 100644 index c0ffbdcba..000000000 --- a/scrapy/contrib/domainsch.py +++ /dev/null @@ -1,37 +0,0 @@ -""" -The Domain Scheduler keeps track of next domains to scrape. They must implement -the following methods: - -* next_domain() - return next domain to scrape and remove it from pending queue - -* add_domain(domain) - add domain to pending queue - -* remove_pending_domain(domain) - remove (all occurrences) of domain from pending queue, do nothing if not - pending - -* has_pending_domain(domain) - Return ``True`` if the domain is pending to scrape, ``False`` otherwise - -""" - -class FifoDomainScheduler(object): - """Basic domain scheduler based on a FIFO queue""" - - def __init__(self): - self.pending_domains = [] - - def next_domain(self) : - if self.pending_domains: - return self.pending_domains.pop(0) - - def add_domain(self, domain): - self.pending_domains.append(domain) - - def remove_pending_domain(self, domain): - self.pending_domains = [d for d in self.pending_domains if d != domain] - - def has_pending_domain(self, domain): - return domain in self.pending_domains diff --git a/scrapy/contrib/itemsampler.py b/scrapy/contrib/itemsampler.py index a27610de1..fc2d830ee 100644 --- a/scrapy/contrib/itemsampler.py +++ b/scrapy/contrib/itemsampler.py @@ -52,7 +52,8 @@ class ItemSamplerPipeline(object): dispatcher.connect(self.domain_closed, signal=signals.domain_closed) dispatcher.connect(self.engine_stopped, signal=signals.engine_stopped) - def process_item(self, domain, item): + def process_item(self, item, spider): + domain = spider.domain_name sampled = stats.get_value("items_sampled", 0, domain=domain) if sampled < items_per_domain: self.items[item.guid] = item @@ -60,7 +61,7 @@ class ItemSamplerPipeline(object): stats.set_value("items_sampled", sampled, domain=domain) log.msg("Sampled %s" % item, domain=domain, level=log.INFO) if close_domain and sampled == items_per_domain: - scrapyengine.close_domain(domain) + scrapyengine.close_spider(spider) return item def engine_stopped(self): diff --git a/scrapy/contrib/pipeline/__init__.py b/scrapy/contrib/pipeline/__init__.py index cb6f3b326..02d6b83cd 100644 --- a/scrapy/contrib/pipeline/__init__.py +++ b/scrapy/contrib/pipeline/__init__.py @@ -41,10 +41,10 @@ class ItemPipelineManager(object): level=log.DEBUG) self.loaded = True - def open_domain(self, domain): + def open_spider(self, spider): pass - def close_domain(self, domain): + def close_spider(self, spider): pass def process_item(self, item, spider): diff --git a/scrapy/contrib/pipeline/images.py b/scrapy/contrib/pipeline/images.py index de596a3dd..6853696b4 100644 --- a/scrapy/contrib/pipeline/images.py +++ b/scrapy/contrib/pipeline/images.py @@ -23,6 +23,7 @@ from scrapy.utils.misc import md5sum from scrapy.core import signals from scrapy.core.engine import scrapyengine from scrapy.core.exceptions import DropItem, NotConfigured +from scrapy.spider import BaseSpider from scrapy.contrib.pipeline.media import MediaPipeline from scrapy.http import Request from scrapy.conf import settings @@ -76,6 +77,25 @@ class FSImagesStore(object): seen.add(dirname) +class _S3AmazonAWSSpider(BaseSpider): + """This spider is used for uploading images to Amazon S3 + + It is basically not a crawling spider like a normal spider is, this spider is + a placeholder that allows us to open a different slot in downloader and use it + for uploads to S3. + + The use of another downloader slot for S3 images avoid the effect of normal + spider downloader slot to be affected by requests to a complete different + domain (s3.amazonaws.com). + + It means that a spider that uses download_delay or alike is not going to be + delayed even more because it is uploading images to s3. + """ + domain_name = "s3.amazonaws.com" + start_urls = ['http://s3.amazonaws.com/'] + max_concurrent_requests = 100 + + class S3ImagesStore(object): request_priority = 1000 @@ -86,10 +106,9 @@ class S3ImagesStore(object): self._set_custom_spider() def _set_custom_spider(self): - domain = settings['IMAGES_S3STORE_SPIDER'] - if domain: - from scrapy.spider import spiders - self.s3_spider = spiders.fromdomain(domain) + use_custom_spider = bool(settings['IMAGES_S3STORE_SPIDER']) + if use_custom_spider: + self.s3_spider = _S3AmazonAWSSpider() else: self.s3_spider = None diff --git a/scrapy/contrib/pipeline/media.py b/scrapy/contrib/pipeline/media.py index 39388205e..5225cdf65 100644 --- a/scrapy/contrib/pipeline/media.py +++ b/scrapy/contrib/pipeline/media.py @@ -6,7 +6,6 @@ from scrapy import log from scrapy.core import signals from scrapy.core.engine import scrapyengine from scrapy.utils.request import request_fingerprint -from scrapy.spider import spiders from scrapy.utils.misc import arg_to_iter @@ -14,9 +13,9 @@ class MediaPipeline(object): DOWNLOAD_PRIORITY = 1000 class DomainInfo(object): - def __init__(self, domain): - self.domain = domain - self.spider = spiders.fromdomain(domain) + def __init__(self, spider): + self.domain = spider.domain_name + self.spider = spider self.downloading = {} self.downloaded = {} self.waiting = {} @@ -26,8 +25,8 @@ class MediaPipeline(object): dispatcher.connect(self.domain_opened, signals.domain_opened) dispatcher.connect(self.domain_closed, signals.domain_closed) - def domain_opened(self, domain): - self.domaininfo[domain] = self.DomainInfo(domain) + def domain_opened(self, spider): + self.domaininfo[spider.domain_name] = self.DomainInfo(spider) def domain_closed(self, domain): del self.domaininfo[domain] diff --git a/scrapy/contrib/spidermanager.py b/scrapy/contrib/spidermanager.py index b09debe89..11a2edf12 100644 --- a/scrapy/contrib/spidermanager.py +++ b/scrapy/contrib/spidermanager.py @@ -80,6 +80,8 @@ class TwistedPluginSpiderManager(object): """Reload spider module to release any resources held on to by the spider """ + if domain not in self._spiders: + return spider = self._spiders[domain] module_name = spider.__module__ module = sys.modules[module_name] diff --git a/scrapy/contrib/spiderscheduler.py b/scrapy/contrib/spiderscheduler.py new file mode 100644 index 000000000..9c569b6c9 --- /dev/null +++ b/scrapy/contrib/spiderscheduler.py @@ -0,0 +1,37 @@ +""" +The Spider Scheduler keeps track of next spiders to scrape. They must implement +the following methods: + +* next_spider() + return next spider to scrape and remove it from pending queue + +* add_spider(spider) + add spider to pending queue + +* remove_pending_spider(spider) + remove (all occurrences) of spider from pending queue, do nothing if not + pending + +* has_pending_spider(spider) + Return ``True`` if the spider is pending to scrape, ``False`` otherwise + +""" + +class FifoSpiderScheduler(object): + """Basic spider scheduler based on a FIFO queue""" + + def __init__(self): + self.pending_spiders = [] + + def next_spider(self) : + if self.pending_spiders: + return self.pending_spiders.pop(0) + + def add_spider(self, spider): + self.pending_spiders.append(spider) + + def remove_pending_spider(self, spider): + self.pending_spiders = [d for d in self.pending_spiders if d != spider] + + def has_pending_spider(self, spider): + return spider in self.pending_spiders diff --git a/scrapy/core/downloader/manager.py b/scrapy/core/downloader/manager.py index 59490b338..c82c9a8f3 100644 --- a/scrapy/core/downloader/manager.py +++ b/scrapy/core/downloader/manager.py @@ -8,7 +8,6 @@ from twisted.internet import reactor, defer from twisted.python.failure import Failure from scrapy.core.exceptions import IgnoreRequest -from scrapy.spider import spiders from scrapy.conf import settings from scrapy.utils.defer import mustbe_deferred from scrapy import log @@ -48,7 +47,7 @@ class SiteInfo(object): class Downloader(object): """Mantain many concurrent downloads and provide an HTTP abstraction. - It supports a limited number of connections per domain and many domains in + It supports a limited number of connections per spider and many spiders in parallel. """ @@ -64,14 +63,14 @@ class Downloader(object): Response object, then request never reach downloader queue, and it will not be downloaded from site. """ - site = self.sites[spider.domain_name] + site = self.sites[spider] if site.closing: - raise IgnoreRequest('Cannot fetch on a closing domain') + raise IgnoreRequest('Cannot fetch on a closing spider') site.active.add(request) def _deactivate(_): site.active.remove(request) - self._close_if_idle(spider.domain_name) + self._close_if_idle(spider) return _ dfd = self.middleware.download(self.enqueue, request, spider) @@ -79,7 +78,7 @@ class Downloader(object): def enqueue(self, request, spider): """Enqueue a Request for a effective download from site""" - site = self.sites[spider.domain_name] + site = self.sites[spider] if site.closing: raise IgnoreRequest deferred = defer.Deferred() @@ -89,8 +88,7 @@ class Downloader(object): def _process_queue(self, spider): """Effective download requests from site queue""" - domain = spider.domain_name - site = self.sites.get(domain) + site = self.sites.get(spider) if not site: return @@ -112,12 +110,12 @@ class Downloader(object): dfd = self._download(site, request, spider) dfd.chainDeferred(deferred) - self._close_if_idle(domain) + self._close_if_idle(spider) - def _close_if_idle(self, domain): - site = self.sites.get(domain) + def _close_if_idle(self, spider): + site = self.sites.get(spider) if site and site.closing and not site.active: - del self.sites[domain] + del self.sites[spider] def _download(self, site, request, spider): # The order is very important for the following deferreds. Do not change! @@ -136,35 +134,35 @@ class Downloader(object): # avoid partially downloaded responses from propagating to the # downloader middleware, to speed-up the closing process if site.closing: - log.msg("Crawled while closing domain: %s" % request, \ + log.msg("Crawled while closing spider: %s" % request, \ level=log.DEBUG) raise IgnoreRequest return _ return dfd.addBoth(finish_transferring) - def open_domain(self, domain): - """Allocate resources to begin processing a domain""" - if domain in self.sites: - raise RuntimeError('Downloader domain already opened: %s' % domain) + def open_spider(self, spider): + """Allocate resources to begin processing a spider""" + domain = spider.domain_name + if spider in self.sites: + raise RuntimeError('Downloader spider already opened: %s' % domain) - spider = spiders.fromdomain(domain) - self.sites[domain] = SiteInfo( + self.sites[spider] = SiteInfo( download_delay=getattr(spider, 'download_delay', None), max_concurrent_requests=getattr(spider, 'max_concurrent_requests', None) ) - def close_domain(self, domain): - """Free any resources associated with the given domain""" - site = self.sites.get(domain) + def close_spider(self, spider): + """Free any resources associated with the given spider""" + domain = spider.domain_name + site = self.sites.get(spider) if not site or site.closing: - raise RuntimeError('Downloader domain already closed: %s' % domain) + raise RuntimeError('Downloader spider already closed: %s' % domain) site.closing = True - spider = spiders.fromdomain(domain) self._process_queue(spider) def has_capacity(self): - """Does the downloader have capacity to handle more domains""" + """Does the downloader have capacity to handle more spiders""" return len(self.sites) < self.concurrent_domains def is_idle(self): diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 60a0893be..44f982e02 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -28,7 +28,7 @@ class ExecutionEngine(object): def __init__(self): self.configured = False self.keep_alive = False - self.closing = {} # dict (domain -> reason) of spiders being closed + self.closing = {} # dict (spider -> reason) of spiders being closed self.running = False self.killed = False self.paused = False @@ -40,7 +40,7 @@ class ExecutionEngine(object): Configure execution engine with the given scheduling policy and downloader. """ self.scheduler = load_object(settings['SCHEDULER'])() - self.domain_scheduler = load_object(settings['DOMAIN_SCHEDULER'])() + self.spider_scheduler = load_object(settings['SPIDER_SCHEDULER'])() self.downloader = Downloader() self.scraper = Scraper(self) self.configured = True @@ -60,9 +60,9 @@ class ExecutionEngine(object): if not self.running: return self.running = False - for domain in self.open_domains: + for spider in self.open_spiders: reactor.addSystemEventTrigger('before', 'shutdown', \ - self.close_domain, domain, reason='shutdown') + self.close_spider, spider, reason='shutdown') if self._mainloop_task.running: self._mainloop_task.stop() try: @@ -90,92 +90,90 @@ class ExecutionEngine(object): return self.scheduler.is_idle() and self.downloader.is_idle() and \ self.scraper.is_idle() - def next_domain(self): - domain = self.domain_scheduler.next_domain() - if domain: - self.open_domain(domain) - return domain + def next_spider(self): + spider = self.spider_scheduler.next_spider() + if spider: + self.open_spider(spider) + return True - def next_request(self, domain, now=False): + def next_request(self, spider, now=False): """Scrape the next request for the domain passed. The next request to be scraped is retrieved from the scheduler and requested from the downloader. - The domain is closed if there are no more pages to scrape. + The spider is closed if there are no more pages to scrape. """ if now: - self._next_request_pending.discard(domain) - elif domain not in self._next_request_pending: - self._next_request_pending.add(domain) - return reactor.callLater(0, self.next_request, domain, now=True) + self._next_request_pending.discard(spider) + elif spider not in self._next_request_pending: + self._next_request_pending.add(spider) + return reactor.callLater(0, self.next_request, spider, now=True) else: return if self.paused: - return reactor.callLater(5, self.next_request, domain) + return reactor.callLater(5, self.next_request, spider) - while not self._needs_backout(domain): - if not self._next_request(domain): + while not self._needs_backout(spider): + if not self._next_request(spider): break - if self.domain_is_idle(domain): - self._domain_idle(domain) + if self.spider_is_idle(spider): + self._spider_idle(spider) - def _needs_backout(self, domain): + def _needs_backout(self, spider): return not self.running \ - or self.domain_is_closed(domain) \ - or self.downloader.sites[domain].needs_backout() \ - or self.scraper.sites[domain].needs_backout() + or self.spider_is_closed(spider) \ + or self.downloader.sites[spider].needs_backout() \ + or self.scraper.sites[spider].needs_backout() - def _next_request(self, domain): + def _next_request(self, spider): # Next pending request from scheduler - request, deferred = self.scheduler.next_request(domain) + request, deferred = self.scheduler.next_request(spider) if request: - spider = spiders.fromdomain(domain) dwld = mustbe_deferred(self.download, request, spider) dwld.chainDeferred(deferred).addBoth(lambda _: deferred) dwld.addErrback(log.err, "Unhandled error on engine._next_request") return dwld - def domain_is_idle(self, domain): - scraper_idle = domain in self.scraper.sites \ - and self.scraper.sites[domain].is_idle() - pending = self.scheduler.domain_has_pending_requests(domain) - downloading = domain in self.downloader.sites \ - and self.downloader.sites[domain].active + def spider_is_idle(self, spider): + scraper_idle = spider in self.scraper.sites \ + and self.scraper.sites[spider].is_idle() + pending = self.scheduler.spider_has_pending_requests(spider) + downloading = spider in self.downloader.sites \ + and self.downloader.sites[spider].active return scraper_idle and not (pending or downloading) - def domain_is_closed(self, domain): - """Return True if the domain is fully closed (ie. not even in the + def spider_is_closed(self, spider): + """Return True if the spider is fully closed (ie. not even in the closing stage)""" - return domain not in self.downloader.sites + return spider not in self.downloader.sites - def domain_is_open(self, domain): - """Return True if the domain is fully opened (ie. not in closing + def spider_is_open(self, spider): + """Return True if the spider is fully opened (ie. not in closing stage)""" - return domain in self.downloader.sites and domain not in self.closing + return spider in self.downloader.sites and spider not in self.closing @property - def open_domains(self): + def open_spiders(self): return self.downloader.sites.keys() def crawl(self, request, spider): schd = mustbe_deferred(self.schedule, request, spider) schd.addBoth(self.scraper.enqueue_scrape, request, spider) schd.addErrback(log.err, "Unhandled error on engine.crawl()") - schd.addBoth(lambda _: self.next_request(spider.domain_name)) + schd.addBoth(lambda _: self.next_request(spider)) def schedule(self, request, spider): - domain = spider.domain_name - if domain in self.closing: + if spider in self.closing: raise IgnoreRequest() - if not self.scheduler.domain_is_open(domain): - self.scheduler.open_domain(domain) - if self.domain_is_closed(domain): # scheduler auto-open - self.domain_scheduler.add_domain(domain) - self.next_request(domain) - return self.scheduler.enqueue_request(domain, request) + if not self.scheduler.spider_is_open(spider): + self.scheduler.open_spider(spider) + if self.spider_is_closed(spider): # scheduler auto-open + self.spider_scheduler.add_spider(spider) + self.next_request(spider) + return self.scheduler.enqueue_request(spider, request) def _mainloop(self): """Add more domains to be scraped if the downloader has the capacity. @@ -186,7 +184,7 @@ class ExecutionEngine(object): return while self.running and self.downloader.has_capacity(): - if not self.next_domain(): + if not self.next_spider(): return self._stop_if_idle() def download(self, request, spider): @@ -222,7 +220,7 @@ class ExecutionEngine(object): return Failure(IgnoreRequest(str(exc))) def _on_complete(_): - self.next_request(domain) + self.next_request(spider) return _ dwld = mustbe_deferred(self.downloader.fetch, request, spider) @@ -230,13 +228,13 @@ class ExecutionEngine(object): dwld.addBoth(_on_complete) return dwld - def open_domain(self, domain): + def open_spider(self, spider): + domain = spider.domain_name log.msg("Domain opened", domain=domain) - spider = spiders.fromdomain(domain) - self.next_request(domain) + self.next_request(spider) - self.downloader.open_domain(domain) - self.scraper.open_domain(domain) + self.downloader.open_spider(spider) + self.scraper.open_spider(spider) stats.open_domain(domain) # XXX: sent for backwards compatibility (will be removed in Scrapy 0.8) @@ -246,7 +244,7 @@ class ExecutionEngine(object): send_catch_log(signals.domain_opened, sender=self.__class__, \ domain=domain, spider=spider) - def _domain_idle(self, domain): + def _spider_idle(self, spider): """Called when a domain gets idle. This function is called when there are no remaining pages to download or schedule. It can be called multiple times. If some extension raises a DontCloseDomain exception @@ -254,57 +252,58 @@ class ExecutionEngine(object): next loop and this function is guaranteed to be called (at least) once again for this domain. """ - spider = spiders.fromdomain(domain) + domain = spider.domain_name try: dispatcher.send(signal=signals.domain_idle, sender=self.__class__, \ domain=domain, spider=spider) except DontCloseDomain: - self.next_request(domain) + self.next_request(spider) return except: log.err("Exception catched on domain_idle signal dispatch") - if self.domain_is_idle(domain): - self.close_domain(domain, reason='finished') + if self.spider_is_idle(spider): + self.close_spider(spider, reason='finished') def _stop_if_idle(self): """Call the stop method if the system has no outstanding tasks. """ if self.is_idle() and not self.keep_alive: self.stop() - def close_domain(self, domain, reason='cancelled'): - """Close (cancel) domain and clear all its outstanding requests""" - if domain not in self.closing: + def close_spider(self, spider, reason='cancelled'): + """Close (cancel) spider and clear all its outstanding requests""" + domain = spider.domain_name + if spider not in self.closing: log.msg("Closing domain (%s)" % reason, domain=domain) - self.closing[domain] = reason - self.downloader.close_domain(domain) - self.scheduler.clear_pending_requests(domain) - return self._finish_closing_domain_if_idle(domain) + self.closing[spider] = reason + self.downloader.close_spider(spider) + self.scheduler.clear_pending_requests(spider) + return self._finish_closing_spider_if_idle(spider) return defer.succeed(None) - def _finish_closing_domain_if_idle(self, domain): - """Call _finish_closing_domain if domain is idle""" - if self.domain_is_idle(domain) or self.killed: - self._finish_closing_domain(domain) + def _finish_closing_spider_if_idle(self, spider): + """Call _finish_closing_spider if domain is idle""" + if self.spider_is_idle(spider) or self.killed: + self._finish_closing_spider(spider) else: dfd = defer.Deferred() - dfd.addCallback(self._finish_closing_domain_if_idle) + dfd.addCallback(self._finish_closing_spider_if_idle) delay = 5 if self.running else 1 - reactor.callLater(delay, dfd.callback, domain) + reactor.callLater(delay, dfd.callback, spider) return dfd - def _finish_closing_domain(self, domain): - """This function is called after the domain has been closed""" - spider = spiders.fromdomain(domain) - self.scheduler.close_domain(domain) - self.scraper.close_domain(domain) - reason = self.closing.pop(domain, 'finished') + def _finish_closing_spider(self, spider): + """This function is called after the spider has been closed""" + domain = spider.domain_name + self.scheduler.close_spider(spider) + self.scraper.close_spider(spider) + reason = self.closing.pop(spider, 'finished') send_catch_log(signal=signals.domain_closed, sender=self.__class__, \ domain=domain, spider=spider, reason=reason) stats.close_domain(domain, reason=reason) - log.msg("Domain closed (%s)" % reason, domain=domain) spiders.close_domain(domain) + log.msg("Domain closed (%s)" % reason, domain=domain) self._mainloop() - if not self.open_domains: + if not self.open_spiders: send_catch_log(signal=signals.engine_stopped, sender=self.__class__) scrapyengine = ExecutionEngine() diff --git a/scrapy/core/manager.py b/scrapy/core/manager.py index 8560cc929..2078aaa95 100644 --- a/scrapy/core/manager.py +++ b/scrapy/core/manager.py @@ -1,4 +1,5 @@ import signal +from collections import defaultdict from twisted.internet import reactor @@ -6,50 +7,44 @@ from scrapy.extension import extensions from scrapy import log from scrapy.http import Request from scrapy.core.engine import scrapyengine -from scrapy.spider import spiders +from scrapy.spider import BaseSpider, spiders from scrapy.utils.misc import arg_to_iter from scrapy.utils.url import is_url from scrapy.utils.ossignal import install_shutdown_handlers, signal_names -def _parse_args(args): - """Parse crawl arguments and return a dict of domains -> list of requests""" - requests, urls, sites = set(), set(), set() - for a in args: - if isinstance(a, Request): - requests.add(a) - elif is_url(a): - urls.add(a) +def _get_spider_requests(*args): + """Collect requests and spiders from the given arguments. Returns a dict of + spider -> list of requests + """ + spider_requests = defaultdict(list) + for arg in args: + if isinstance(arg, tuple): + request, spider = arg + spider_requests[spider] = request + elif isinstance(arg, Request): + spider = spiders.fromurl(arg.url) or BaseSpider('default') + if spider: + spider_requests[spider] += [arg] + else: + log.msg('Could not find spider for request: %s' % arg, log.ERROR) + elif isinstance(arg, BaseSpider): + spider_requests[arg] += arg.start_requests() + elif is_url(arg): + spider = spiders.fromurl(arg) or BaseSpider('default') + if spider: + for req in arg_to_iter(spider.make_requests_from_url(arg)): + spider_requests[spider] += [req] + else: + log.msg('Could not find spider for url: %s' % arg, log.ERROR) + elif isinstance(arg, basestring): + spider = spiders.fromdomain(arg) + if spider: + spider_requests[spider] += spider.start_requests() + else: + log.msg('Could not find spider for domain: %s' % arg, log.ERROR) else: - sites.add(a) - - perdomain = {} - - # sites - for domain in sites: - spider = spiders.fromdomain(domain) - if not spider: - log.msg('Could not find spider for %s' % domain, log.ERROR) - continue - reqs = spider.start_requests() - perdomain.setdefault(domain, []).extend(reqs) - - # urls - for url in urls: - spider = spiders.fromurl(url) - if spider: - for req in arg_to_iter(spider.make_requests_from_url(url)): - perdomain.setdefault(spider.domain_name, []).append(req) - else: - log.msg('Could not find spider for <%s>' % url, log.ERROR) - - # requests - for request in requests: - spider = spiders.fromurl(request.url) - if not spider: - log.msg('Could not find spider for %s' % request, log.ERROR) - continue - perdomain.setdefault(spider.domain_name, []).append(request) - return perdomain + raise TypeError("Unsupported argument: %r" % arg) + return spider_requests class ExecutionManager(object): @@ -85,17 +80,14 @@ class ExecutionManager(object): def crawl(self, *args): """Schedule the given args for crawling. args is a list of urls or domains""" - - requests = _parse_args(args) - # schedule initial requests to be scraped at engine start - for domain in requests or (): - spider = spiders.fromdomain(domain) - for request in requests[domain]: + assert self.configured, "Scrapy Manager not yet configured" + spider_requests = _get_spider_requests(*args) + for spider, requests in spider_requests.iteritems(): + for request in requests: scrapyengine.crawl(request, spider) def runonce(self, *args): """Run the engine until it finishes scraping all domains and then exit""" - assert self.configured, "Scrapy Manger not yet configured" self.crawl(*args) scrapyengine.start() if self.control_reactor: @@ -103,7 +95,6 @@ class ExecutionManager(object): def start(self): """Start the scrapy server, without scheduling any domains""" - assert self.configured, "Scrapy Manger not yet configured" scrapyengine.keep_alive = True scrapyengine.start() if self.control_reactor: diff --git a/scrapy/core/scheduler/middleware.py b/scrapy/core/scheduler/middleware.py index 483a7aa5d..7587f8504 100644 --- a/scrapy/core/scheduler/middleware.py +++ b/scrapy/core/scheduler/middleware.py @@ -46,29 +46,30 @@ class SchedulerMiddlewareManager(object): self.loaded = True def _add_middleware(self, mw): - for name in ('enqueue_request', 'open_domain', 'close_domain'): + for name in ['enqueue_request', 'open_domain', 'close_domain']: mwfunc = getattr(mw, name, None) if mwfunc: self.mw_cbs[name].append(mwfunc) - def enqueue_request(self, wrappedfunc, domain, request): + def enqueue_request(self, wrappedfunc, spider, request): def _enqueue_request(request): for mwfunc in self.mw_cbs['enqueue_request']: - result = mwfunc(domain=domain, request=request) + result = mwfunc(domain=spider.domain_name, request=request) assert result is None or isinstance(result, (Response, Deferred)), \ 'Middleware %s.enqueue_request must return None, Response or Deferred, got %s' % \ (mwfunc.im_self.__class__.__name__, result.__class__.__name__) if result: return result - return wrappedfunc(domain=domain, request=request) + return wrappedfunc(spider=spider, request=request) deferred = mustbe_deferred(_enqueue_request, request) return deferred - def open_domain(self, domain): + def open_spider(self, spider): for mwfunc in self.mw_cbs['open_domain']: - mwfunc(domain) + mwfunc(spider.domain_name) - def close_domain(self, domain): + def close_spider(self, spider): for mwfunc in self.mw_cbs['close_domain']: - mwfunc(domain) + mwfunc(spider.domain_name) + diff --git a/scrapy/core/scheduler/schedulers.py b/scrapy/core/scheduler/schedulers.py index 22af11b02..a30ae9d98 100644 --- a/scrapy/core/scheduler/schedulers.py +++ b/scrapy/core/scheduler/schedulers.py @@ -23,60 +23,60 @@ class Scheduler(object): self.dfo = settings['SCHEDULER_ORDER'].upper() == 'DFO' self.middleware = SchedulerMiddlewareManager() - def domain_is_open(self, domain): - """Check if scheduler's resources were allocated for a domain""" - return domain in self.pending_requests + def spider_is_open(self, spider): + """Check if scheduler's resources were allocated for a spider""" + return spider in self.pending_requests - def domain_has_pending_requests(self, domain): - """Check if are there pending requests for a domain""" - if domain in self.pending_requests: - return bool(self.pending_requests[domain]) + def spider_has_pending_requests(self, spider): + """Check if are there pending requests for a spider""" + if spider in self.pending_requests: + return bool(self.pending_requests[spider]) - def open_domain(self, domain): - """Allocates scheduling resources for the given domain""" - if domain in self.pending_requests: - raise RuntimeError('Scheduler domain already opened: %s' % domain) + def open_spider(self, spider): + """Allocates scheduling resources for the given spider""" + if spider in self.pending_requests: + raise RuntimeError('Scheduler spider already opened: %s' % spider) Priority = PriorityStack if self.dfo else PriorityQueue - self.pending_requests[domain] = Priority() - self.middleware.open_domain(domain) + self.pending_requests[spider] = Priority() + self.middleware.open_spider(spider) - def close_domain(self, domain): + def close_spider(self, spider): """Called when a spider has finished scraping to free any resources - associated with the domain. + associated with the spider. """ - if domain not in self.pending_requests: - raise RuntimeError('Scheduler domain is not open: %s' % domain) - self.middleware.close_domain(domain) - self.pending_requests.pop(domain, None) + if spider not in self.pending_requests: + raise RuntimeError('Scheduler spider is not open: %s' % spider) + self.middleware.close_spider(spider) + self.pending_requests.pop(spider, None) - def enqueue_request(self, domain, request): - """Enqueue a request to be downloaded for a domain that is currently being scraped.""" - return self.middleware.enqueue_request(self._enqueue_request, domain, request) + def enqueue_request(self, spider, request): + """Enqueue a request to be downloaded for a spider that is currently being scraped.""" + return self.middleware.enqueue_request(self._enqueue_request, spider, request) - def _enqueue_request(self, domain, request): + def _enqueue_request(self, spider, request): dfd = defer.Deferred() - self.pending_requests[domain].push((request, dfd), -request.priority) + self.pending_requests[spider].push((request, dfd), -request.priority) return dfd - def clear_pending_requests(self, domain): - """Remove all pending requests for the given domain""" - q = self.pending_requests[domain] + def clear_pending_requests(self, spider): + """Remove all pending requests for the given spider""" + q = self.pending_requests[spider] while q: _, dfd = q.pop()[0] dfd.errback(Failure(IgnoreRequest())) - def next_request(self, domain): - """Return the next available request to be downloaded for a domain. + def next_request(self, spider): + """Return the next available request to be downloaded for a spider. Returns a pair ``(request, deferred)`` where ``deferred`` is the `Deferred` instance returned to the original requester. ``(None, None)`` is returned if there aren't any request pending for - the given domain. + the given spider. """ try: - return self.pending_requests[domain].pop()[0] # [1] is priority + return self.pending_requests[spider].pop()[0] # [1] is priority except (KeyError, IndexError): return (None, None) diff --git a/scrapy/core/scraper.py b/scrapy/core/scraper.py index bed4dd877..aa2473787 100644 --- a/scrapy/core/scraper.py +++ b/scrapy/core/scraper.py @@ -16,7 +16,7 @@ from scrapy import log from scrapy.stats import stats from scrapy.conf import settings -class SiteInfo(object): +class SpiderInfo(object): """Object for holding data of the responses being scraped""" MIN_RESPONSE_SIZE = 1024 @@ -64,26 +64,26 @@ class Scraper(object): self.concurrent_items = settings.getint('CONCURRENT_ITEMS') self.engine = engine - def open_domain(self, domain): - """Open the given domain for scraping and allocate resources for it""" - if domain in self.sites: - raise RuntimeError('Scraper domain already opened: %s' % domain) - self.sites[domain] = SiteInfo() - self.itemproc.open_domain(domain) + def open_spider(self, spider): + """Open the given spider for scraping and allocate resources for it""" + if spider in self.sites: + raise RuntimeError('Scraper spider already opened: %s' % spider) + self.sites[spider] = SpiderInfo() + self.itemproc.open_spider(spider) - def close_domain(self, domain): - """Close a domain being scraped and release its resources""" - if domain not in self.sites: - raise RuntimeError('Scraper domain already closed: %s' % domain) - self.sites.pop(domain) - self.itemproc.open_domain(domain) + def close_spider(self, spider): + """Close a spider being scraped and release its resources""" + if spider not in self.sites: + raise RuntimeError('Scraper spider already closed: %s' % spider) + self.sites.pop(spider) + self.itemproc.close_spider(spider) def is_idle(self): """Return True if there isn't any more spiders to process""" return not self.sites def enqueue_scrape(self, response, request, spider): - site = self.sites[spider.domain_name] + site = self.sites[spider] dfd = site.add_response_request(response, request) # FIXME: this can't be called here because the stats domain may be # already closed @@ -154,7 +154,7 @@ class Scraper(object): """ # TODO: keep closing state internally instead of checking engine domain = spider.domain_name - if domain in self.engine.closing: + if spider in self.engine.closing: return elif isinstance(output, Request): send_catch_log(signal=signals.request_received, request=output, \ @@ -165,7 +165,7 @@ class Scraper(object): domain=domain) send_catch_log(signal=signals.item_scraped, sender=self.__class__, \ item=output, spider=spider, response=response) - self.sites[domain].itemproc_size += 1 + self.sites[spider].itemproc_size += 1 # FIXME: this can't be called here because the stats domain may be # already closed #stats.max_value('scraper/max_itemproc_size', \ @@ -197,7 +197,7 @@ class Scraper(object): """ItemProcessor finished for the given ``item`` and returned ``output`` """ domain = spider.domain_name - self.sites[domain].itemproc_size -= 1 + self.sites[spider].itemproc_size -= 1 if isinstance(output, Failure): ex = output.value if isinstance(ex, DropItem): diff --git a/scrapy/fetcher.py b/scrapy/fetcher.py deleted file mode 100644 index a24192503..000000000 --- a/scrapy/fetcher.py +++ /dev/null @@ -1,30 +0,0 @@ -from urlparse import urlparse - -from scrapy.spider import spiders -from scrapy.http import Request -from scrapy.core.manager import scrapymanager -from scrapy.spider import BaseSpider - -def fetch(urls): - """Download the given urls and return a list of the successfully downloaded - responses. - - Suitable for for calling from a script, shouldn't be called from spiders. - """ - map(get_or_create_spider, urls) - responses = [] - requests = [Request(url, callback=responses.append, dont_filter=True) \ - for url in urls] - scrapymanager.runonce(*requests) - return responses - -def get_or_create_spider(url): - # XXX: hack to allow downloading pages from unknown domains - spider = spiders.fromurl(url) - if not spider: - domain = urlparse(url).hostname - spider = BaseSpider() - spider.domain_name = domain - spiders.add_spider(spider) - return spider - diff --git a/scrapy/http/request/__init__.py b/scrapy/http/request/__init__.py index 28c173505..8cd55fe99 100644 --- a/scrapy/http/request/__init__.py +++ b/scrapy/http/request/__init__.py @@ -91,15 +91,9 @@ class Request(object_ref): return "<%s %s>" % (self.method, self.url) def __repr__(self): - d = { - 'method': self.method, - 'url': self.url, - 'headers': self.headers, - 'body': self.body, - 'cookies': self.cookies, - 'meta': self.meta, - } - return "%s(%s)" % (self.__class__.__name__, repr(d)) + attrs = ['url', 'method', 'body', 'headers', 'cookies', 'meta'] + args = ", ".join(["%s=%r" % (a, getattr(self, a)) for a in attrs]) + return "%s(%s)" % (self.__class__.__name__, args) def copy(self): """Return a copy of this Request""" diff --git a/scrapy/http/response/__init__.py b/scrapy/http/response/__init__.py index 9e4994642..11a9bfcde 100644 --- a/scrapy/http/response/__init__.py +++ b/scrapy/http/response/__init__.py @@ -60,15 +60,9 @@ class Response(object_ref): body = property(_get_body, _set_body) def __repr__(self): - d = { - 'status': self.status, - 'url': self.url, - 'headers': self.headers, - 'body': self.body, - 'meta': self.meta, - 'flags': self.flags, - } - return "%s(%s)" % (self.__class__.__name__, repr(d)) + attrs = ['url', 'status', 'body', 'headers', 'meta', 'flags'] + args = ", ".join(["%s=%r" % (a, getattr(self, a)) for a in attrs]) + return "%s(%s)" % (self.__class__.__name__, args) def __str__(self): flags = "(%s) " % ",".join(self.flags) if self.flags else "" diff --git a/scrapy/shell.py b/scrapy/shell.py index 199a91437..9c0298502 100644 --- a/scrapy/shell.py +++ b/scrapy/shell.py @@ -11,7 +11,7 @@ import signal from twisted.internet import reactor, threads -from scrapy.spider import spiders +from scrapy.spider import BaseSpider, spiders from scrapy.selector import XmlXPathSelector, HtmlXPathSelector from scrapy.utils.misc import load_object from scrapy.utils.response import open_in_browser @@ -19,7 +19,6 @@ from scrapy.conf import settings from scrapy.core.manager import scrapymanager from scrapy.core.engine import scrapyengine from scrapy.http import Request -from scrapy.fetcher import get_or_create_spider def relevant_var(varname): return varname not in ['shelp', 'fetch', 'view', '__builtins__', 'In', \ @@ -53,7 +52,7 @@ class Shell(object): else: url = parse_url(request_or_url) request = Request(url) - spider = get_or_create_spider(url) + spider = spiders.fromurl(url) or BaseSpider('default') print "Fetching %s..." % request response = threads.blockingCallFromThread(reactor, scrapyengine.schedule, \ request, spider) diff --git a/scrapy/spider/models.py b/scrapy/spider/models.py index a065f484f..bbec56a6c 100644 --- a/scrapy/spider/models.py +++ b/scrapy/spider/models.py @@ -47,10 +47,21 @@ class BaseSpider(object): implements(ISpider) - start_urls = [] + # XXX: class attributes kept for backwards compatibility domain_name = None + start_urls = [] extra_domain_names = [] + def __init__(self, domain_name=None): + if domain_name is not None: + self.domain_name = domain_name + # XXX: create instance attributes (class attributes were kept for + # backwards compatibility) + if not self.start_urls: + self.start_urls = [] + if not self.extra_domain_names: + self.extra_domain_names = [] + def log(self, message, level=log.DEBUG): """Log the given messages at the given log level. Always use this method to send log messages from your spider @@ -71,3 +82,8 @@ class BaseSpider(object): requests, although it can be overrided in descendant spiders. """ pass + + def __str__(self): + return "<%s %r>" % (type(self).__name__, self.domain_name) + + __repr__ = __str__ diff --git a/scrapy/tests/test_engine.py b/scrapy/tests/test_engine.py index 4f96413b5..c4b51b0fd 100644 --- a/scrapy/tests/test_engine.py +++ b/scrapy/tests/test_engine.py @@ -2,15 +2,51 @@ Scrapy engine tests """ -import sys -import os -import urlparse -import unittest +import sys, os, re, urlparse, unittest from twisted.internet import reactor from twisted.web import server, resource, static, util +from scrapy.core import signals +from scrapy.core.manager import scrapymanager +from scrapy.xlib.pydispatch import dispatcher from scrapy.tests import tests_datadir +from scrapy.spider import BaseSpider +from scrapy.item import Item, Field +from scrapy.contrib.linkextractors.sgml import SgmlLinkExtractor +from scrapy.http import Request + +class TestItem(Item): + name = Field() + url = Field() + price = Field() + +class TestSpider(BaseSpider): + domain_name = "scrapytest.org" + extra_domain_names = ["localhost"] + start_urls = ['http://localhost'] + + itemurl_re = re.compile("item\d+.html") + name_re = re.compile("