mirror of https://github.com/scrapy/scrapy.git
* Renamed domain_{opened,closed,idle} signals to spider_{opened,closed,idle}
* Changed them to pass spider instances only (no domains) (refs #105)
This commit is contained in:
parent
a3f2933912
commit
97c322707a
|
|
@ -46,17 +46,19 @@ Exporter to export scraped items to different files, one per spider::
|
|||
class XmlExportPipeline(object):
|
||||
|
||||
def __init__(self):
|
||||
dispatcher.connect(self.domain_opened, signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signals.spider_closed)
|
||||
self.files = {}
|
||||
|
||||
def domain_opened(self, domain):
|
||||
def spider_opened(self, spider):
|
||||
domain = spider.domain_name
|
||||
file = open('%s_products.xml' % domain, 'w+b')
|
||||
self.files[domain] = file
|
||||
self.exporter = XmlItemExporter(file)
|
||||
self.exporter.start_exporting()
|
||||
|
||||
def domain_closed(self, domain):
|
||||
def spider_closed(self, spider):
|
||||
domain = spider.domain_name
|
||||
self.exporter.finish_exporting()
|
||||
file = self.files.pop(domain)
|
||||
file.close()
|
||||
|
|
|
|||
|
|
@ -101,14 +101,14 @@ everytime a domain/spider is opened and closed::
|
|||
class SpiderOpenCloseLogging(object):
|
||||
|
||||
def __init__(self):
|
||||
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
def domain_opened(self, domain, spider):
|
||||
log.msg("opened domain %s" % domain)
|
||||
def spider_opened(self, spider):
|
||||
log.msg("opened spider %s" % spider.domain_name)
|
||||
|
||||
def domain_closed(self, domain, spider):
|
||||
log.msg("closed domain %s" % domain)
|
||||
def spider_closed(self, spider):
|
||||
log.msg("closed spider %s" % spider.domain_name)
|
||||
|
||||
|
||||
.. _topics-extensions-ref-manager:
|
||||
|
|
|
|||
|
|
@ -85,15 +85,15 @@ spider returns multiples items with the same id::
|
|||
|
||||
class DuplicatesPipeline(object):
|
||||
def __init__(self):
|
||||
self.domaininfo = {}
|
||||
dispatcher.connect(self.domain_opened, signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signals.domain_closed)
|
||||
self.duplicates = {}
|
||||
dispatcher.connect(self.spider_opened, signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signals.spider_closed)
|
||||
|
||||
def domain_opened(self, domain):
|
||||
self.duplicates[domain] = set()
|
||||
def spider_opened(self, spider):
|
||||
self.duplicates[spider.domain_name] = set()
|
||||
|
||||
def domain_closed(self, domain):
|
||||
del self.duplicates[domain]
|
||||
def spider_closed(self, spider):
|
||||
del self.duplicates[spider.domain_name]
|
||||
|
||||
def process_item(self, domain, item):
|
||||
if item.id in self.duplicates[domain]:
|
||||
|
|
|
|||
|
|
@ -45,7 +45,7 @@ which quite often consists in answer the question: *which spider is leaking?*.
|
|||
The leak could also come from a custom middleware, pipeline or extension that
|
||||
you have written, if you are not releasing the (previously allocated) resources
|
||||
properly. For example, if you're allocating resources on
|
||||
:signal:`domain_opened` but not releasing them on :signal:`domain_closed`.
|
||||
:signal:`spider_opened` but not releasing them on :signal:`spider_closed`.
|
||||
|
||||
.. _topics-leaks-trackrefs:
|
||||
|
||||
|
|
|
|||
|
|
@ -29,68 +29,58 @@ Built-in signals reference
|
|||
Here's a list of signals used in Scrapy and their meaning, in alphabetical
|
||||
order.
|
||||
|
||||
domain_closed
|
||||
spider_closed
|
||||
-------------
|
||||
|
||||
.. signal:: domain_closed
|
||||
.. function:: domain_closed(domain, spider, reason)
|
||||
.. signal:: spider_closed
|
||||
.. function:: spider_closed(spider, reason)
|
||||
|
||||
Sent after a spider/domain has been closed. This can be used to release
|
||||
per-spider resources reserved on :signal:`domain_opened`.
|
||||
|
||||
:param domain: a string which contains the domain of the spider which has
|
||||
been closed
|
||||
:type domain: str
|
||||
Sent after a spider has been closed. This can be used to release per-spider
|
||||
resources reserved on :signal:`spider_opened`.
|
||||
|
||||
:param spider: the spider which has been closed
|
||||
:type spider: :class:`~scrapy.spider.BaseSpider` object
|
||||
|
||||
:param reason: a string which describes the reason why the domain was closed. If
|
||||
it was closed because the domain has completed scraping, it the reason
|
||||
is ``'finished'``. Otherwise, if the domain was manually closed by
|
||||
calling the ``close_domain`` engine method, then the reason is the one
|
||||
:param reason: a string which describes the reason why the spider was closed. If
|
||||
it was closed because the spider has completed scraping, it the reason
|
||||
is ``'finished'``. Otherwise, if the spider was manually closed by
|
||||
calling the ``close_spider`` engine method, then the reason is the one
|
||||
passed in the ``reason`` argument of that method (which defaults to
|
||||
``'cancelled'``). If the engine was shutdown (for example, by hitting
|
||||
Ctrl-C to stop it) the reason will be ``'shutdown'``.
|
||||
:type reason: str
|
||||
|
||||
domain_opened
|
||||
spider_opened
|
||||
-------------
|
||||
|
||||
.. signal:: domain_opened
|
||||
.. function:: domain_opened(domain, spider)
|
||||
.. signal:: spider_opened
|
||||
.. function:: spider_opened(spider)
|
||||
|
||||
Sent after a spider/domain has been opened for crawling. This is typically
|
||||
used to reserve per-spider resources, but can be used for any task that
|
||||
needs to be performed when a spider/domain is opened.
|
||||
|
||||
:param domain: a string with the domain of the spider which has been opened
|
||||
:type domain: str
|
||||
Sent after a spider has been opened for crawling. This is typically used to
|
||||
reserve per-spider resources, but can be used for any task that needs to be
|
||||
performed when a spider is opened.
|
||||
|
||||
:param spider: the spider which has been opened
|
||||
:type spider: :class:`~scrapy.spider.BaseSpider` object
|
||||
|
||||
domain_idle
|
||||
spider_idle
|
||||
-----------
|
||||
|
||||
.. signal:: domain_idle
|
||||
.. function:: domain_idle(domain, spider)
|
||||
.. signal:: spider_idle
|
||||
.. function:: spider_idle(spider)
|
||||
|
||||
Sent when a domain has gone idle, which means the spider has no further:
|
||||
Sent when a spider has gone idle, which means the spider has no further:
|
||||
|
||||
* requests waiting to be downloaded
|
||||
* requests scheduled
|
||||
* items being processed in the item pipeline
|
||||
|
||||
If the idle state persists after all handlers of this signal have finished,
|
||||
the engine starts closing the domain. After the domain has finished
|
||||
closing, the :signal:`domain_closed` signal is sent.
|
||||
the engine starts closing the spider. After the spider has finished
|
||||
closing, the :signal:`spider_closed` signal is sent.
|
||||
|
||||
You can, for example, schedule some requests in your :signal:`domain_idle`
|
||||
handler to prevent the domain from being closed.
|
||||
|
||||
:param domain: is a string with the domain of the spider which has gone idle
|
||||
:type domain: str
|
||||
You can, for example, schedule some requests in your :signal:`spider_idle`
|
||||
handler to prevent the spider from being closed.
|
||||
|
||||
:param spider: the spider which has gone idle
|
||||
:type spider: :class:`~scrapy.spider.BaseSpider` object
|
||||
|
|
|
|||
|
|
@ -176,7 +176,7 @@ class (which they all inherit from).
|
|||
|
||||
Close the given domain. After this is called, no more specific stats
|
||||
for this domain can be accessed. This method is called automatically on
|
||||
the :signal:`domain_closed` signal.
|
||||
the :signal:`spider_closed` signal.
|
||||
|
||||
Available Stats Collectors
|
||||
==========================
|
||||
|
|
@ -302,7 +302,7 @@ functionality:
|
|||
:type domain: str
|
||||
|
||||
:param reason: the reason why the domain is being closed. See
|
||||
:signal:`domain_closed` signal for more info.
|
||||
:signal:`spider_closed` signal for more info.
|
||||
:type reason: str
|
||||
|
||||
.. signal:: stats_domain_closed
|
||||
|
|
@ -316,7 +316,7 @@ functionality:
|
|||
:type domain: str
|
||||
|
||||
:param reason: the reason why the domain was closed. See
|
||||
:signal:`domain_closed` signal for more info.
|
||||
:signal:`spider_closed` signal for more info.
|
||||
:type reason: str
|
||||
|
||||
:param domain_stats: the stats of the domain just closed.
|
||||
|
|
|
|||
|
|
@ -23,12 +23,12 @@ class CloseDomain(object):
|
|||
self.tasks = {}
|
||||
|
||||
if self.timeout:
|
||||
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
if self.itempassed:
|
||||
dispatcher.connect(self.item_passed, signal=signals.item_passed)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
def domain_opened(self, spider):
|
||||
def spider_opened(self, spider):
|
||||
self.tasks[spider] = reactor.callLater(self.timeout, scrapyengine.close_spider, \
|
||||
spider=spider, reason='closedomain_timeout')
|
||||
|
||||
|
|
@ -37,7 +37,7 @@ class CloseDomain(object):
|
|||
if self.counts[spider] == self.itempassed:
|
||||
scrapyengine.close_spider(spider, 'closedomain_itempassed')
|
||||
|
||||
def domain_closed(self, spider):
|
||||
def spider_closed(self, spider):
|
||||
self.counts.pop(spider, None)
|
||||
tsk = self.tasks.pop(spider, None)
|
||||
if tsk and not tsk.called:
|
||||
|
|
|
|||
|
|
@ -21,19 +21,19 @@ class DelayedCloseDomain(object):
|
|||
raise NotConfigured
|
||||
|
||||
self.opened_at = defaultdict(time)
|
||||
dispatcher.connect(self.domain_idle, signal=signals.domain_idle)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_idle, signal=signals.spider_idle)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
def domain_idle(self, domain):
|
||||
def spider_idle(self, spider):
|
||||
try:
|
||||
lastseen = scrapyengine.downloader.sites[domain].lastseen
|
||||
lastseen = scrapyengine.downloader.sites[spider].lastseen
|
||||
except KeyError:
|
||||
lastseen = None
|
||||
if not lastseen:
|
||||
lastseen = self.opened_at[domain]
|
||||
lastseen = self.opened_at[spider]
|
||||
|
||||
if time() < lastseen + self.delay:
|
||||
raise DontCloseDomain
|
||||
|
||||
def domain_closed(self, domain):
|
||||
self.opened_at.pop(domain, None)
|
||||
def spider_closed(self, spider):
|
||||
self.opened_at.pop(spider, None)
|
||||
|
|
|
|||
|
|
@ -16,13 +16,13 @@ class CookiesMiddleware(object):
|
|||
|
||||
def __init__(self):
|
||||
self.jars = defaultdict(CookieJar)
|
||||
dispatcher.connect(self.domain_closed, signals.domain_closed)
|
||||
dispatcher.connect(self.spider_closed, signals.spider_closed)
|
||||
|
||||
def process_request(self, request, spider):
|
||||
if request.meta.get('dont_merge_cookies', False):
|
||||
return
|
||||
|
||||
jar = self.jars[spider.domain_name]
|
||||
jar = self.jars[spider]
|
||||
cookies = self._get_request_cookies(jar, request)
|
||||
for cookie in cookies:
|
||||
jar.set_cookie_if_ok(cookie, request)
|
||||
|
|
@ -37,14 +37,14 @@ class CookiesMiddleware(object):
|
|||
return response
|
||||
|
||||
# extract cookies from Set-Cookie and drop invalid/expired cookies
|
||||
jar = self.jars[spider.domain_name]
|
||||
jar = self.jars[spider]
|
||||
jar.extract_cookies(response, request)
|
||||
self._debug_set_cookie(response)
|
||||
|
||||
return response
|
||||
|
||||
def domain_closed(self, domain):
|
||||
self.jars.pop(domain, None)
|
||||
def spider_closed(self, spider):
|
||||
self.jars.pop(spider, None)
|
||||
|
||||
def _debug_cookie(self, request):
|
||||
"""log Cookie header for request"""
|
||||
|
|
|
|||
|
|
@ -24,10 +24,10 @@ class HttpCacheMiddleware(object):
|
|||
raise NotConfigured
|
||||
self.cache = Cache(settings['HTTPCACHE_DIR'], sectorize=settings.getbool('HTTPCACHE_SECTORIZE'))
|
||||
self.ignore_missing = settings.getbool('HTTPCACHE_IGNORE_MISSING')
|
||||
dispatcher.connect(self.open_domain, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.open_domain, signal=signals.spider_opened)
|
||||
|
||||
def open_domain(self, domain):
|
||||
self.cache.open_domain(domain)
|
||||
def open_domain(self, spider):
|
||||
self.cache.open_domain(spider.domain_name)
|
||||
|
||||
def process_request(self, request, spider):
|
||||
if not is_cacheable(request):
|
||||
|
|
|
|||
|
|
@ -26,8 +26,8 @@ class RobotsTxtMiddleware(object):
|
|||
self._spider_netlocs = {}
|
||||
self._useragents = {}
|
||||
self._pending = {}
|
||||
dispatcher.connect(self.domain_opened, signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signals.spider_closed)
|
||||
|
||||
def process_request(self, request, spider):
|
||||
useragent = self._useragents[spider]
|
||||
|
|
@ -52,12 +52,12 @@ class RobotsTxtMiddleware(object):
|
|||
rp.parse(response.body.splitlines())
|
||||
self._parsers[urlparse_cached(response).netloc] = rp
|
||||
|
||||
def domain_opened(self, spider):
|
||||
def spider_opened(self, spider):
|
||||
self._spider_netlocs[spider] = set()
|
||||
self._useragents[spider] = getattr(spider, 'user_agent', None) \
|
||||
or settings['USER_AGENT']
|
||||
|
||||
def domain_closed(self, domain, spider):
|
||||
def spider_closed(self, spider):
|
||||
for netloc in self._spider_netlocs[domain]:
|
||||
del self._parsers[netloc]
|
||||
del self._spider_netlocs[domain]
|
||||
|
|
|
|||
|
|
@ -49,7 +49,7 @@ class ItemSamplerPipeline(object):
|
|||
self.items = {}
|
||||
self.domains_count = 0
|
||||
self.empty_domains = set()
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
dispatcher.connect(self.engine_stopped, signal=signals.engine_stopped)
|
||||
|
||||
def process_item(self, item, spider):
|
||||
|
|
@ -70,7 +70,8 @@ class ItemSamplerPipeline(object):
|
|||
if self.empty_domains:
|
||||
log.msg("No products sampled for: %s" % " ".join(self.empty_domains), level=log.WARNING)
|
||||
|
||||
def domain_closed(self, domain, spider, reason):
|
||||
def spider_closed(self, spider, reason):
|
||||
domain = spider.domain_name
|
||||
if reason == 'finished' and not stats.get_value("items_sampled", domain=domain):
|
||||
self.empty_domains.add(domain)
|
||||
self.domains_count += 1
|
||||
|
|
|
|||
|
|
@ -44,10 +44,10 @@ class FSImagesStore(object):
|
|||
self.basedir = basedir
|
||||
self._mkdir(self.basedir)
|
||||
self.created_directories = defaultdict(set)
|
||||
dispatcher.connect(self.domain_closed, signals.domain_closed)
|
||||
dispatcher.connect(self.spider_closed, signals.spider_closed)
|
||||
|
||||
def domain_closed(self, domain):
|
||||
self.created_directories.pop(domain, None)
|
||||
def spider_closed(self, spider):
|
||||
self.created_directories.pop(spider.domain_name, None)
|
||||
|
||||
def persist_image(self, key, image, buf, info):
|
||||
absolute_path = self._get_filesystem_path(key)
|
||||
|
|
|
|||
|
|
@ -22,14 +22,14 @@ class MediaPipeline(object):
|
|||
|
||||
def __init__(self):
|
||||
self.domaininfo = {}
|
||||
dispatcher.connect(self.domain_opened, signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signals.spider_closed)
|
||||
|
||||
def domain_opened(self, spider):
|
||||
def spider_opened(self, spider):
|
||||
self.domaininfo[spider.domain_name] = self.DomainInfo(spider)
|
||||
|
||||
def domain_closed(self, domain):
|
||||
del self.domaininfo[domain]
|
||||
def spider_closed(self, spider):
|
||||
del self.domaininfo[spider.domain_name]
|
||||
|
||||
def process_item(self, domain, item):
|
||||
info = self.domaininfo[domain]
|
||||
|
|
|
|||
|
|
@ -16,14 +16,14 @@ class CachingResolver(object):
|
|||
self.resolver = _CachingThreadedResolver(reactor)
|
||||
reactor.installResolver(self.resolver)
|
||||
dispatcher.connect(self.request_received, signals.request_received)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
def request_received(self, request, spider):
|
||||
url_hostname = urlparse_cached(request).hostname
|
||||
self.spider_hostnames[spider.domain_name].add(url_hostname)
|
||||
self.spider_hostnames[spider].add(url_hostname)
|
||||
|
||||
def domain_closed(self, spider):
|
||||
for hostname in self.spider_hostnames:
|
||||
def spider_closed(self, spider):
|
||||
for hostname in self.spider_hostnames[spider]:
|
||||
self.resolver._cache.pop(hostname, None)
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -15,8 +15,8 @@ class OffsiteMiddleware(object):
|
|||
|
||||
def __init__(self):
|
||||
self.host_regexes = {}
|
||||
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
def process_spider_output(self, response, result, spider):
|
||||
return (x for x in result if not isinstance(x, Request) or \
|
||||
|
|
@ -34,9 +34,9 @@ class OffsiteMiddleware(object):
|
|||
regex = r'^(.*\.)?(%s)$' % '|'.join(domains)
|
||||
return re.compile(regex)
|
||||
|
||||
def domain_opened(self, spider):
|
||||
def spider_opened(self, spider):
|
||||
domains = [spider.domain_name] + spider.extra_domain_names
|
||||
self.host_regexes[spider] = self.get_host_regex(domains)
|
||||
|
||||
def domain_closed(self, spider):
|
||||
def spider_closed(self, spider):
|
||||
del self.host_regexes[spider]
|
||||
|
|
|
|||
|
|
@ -23,44 +23,43 @@ class RequestLimitMiddleware(object):
|
|||
self.max_pending = {}
|
||||
self.dropped_count = {}
|
||||
|
||||
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
def domain_opened(self, domain, spider):
|
||||
self.max_pending[domain] = getattr(spider, 'requests_queue_size', self.max_queue_size)
|
||||
self.dropped_count[domain] = 0
|
||||
def spider_opened(self, spider):
|
||||
self.max_pending[spider] = getattr(spider, 'requests_queue_size', self.max_queue_size)
|
||||
self.dropped_count[spider] = 0
|
||||
|
||||
def domain_closed(self, domain):
|
||||
dropped_count = self.dropped_count[domain]
|
||||
def spider_closed(self, spider):
|
||||
dropped_count = self.dropped_count[spider]
|
||||
if dropped_count:
|
||||
max_pending = self.max_pending[domain]
|
||||
max_pending = self.max_pending[spider]
|
||||
log.msg('Dropped %d request(s) because the scheduler queue size limit (%d requests) was exceeded' % \
|
||||
(dropped_count, max_pending), level=log.DEBUG, domain=domain)
|
||||
del self.dropped_count[domain]
|
||||
del self.max_pending[domain]
|
||||
(dropped_count, max_pending), level=log.DEBUG, spider=spider)
|
||||
del self.dropped_count[spider]
|
||||
del self.max_pending[spider]
|
||||
|
||||
def process_spider_output(self, response, result, spider):
|
||||
domain = spider.domain_name
|
||||
max_pending = self.max_pending.get(domain, 0)
|
||||
max_pending = self.max_pending.get(spider, 0)
|
||||
if max_pending:
|
||||
return imap(lambda v: self._limit_requests(v, domain, max_pending), result)
|
||||
return imap(lambda v: self._limit_requests(v, spider, max_pending), result)
|
||||
else:
|
||||
return result
|
||||
|
||||
def _limit_requests(self, request_or_other, domain, max_pending):
|
||||
def _limit_requests(self, request_or_other, spider, max_pending):
|
||||
if isinstance(request_or_other, Request):
|
||||
free_slots = max_pending - self._pending_count(domain)
|
||||
free_slots = max_pending - self._pending_count(spider)
|
||||
if free_slots > 0:
|
||||
# Scheduler isn't saturated and it is fine to schedule more requests.
|
||||
return request_or_other
|
||||
else:
|
||||
# Skip the request and give engine time to handle other tasks.
|
||||
self.dropped_count[domain] += 1
|
||||
self.dropped_count[spider] += 1
|
||||
return None
|
||||
else:
|
||||
# Return others (non-requests) as is.
|
||||
return request_or_other
|
||||
|
||||
def _pending_count(self, domain):
|
||||
pending = scrapyengine.scheduler.pending_requests.get(domain, [])
|
||||
def _pending_count(self, spider):
|
||||
pending = scrapyengine.scheduler.pending_requests.get(spider, [])
|
||||
return len(pending)
|
||||
|
|
|
|||
|
|
@ -20,49 +20,49 @@ class LiveStats(object):
|
|||
|
||||
def __init__(self):
|
||||
self.domains = {}
|
||||
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
dispatcher.connect(self.item_scraped, signal=signals.item_scraped)
|
||||
dispatcher.connect(self.response_downloaded, signal=signals.response_downloaded)
|
||||
|
||||
dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module)
|
||||
|
||||
def domain_opened(self, domain, spider):
|
||||
def spider_opened(self, spider):
|
||||
pstats = SpiderStats()
|
||||
self.domains[spider.domain_name] = pstats
|
||||
self.domains[spider] = pstats
|
||||
pstats.started = datetime.now().replace(microsecond=0)
|
||||
pstats.finished = None
|
||||
|
||||
def domain_closed(self, domain, spider):
|
||||
self.domains[spider.domain_name].finished = datetime.now().replace(microsecond=0)
|
||||
def spider_closed(self, spider):
|
||||
self.domains[spider].finished = datetime.now().replace(microsecond=0)
|
||||
|
||||
def item_scraped(self, item, spider):
|
||||
self.domains[spider.domain_name].scraped += 1
|
||||
self.domains[spider].scraped += 1
|
||||
|
||||
def response_downloaded(self, response, spider):
|
||||
# sometimes we download responses without opening/closing domains,
|
||||
# for example from scrapy shell
|
||||
if self.domains.get(spider.domain_name):
|
||||
self.domains[spider.domain_name].crawled += 1
|
||||
if self.domains.get(spider):
|
||||
self.domains[spider].crawled += 1
|
||||
|
||||
def webconsole_render(self, wc_request):
|
||||
sch = scrapyengine.scheduler
|
||||
dwl = scrapyengine.downloader
|
||||
|
||||
totdomains = totscraped = totcrawled = totscheduled = totactive = totpending = totdqueued = tottransf = 0
|
||||
totdomains = totscraped = totcrawled = totscheduled = totactive = totdqueued = tottransf = 0
|
||||
s = banner(self)
|
||||
s += "<table border='1'>\n"
|
||||
s += "<tr><th>Domain</th><th>Items<br>Scraped</th><th>Pages<br>Crawled</th><th>Scheduler<br>Pending</th><th>Downloader<br/>Queued</th><th>Downloader<br/>Active</th><th>Downloader<br/>Transferring</th><th>Start time</th><th>Finish time</th><th>Run time</th></tr>\n"
|
||||
for d in sorted(self.domains.keys()):
|
||||
scheduled = len(sch.pending_requests[d]) if d in sch.pending_requests else 0
|
||||
active = len(dwl.sites[d].active) if d in dwl.sites else 0
|
||||
dqueued = len(dwl.sites[d].queue) if d in dwl.sites else 0
|
||||
transf = len(dwl.sites[d].transferring) if d in dwl.sites else 0
|
||||
stats = self.domains[d]
|
||||
for spider in sorted(self.domains.keys()):
|
||||
scheduled = len(sch.pending_requests[spider]) if spider in sch.pending_requests else 0
|
||||
active = len(dwl.sites[spider].active) if spider in dwl.sites else 0
|
||||
dqueued = len(dwl.sites[spider].queue) if spider in dwl.sites else 0
|
||||
transf = len(dwl.sites[spider].transferring) if spider in dwl.sites else 0
|
||||
stats = self.domains[spider]
|
||||
runtime = stats.finished - stats.started if stats.finished else datetime.now() - stats.started
|
||||
|
||||
s += '<tr><td>%s</td><td align="right">%d</td><td align="right">%d</td><td align="right">%d</td><td align="right">%d</td><td align="right">%d</td><td align="right">%d</td><td>%s</td><td>%s</td><td>%s</td></tr>\n' % \
|
||||
(d, stats.scraped, stats.crawled, scheduled, dqueued, active, transf, str(stats.started), str(stats.finished), str(runtime))
|
||||
(spider.domain_name, stats.scraped, stats.crawled, scheduled, dqueued, active, transf, str(stats.started), str(stats.finished), str(runtime))
|
||||
|
||||
totdomains += 1
|
||||
totscraped += stats.scraped
|
||||
|
|
|
|||
|
|
@ -18,16 +18,16 @@ class Spiderctl(object):
|
|||
def __init__(self):
|
||||
self.running = {}
|
||||
self.finished = set()
|
||||
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
from scrapy.management.web import webconsole_discover_module
|
||||
dispatcher.connect(self.webconsole_discover_module, signal=webconsole_discover_module)
|
||||
|
||||
def domain_opened(self, spider):
|
||||
def spider_opened(self, spider):
|
||||
self.running[spider.domain_name] = spider
|
||||
|
||||
def domain_closed(self, spider):
|
||||
def spider_closed(self, spider):
|
||||
del self.running[spider.domain_name]
|
||||
self.finished.add(spider.domain_name)
|
||||
|
||||
|
|
|
|||
|
|
@ -24,8 +24,8 @@ class ShoveItemPipeline(object):
|
|||
self.opts = settings['SHOVEITEM_STORE_OPT'] or {}
|
||||
self.stores = {}
|
||||
|
||||
dispatcher.connect(self.domain_opened, signal=signals.domain_opened)
|
||||
dispatcher.connect(self.domain_closed, signal=signals.domain_closed)
|
||||
dispatcher.connect(self.spider_opened, signal=signals.spider_opened)
|
||||
dispatcher.connect(self.spider_closed, signal=signals.spider_closed)
|
||||
|
||||
def process_item(self, domain, item):
|
||||
guid = str(item.guid)
|
||||
|
|
@ -43,12 +43,13 @@ class ShoveItemPipeline(object):
|
|||
self.log(domain, item, status)
|
||||
return item
|
||||
|
||||
def domain_opened(self, domain):
|
||||
def spider_opened(self, spider):
|
||||
domain = spider.domain_name
|
||||
uri = Template(self.uritpl).substitute(domain=domain)
|
||||
self.stores[domain] = Shove(uri, **self.opts)
|
||||
|
||||
def domain_closed(self, domain):
|
||||
self.stores[domain].sync()
|
||||
def spider_closed(self, spider):
|
||||
self.stores[spider.domain_name].sync()
|
||||
|
||||
def log(self, domain, item, status):
|
||||
log.msg("Shove (%s): Item guid=%s" % (status, item.guid), level=log.DEBUG, domain=domain)
|
||||
|
|
|
|||
|
|
@ -99,7 +99,7 @@ class ExecutionEngine(object):
|
|||
return True
|
||||
|
||||
def next_request(self, spider, now=False):
|
||||
"""Scrape the next request for the domain passed.
|
||||
"""Scrape the next request for the spider passed.
|
||||
|
||||
The next request to be scraped is retrieved from the scheduler and
|
||||
requested from the downloader.
|
||||
|
|
@ -164,7 +164,7 @@ class ExecutionEngine(object):
|
|||
def crawl(self, request, spider):
|
||||
if not request.deferred.callbacks:
|
||||
log.msg("Unable to crawl Request with no callback: %s" % request,
|
||||
level=log.ERROR, domain=spider.domain_name)
|
||||
level=log.ERROR, spider=spider)
|
||||
return
|
||||
schd = mustbe_deferred(self.schedule, request, spider)
|
||||
# FIXME: we can't log errors because we would be preventing them from
|
||||
|
|
@ -186,7 +186,7 @@ class ExecutionEngine(object):
|
|||
return self.scheduler.enqueue_request(spider, request)
|
||||
|
||||
def _mainloop(self):
|
||||
"""Add more domains to be scraped if the downloader has the capacity.
|
||||
"""Add more spiders to be scraped if the downloader has the capacity.
|
||||
|
||||
If there is nothing else scheduled then stop the execution engine.
|
||||
"""
|
||||
|
|
@ -198,15 +198,13 @@ class ExecutionEngine(object):
|
|||
return self._stop_if_idle()
|
||||
|
||||
def download(self, request, spider):
|
||||
domain = spider.domain_name
|
||||
|
||||
def _on_success(response):
|
||||
"""handle the result of a page download"""
|
||||
assert isinstance(response, (Response, Request))
|
||||
if isinstance(response, Response):
|
||||
response.request = request # tie request to response received
|
||||
log.msg(self._crawled_logline(request, response), \
|
||||
level=log.DEBUG, domain=spider.domain_name)
|
||||
level=log.DEBUG, spider=spider)
|
||||
return response
|
||||
elif isinstance(response, Request):
|
||||
newrequest = response
|
||||
|
|
@ -225,7 +223,7 @@ class ExecutionEngine(object):
|
|||
level = log.ERROR
|
||||
if errmsg:
|
||||
log.msg("Crawling <%s>: %s" % (request.url, errmsg), \
|
||||
level=level, domain=domain)
|
||||
level=level, spider=spider)
|
||||
return Failure(IgnoreRequest(str(exc)))
|
||||
|
||||
def _on_complete(_):
|
||||
|
|
@ -241,38 +239,31 @@ class ExecutionEngine(object):
|
|||
return dwld
|
||||
|
||||
def open_spider(self, spider):
|
||||
domain = spider.domain_name
|
||||
log.msg("Domain opened", domain=domain)
|
||||
log.msg("Spider opened", spider=spider)
|
||||
self.next_request(spider)
|
||||
|
||||
self.downloader.open_spider(spider)
|
||||
self.scraper.open_spider(spider)
|
||||
stats.open_domain(domain)
|
||||
stats.open_domain(spider.domain_name)
|
||||
|
||||
# XXX: sent for backwards compatibility (will be removed in Scrapy 0.8)
|
||||
send_catch_log(signals.domain_open, sender=self.__class__, \
|
||||
domain=domain, spider=spider)
|
||||
|
||||
send_catch_log(signals.domain_opened, sender=self.__class__, \
|
||||
domain=domain, spider=spider)
|
||||
send_catch_log(signals.spider_opened, sender=self.__class__, spider=spider)
|
||||
|
||||
def _spider_idle(self, spider):
|
||||
"""Called when a domain gets idle. This function is called when there
|
||||
"""Called when a spider 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
|
||||
(in the domain_idle signal handler) the domain is not closed until the
|
||||
(in the spider_idle signal handler) the spider is not closed until the
|
||||
next loop and this function is guaranteed to be called (at least) once
|
||||
again for this domain.
|
||||
again for this spider.
|
||||
"""
|
||||
domain = spider.domain_name
|
||||
try:
|
||||
dispatcher.send(signal=signals.domain_idle, sender=self.__class__, \
|
||||
domain=domain, spider=spider)
|
||||
dispatcher.send(signal=signals.spider_idle, sender=self.__class__, \
|
||||
spider=spider)
|
||||
except DontCloseDomain:
|
||||
self.next_request(spider)
|
||||
return
|
||||
except:
|
||||
log.err("Exception catched on domain_idle signal dispatch")
|
||||
log.err("Exception catched on spider_idle signal dispatch")
|
||||
if self.spider_is_idle(spider):
|
||||
self.close_spider(spider, reason='finished')
|
||||
|
||||
|
|
@ -283,9 +274,8 @@ class ExecutionEngine(object):
|
|||
|
||||
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)
|
||||
log.msg("Closing spider (%s)" % reason, spider=spider)
|
||||
self.closing[spider] = reason
|
||||
self.downloader.close_spider(spider)
|
||||
self.scheduler.clear_pending_requests(spider)
|
||||
|
|
@ -298,7 +288,7 @@ class ExecutionEngine(object):
|
|||
return dlist
|
||||
|
||||
def _finish_closing_spider_if_idle(self, spider):
|
||||
"""Call _finish_closing_spider if domain is idle"""
|
||||
"""Call _finish_closing_spider if spider is idle"""
|
||||
if self.spider_is_idle(spider) or self.killed:
|
||||
return self._finish_closing_spider(spider)
|
||||
else:
|
||||
|
|
@ -310,15 +300,14 @@ class ExecutionEngine(object):
|
|||
|
||||
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)
|
||||
send_catch_log(signal=signals.spider_closed, sender=self.__class__, \
|
||||
spider=spider, reason=reason)
|
||||
stats.close_domain(spider.domain_name, reason=reason)
|
||||
dfd = defer.maybeDeferred(spiders.close_spider, spider)
|
||||
dfd.addBoth(log.msg, "Domain closed (%s)" % reason, domain=domain)
|
||||
dfd.addBoth(log.msg, "Spider closed (%s)" % reason, spider=spider)
|
||||
reactor.callLater(0, self._mainloop)
|
||||
return dfd
|
||||
|
||||
|
|
|
|||
|
|
@ -7,9 +7,9 @@ signals here without documenting them there.
|
|||
|
||||
engine_started = object()
|
||||
engine_stopped = object()
|
||||
domain_opened = object()
|
||||
domain_idle = object()
|
||||
domain_closed = object()
|
||||
spider_opened = object()
|
||||
spider_idle = object()
|
||||
spider_closed = object()
|
||||
request_received = object()
|
||||
request_uploaded = object()
|
||||
response_received = object()
|
||||
|
|
@ -17,7 +17,3 @@ response_downloaded = object()
|
|||
item_scraped = object()
|
||||
item_passed = object()
|
||||
item_dropped = object()
|
||||
|
||||
# XXX: deprecated signals (will be removed in Scrapy 0.8)
|
||||
domain_open = object()
|
||||
|
||||
|
|
|
|||
|
|
@ -63,26 +63,27 @@ def start(logfile=None, loglevel=None, logstdout=None):
|
|||
file = open(logfile, 'a') if logfile else sys.stderr
|
||||
log.startLogging(file, setStdout=logstdout)
|
||||
|
||||
def msg(message, level=INFO, component=BOT_NAME, domain=None):
|
||||
def msg(message, level=INFO, component=BOT_NAME, domain=None, spider=None):
|
||||
"""Log message according to the level"""
|
||||
if level > log_level:
|
||||
return
|
||||
dispatcher.send(signal=logmessage_received, message=message, level=level, \
|
||||
domain=domain)
|
||||
system = domain if domain else component
|
||||
domain=domain, spider=spider)
|
||||
system = domain or spider.domain_name if spider else component
|
||||
msg_txt = unicode_to_str("%s: %s" % (level_names[level], message))
|
||||
log.msg(msg_txt, system=system)
|
||||
|
||||
def exc(message, level=ERROR, component=BOT_NAME, domain=None):
|
||||
def exc(message, level=ERROR, component=BOT_NAME, domain=None, spider=None):
|
||||
message = message + '\n' + format_exc()
|
||||
msg(message, level, component, domain)
|
||||
msg(message, level, component, domain, spider)
|
||||
|
||||
def err(_stuff=None, _why=None, **kwargs):
|
||||
if ERROR > log_level:
|
||||
return
|
||||
domain = kwargs.pop('domain', None)
|
||||
spider = kwargs.pop('spider', None)
|
||||
component = kwargs.pop('component', BOT_NAME)
|
||||
kwargs['system'] = domain if domain else component
|
||||
kwargs['system'] = domain or spider.domain_name if spider else component
|
||||
if _why:
|
||||
_why = unicode_to_str("ERROR: %s" % _why)
|
||||
log.err(_stuff, _why, **kwargs)
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ class CookiesMiddlewareTest(TestCase):
|
|||
self.mw = CookiesMiddleware()
|
||||
|
||||
def tearDown(self):
|
||||
self.mw.domain_closed('scrapytest.org')
|
||||
self.mw.spider_closed('scrapytest.org')
|
||||
del self.mw
|
||||
|
||||
def test_basic(self):
|
||||
|
|
|
|||
|
|
@ -89,9 +89,9 @@ class CrawlingSession(object):
|
|||
|
||||
dispatcher.connect(self.record_signal, signals.engine_started)
|
||||
dispatcher.connect(self.record_signal, signals.engine_stopped)
|
||||
dispatcher.connect(self.record_signal, signals.domain_opened)
|
||||
dispatcher.connect(self.record_signal, signals.domain_idle)
|
||||
dispatcher.connect(self.record_signal, signals.domain_closed)
|
||||
dispatcher.connect(self.record_signal, signals.spider_opened)
|
||||
dispatcher.connect(self.record_signal, signals.spider_idle)
|
||||
dispatcher.connect(self.record_signal, signals.spider_closed)
|
||||
dispatcher.connect(self.item_scraped, signals.item_scraped)
|
||||
dispatcher.connect(self.request_received, signals.request_received)
|
||||
dispatcher.connect(self.response_downloaded, signals.response_downloaded)
|
||||
|
|
@ -201,16 +201,16 @@ class EngineTest(unittest.TestCase):
|
|||
|
||||
assert signals.engine_started in session.signals_catched
|
||||
assert signals.engine_stopped in session.signals_catched
|
||||
assert signals.domain_opened in session.signals_catched
|
||||
assert signals.domain_idle in session.signals_catched
|
||||
assert signals.domain_closed in session.signals_catched
|
||||
assert signals.spider_opened in session.signals_catched
|
||||
assert signals.spider_idle in session.signals_catched
|
||||
assert signals.spider_closed in session.signals_catched
|
||||
|
||||
self.assertEqual({'domain': session.domain, 'spider': session.spider},
|
||||
session.signals_catched[signals.domain_opened])
|
||||
self.assertEqual({'domain': session.domain, 'spider': session.spider},
|
||||
session.signals_catched[signals.domain_idle])
|
||||
self.assertEqual({'domain': session.domain, 'spider': session.spider, 'reason': 'finished'},
|
||||
session.signals_catched[signals.domain_closed])
|
||||
self.assertEqual({'spider': session.spider},
|
||||
session.signals_catched[signals.spider_opened])
|
||||
self.assertEqual({'spider': session.spider},
|
||||
session.signals_catched[signals.spider_idle])
|
||||
self.assertEqual({'spider': session.spider, 'reason': 'finished'},
|
||||
session.signals_catched[signals.spider_closed])
|
||||
|
||||
if __name__ == "__main__":
|
||||
if len(sys.argv) > 1 and sys.argv[1] == 'runserver':
|
||||
|
|
|
|||
|
|
@ -13,7 +13,7 @@ class TestOffsiteMiddleware(TestCase):
|
|||
self.spider.extra_domain_names = ['scrapy.org']
|
||||
|
||||
self.mw = OffsiteMiddleware()
|
||||
self.mw.domain_opened(self.spider)
|
||||
self.mw.spider_opened(self.spider)
|
||||
|
||||
def test_process_spider_output(self):
|
||||
res = Response('http://scrapytest.org')
|
||||
|
|
@ -28,5 +28,5 @@ class TestOffsiteMiddleware(TestCase):
|
|||
self.assertEquals(out, onsite_reqs)
|
||||
|
||||
def tearDown(self):
|
||||
self.mw.domain_closed(self.spider)
|
||||
self.mw.spider_closed(self.spider)
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue