Big Scrapy core refactoring to pass around spider references instead of domains.

This is to avoid accessing the scrapy.spider.spiders singleton for "resolving"
spiders, which is considered an "evil" practice because it ties us to the
singleton model for the spider resolver, which is a bad thing.

This change will also work as the foundation for the API cleaning that we'll
perform for 0.8. We decided to introduce this change now to have a more common
basecode between 0.7 and 0.8, which will allow us to better support 0.7 until
0.8 is released.

However, this change doesn't modify the stable/documented API, nor does it
change the core logic. Those changes will land on the 0.8 branch, after 0.7 is
released.

--HG--
rename : scrapy/contrib/domainsch.py => scrapy/contrib/spiderscheduler.py
This commit is contained in:
Pablo Hoffman 2009-09-12 14:34:18 -03:00
parent 655cfe138d
commit 921fc4f3bf
28 changed files with 392 additions and 410 deletions

View File

@ -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

View File

@ -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)

View File

@ -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

View File

@ -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

View File

@ -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()

View File

@ -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

View File

@ -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):

View File

@ -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):

View File

@ -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

View File

@ -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]

View File

@ -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]

View File

@ -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

View File

@ -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):

View File

@ -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()

View File

@ -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:

View File

@ -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)

View File

@ -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)

View File

@ -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):

View File

@ -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

View File

@ -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"""

View File

@ -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 ""

View File

@ -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)

View File

@ -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__

View File

@ -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("<h1>(.*?)</h1>", re.M)
price_re = re.compile(">Price: \$(.*?)<", re.M)
def parse(self, response):
xlink = SgmlLinkExtractor()
itemre = re.compile(self.itemurl_re)
for link in xlink.extract_links(response):
if itemre.search(link.url):
yield Request(url=link.url, callback=self.parse_item)
def parse_item(self, response):
item = TestItem()
m = self.name_re.search(response.body)
if m:
item['name'] = m.group(1)
item['url'] = response.url
m = self.price_re.search(response.body)
if m:
item['price'] = m.group(1)
return item
#class TestResource(resource.Resource):
# isLeaf = True
@ -44,21 +80,13 @@ class CrawlingSession(object):
self.port = start_test_site()
self.portno = self.port.getHost().port
from scrapy.spider import spiders
spiders.load(['scrapy.tests.test_spiders'])
self.spider = spiders.fromdomain(self.domain)
self.spider = TestSpider()
if self.spider:
self.spider.start_urls = [
self.geturl("/"),
self.geturl("/redirect"),
]
from scrapy.core import signals
from scrapy.core.manager import scrapymanager
from scrapy.core.engine import scrapyengine
from scrapy.xlib.pydispatch import dispatcher
dispatcher.connect(self.record_signal, signals.engine_started)
dispatcher.connect(self.record_signal, signals.engine_stopped)
dispatcher.connect(self.record_signal, signals.domain_opened)
@ -69,7 +97,7 @@ class CrawlingSession(object):
dispatcher.connect(self.response_downloaded, signals.response_downloaded)
scrapymanager.configure()
scrapymanager.runonce(self.domain)
scrapymanager.runonce(self.spider)
self.port.stopListening()
self.wasrun = True

View File

@ -1,46 +0,0 @@
"""
This is a spider for the unittest sample site.
See scrapy/tests/test_engine.py for more info.
"""
import re
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("<h1>(.*?)</h1>", re.M)
price_re = re.compile(">Price: \$(.*?)<", re.M)
def parse(self, response):
xlink = SgmlLinkExtractor()
itemre = re.compile(self.itemurl_re)
for link in xlink.extract_links(response):
if itemre.search(link.url):
yield Request(url=link.url, callback=self.parse_item)
def parse_item(self, response):
item = TestItem()
m = self.name_re.search(response.body)
if m:
item['name'] = m.group(1)
item['url'] = response.url
m = self.price_re.search(response.body)
if m:
item['price'] = m.group(1)
return item
SPIDER = TestSpider()

View File

@ -18,21 +18,21 @@ def get_engine_status(engine=None):
"engine.scraper.is_idle()",
"len(engine.scraper.sites)",
]
domain_tests = [
"engine.domain_is_idle(domain)",
"engine.closing.get(domain)",
"engine.scheduler.domain_has_pending_requests(domain)",
"len(engine.scheduler.pending_requests[domain])",
"len(engine.downloader.sites[domain].queue)",
"len(engine.downloader.sites[domain].active)",
"len(engine.downloader.sites[domain].transferring)",
"engine.downloader.sites[domain].closing",
"engine.downloader.sites[domain].lastseen",
"len(engine.scraper.sites[domain].queue)",
"len(engine.scraper.sites[domain].active)",
"engine.scraper.sites[domain].active_size",
"engine.scraper.sites[domain].itemproc_size",
"engine.scraper.sites[domain].needs_backout()",
spider_tests = [
"engine.spider_is_idle(spider)",
"engine.closing.get(spider)",
"engine.scheduler.spider_has_pending_requests(spider)",
"len(engine.scheduler.pending_requests[spider])",
"len(engine.downloader.sites[spider].queue)",
"len(engine.downloader.sites[spider].active)",
"len(engine.downloader.sites[spider].transferring)",
"engine.downloader.sites[spider].closing",
"engine.downloader.sites[spider].lastseen",
"len(engine.scraper.sites[spider].queue)",
"len(engine.scraper.sites[spider].active)",
"engine.scraper.sites[spider].active_size",
"engine.scraper.sites[spider].itemproc_size",
"engine.scraper.sites[spider].needs_backout()",
]
s = "Execution engine status\n\n"
@ -43,9 +43,9 @@ def get_engine_status(engine=None):
except Exception, e:
s += "%-47s : %s (exception)\n" % (test, type(e).__name__)
s += "\n"
for domain in engine.downloader.sites:
s += "%s\n" % domain
for test in domain_tests:
for spider in engine.downloader.sites:
s += "Spider: %s\n" % spider
for test in spider_tests:
try:
s += " %-50s : %s\n" % (test, eval(test))
except Exception, e:

17
scrapy/utils/fetch.py Normal file
View File

@ -0,0 +1,17 @@
from scrapy.http import Request
from scrapy.core.manager import scrapymanager
def fetch(urls):
"""Fetch a list of urls and return a list of the downloaded Scrapy
Responses.
This is a blocking function not suitable for calling from spiders. Instead,
it is indended to be called from outside the framework such as Scrapy
commands or standalone scripts.
"""
responses = []
requests = [Request(url, callback=responses.append, dont_filter=True) \
for url in urls]
scrapymanager.runonce(*requests)
return responses