Added domain schedulers (whose functionality was previously mixed with the

Scrapy Scheduler) and removed domain prioritizers whose functionality became
duplicated by the new domain schedulers.
This commit is contained in:
Pablo Hoffman 2009-06-19 17:55:54 -03:00
parent 3cb12e5dd3
commit 161335fe78
9 changed files with 65 additions and 120 deletions

View File

@ -297,6 +297,16 @@ Default: ``True``
Whether to collect depth stats.
.. setting:: DOMAIN_SCHEDULER
DOMAIN_SCHEDULER
----------------
Default: ``'scrapy.contrib.domainsch.FifoDomainScheduler'``
The Domain Scheduler to use. The domain scheduler returns the next domain
(spider) to scrape.
.. setting:: DOWNLOADER_DEBUG
DOWNLOADER_DEBUG

View File

@ -54,6 +54,8 @@ DEFAULT_SPIDER = None
DEPTH_LIMIT = 0
DEPTH_STATS = True
DOMAIN_SCHEDULER = 'scrapy.contrib.domainsch.FifoDomainScheduler'
DOWNLOAD_DELAY = 0
DOWNLOAD_TIMEOUT = 180 # 3mins
@ -137,8 +139,6 @@ MYSQL_CONNECTION_SETTINGS = {}
NEWSPIDER_MODULE = ''
PRIORITIZER = 'scrapy.core.prioritizers.RandomPrioritizer'
REDIRECT_MAX_METAREFRESH_DELAY = 100
REDIRECT_MAX_TIMES = 20 # uses Firefox default setting
@ -149,7 +149,6 @@ REQUESTS_PER_DOMAIN = 8 # max simultaneous requests per domain
RETRY_TIMES = 2 # initial response + 2 retries = 3 requests
RETRY_HTTP_CODES = ['500', '503', '504', '400', '408']
ROBOTSTXT_OBEY = False
SCHEDULER = 'scrapy.core.scheduler.Scheduler'

View File

@ -0,0 +1,36 @@
"""
Domain Schedulers keep 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 domains to scrape
* remove_pending_domain(domain)
remove domain from pendings, do nothing if not pending
* has_pending_domain(domain)
Return ``True`` if the domain is pending, ``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.remove(domain)
def has_pending_domain(self, domain):
return domain in self.pending_domains

View File

@ -1,37 +0,0 @@
import time
from scrapy.core.exceptions import NotConfigured
from scrapy.store.db import DomainDataHistory
from scrapy.conf import settings
class LessScrapedPrioritizer(object):
"""
A spider prioritizer based on these few simple rules:
1. if spider was never scraped, it has top priority
2. if spider was scraped before, then the less recently the spider
has been scraped, the more priority it has
"""
def __init__(self):
# FIXME this prioritizer must be refactored
raise NotImplemented
if not settings['SCRAPING_DB']:
raise NotConfigured("SCRAPING_DB setting is required")
self.ddh = DomainDataHistory(settings['SCRAPING_DB'], 'domain_data_history')
domains_to_scrape = set(elements)
self.priorities = {}
for domain in domains_to_scrape:
stat = self.ddh.getlast(domain, path="start_time")
if stat and stat[1]:
last_started = stat[1]
# spider is the timestamp of last start time
self.priorities[domain] = time.mktime(last_started.timetuple())
else:
# if domain was never scraped, it has top priority
self.priorities[domain] = 1
def get_priority(self, element):
return self.priorities[element]

View File

@ -155,7 +155,7 @@ class Spiderctl(object):
if "remove_pending_domains" in args:
removed = []
for domain in args["remove_pending_domains"]:
if scrapyengine.scheduler.remove_pending_domain(domain):
if scrapyengine.domain_scheduler.remove_pending_domain(domain):
removed.append(domain)
if removed:
s += "<p>"

View File

@ -25,6 +25,7 @@ from scrapy.spider import spiders
from scrapy.spider.middleware import SpiderMiddlewareManager
from scrapy.utils.defer import chain_deferred, deferred_imap
from scrapy.utils.request import request_info
from scrapy.utils.misc import load_object
class ExecutionEngine(object):
"""
@ -64,6 +65,7 @@ class ExecutionEngine(object):
"""
self.scheduler = scheduler or Scheduler()
self.schedulermiddleware = SchedulerMiddlewareManager(self.scheduler)
self.domain_scheduler = load_object(settings['DOMAIN_SCHEDULER'])()
self.downloader = downloader or Downloader(self)
self.spidermiddleware = SpiderMiddlewareManager()
self._scraping = {}
@ -163,7 +165,7 @@ class ExecutionEngine(object):
return self.scheduler.is_idle() and self.pipeline.is_idle() and self.downloader.is_idle() and not self._scraping
def next_domain(self):
domain = self.scheduler.next_domain()
domain = self.domain_scheduler.next_domain()
if domain:
spider = spiders.fromdomain(domain)
self.open_domain(domain, spider)
@ -224,7 +226,7 @@ class ExecutionEngine(object):
def open_domains(self):
return self.downloader.sites.keys()
def crawl(self, request, spider, domain_priority=0):
def crawl(self, request, spider):
domain = spider.domain_name
def _process_response(response):
@ -282,16 +284,16 @@ class ExecutionEngine(object):
request.deferred.addErrback(lambda _:None)
request.deferred.errback(_failure) # TODO: merge into spider middleware.
schd = self.schedule(request, spider, domain_priority)
schd = self.schedule(request, spider)
schd.addCallbacks(_process_response, _cleanfailure)
return schd
def schedule(self, request, spider, domain_priority=0):
def schedule(self, request, spider):
domain = spider.domain_name
if not self.scheduler.domain_is_open(domain):
if self.debug_mode:
log.msg('Scheduling %s (delayed)' % request_info(request), log.DEBUG)
return self._add_starter(request, spider, domain_priority)
return self._add_starter(request, spider)
if self.debug_mode:
log.msg('Scheduling %s (now)' % request_info(request), log.DEBUG)
schd = self.schedulermiddleware.enqueue_request(domain, request)
@ -311,10 +313,10 @@ class ExecutionEngine(object):
if not self.next_domain():
return self._stop_if_idle()
def _add_starter(self, request, spider, domain_priority):
def _add_starter(self, request, spider):
domain = spider.domain_name
if not self.scheduler.domain_is_pending(domain):
self.scheduler.add_domain(domain, priority=domain_priority)
if not self.domain_scheduler.has_pending_domain(domain):
self.domain_scheduler.add_domain(domain)
self.starters[domain] = []
deferred = defer.Deferred()
self.starters[domain] += [(request, deferred)]

View File

@ -40,7 +40,6 @@ class ExecutionManager(object):
scheduler = load_object(settings['SCHEDULER'])()
scrapyengine.configure(scheduler=scheduler)
self.domainprio = load_object(settings['PRIORITIZER'])()
def crawl(self, *args):
"""Schedule the given args for crawling. args is a list of urls or domains"""
@ -49,9 +48,8 @@ class ExecutionManager(object):
# schedule initial requests to be scraped at engine start
for domain in requests or ():
spider = spiders.fromdomain(domain)
priority = self.domainprio.get_priority(domain)
for request in requests[domain]:
scrapyengine.crawl(request, spider, domain_priority=priority)
scrapyengine.crawl(request, spider)
def runonce(self, *args):
"""Run the engine until it finishes scraping all domains and then exit"""

View File

@ -1,29 +0,0 @@
"""
The Prioritizer is a class which receives a list of elements and prioritizes
it. It's used for defining the order in which domains are to be scraped.
A Prioritizer basically consists of a class which receives a list of elements
in its constructs and contains only one method: get_priority() which returns
the priority of the given element. The element passed to get_priority() must
exists in the list of elements passed in the constructor.
This module contains several basic prioritizers.
For more advanced prioritizers see: scrapy.contrib.prioritizers
"""
import random
class NullPrioritizer(object):
"""
This prioritizer always return the same priority (1)
"""
def get_priority(self, element):
return 1
class RandomPrioritizer(object):
"""
This prioritizer always return a random priority
"""
def get_priority(self, element):
return random.randrange(0, 1000)

View File

@ -16,8 +16,6 @@ class Scheduler(object):
"""
def __init__(self):
self.pending_domains = set()
self.domains_queue = PriorityQueue()
self.pending_requests = {}
self.dfo = settings['SCHEDULER_ORDER'].upper() == 'DFO'
@ -25,37 +23,20 @@ class Scheduler(object):
"""Check if scheduler's resources were allocated for a domain"""
return domain in self.pending_requests
def domain_is_pending(self, domain):
"""Check if a domain is waiting to be scraped in domain's queue."""
return domain in self.pending_domains
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 next_domain(self) :
"""Return next domain available to scrape and remove it from available domains queue"""
if self.pending_domains:
domain = self.domains_queue.pop()[0]
self.pending_domains.remove(domain)
return domain
def add_domain(self, domain, priority=0):
"""Add a new domain to be scraped, with the given priority. If the
domain is already scheduled, it does nothing.
"""
if domain not in self.pending_domains:
self.domains_queue.push(domain, priority)
self.pending_domains.add(domain)
def open_domain(self, domain):
"""Allocates scheduling resources for the given domain"""
Priority = PriorityStack if self.dfo else PriorityQueue
self.pending_requests[domain] = Priority()
def enqueue_request(self, domain, request):
"""Enqueue a request to be downloaded for a domain that is currently being scraped."""
"""Enqueue a request to be downloaded for a domain that is currently
being scraped.
"""
dfd = defer.Deferred()
self.pending_requests[domain].push((request, dfd), request.priority)
return dfd
@ -66,8 +47,8 @@ class Scheduler(object):
Returns a pair ``(request, deferred)`` where ``deferred`` is the
`Deferred` instance returned to the original requester.
``(None, None)`` should be returned if there aren't requests pending
for the domain.
``(None, None)`` is returned if there aren't any request pending for
the given domain.
"""
try:
return self.pending_requests[domain].pop()[0] # [1] is priority
@ -80,21 +61,6 @@ class Scheduler(object):
"""
self.pending_requests.pop(domain, None)
def remove_pending_domain(self, domain):
"""
Remove a pending domain not yet started. If the domain was enqueued
several times, all those instances are removed.
Returns the number of times the domain was enqueued. 0 if domains was
not pending.
If domain is open (not pending) it is not removed and returns None. You
need to call close_domain for open domains.
"""
if domain in self.pending_domains and not self.domain_is_open(domain):
return self.pending_domains.remove(domain)
def is_idle(self):
"""Checks if the schedulers has any request pendings"""
return not self.pending_requests