From 6d989e3fb01c8a184b84d34a3c083ed40006b0ae Mon Sep 17 00:00:00 2001 From: Pablo Hoffman Date: Sun, 31 Jul 2011 03:32:25 -0300 Subject: [PATCH] imported patch scheduler_single_spider.patch --- scrapy/contrib/dupefilter.py | 38 +++++++++---------------- scrapy/core/engine.py | 27 ++++++++---------- scrapy/core/scheduler.py | 54 ++++++++++++++---------------------- scrapy/utils/engine.py | 1 - 4 files changed, 45 insertions(+), 75 deletions(-) diff --git a/scrapy/contrib/dupefilter.py b/scrapy/contrib/dupefilter.py index c6b9d185a..7bb449665 100644 --- a/scrapy/contrib/dupefilter.py +++ b/scrapy/contrib/dupefilter.py @@ -1,14 +1,8 @@ """ Dupe Filter classes implement a mechanism for filtering duplicate requests. -They must implement the following methods: +They must implement the following method: -* open_spider(spider) - open a spider for tracking duplicates (typically used to reserve resources) - -* close_spider(spider) - close a spider (typically used for freeing resources) - -* request_seen(spider, request, dont_record=False) +* request_seen(request, dont_record=False) return ``True`` if the request was seen before, or ``False`` otherwise. If ``dont_record`` is ``True`` the request must not be recorded as seen. @@ -17,33 +11,27 @@ They must implement the following methods: from scrapy.utils.request import request_fingerprint -class NullDupeFilter(dict): - def open_spider(self, spider): - pass +class BaseDupeFilter(object): - def close_spider(self, spider): - pass + @classmethod + def from_settings(cls, settings): + return cls() - def request_seen(self, spider, request, dont_record=False): + def request_seen(self, request, dont_record=False): return False -class RequestFingerprintDupeFilter(object): +class RequestFingerprintDupeFilter(BaseDupeFilter): """Duplicate filter using scrapy.utils.request.request_fingerprint""" def __init__(self): - self.fingerprints = {} + super(RequestFingerprintDupeFilter, self).__init__() + self.fingerprints = set() - def open_spider(self, spider): - self.fingerprints[spider] = set() - - def close_spider(self, spider): - del self.fingerprints[spider] - - def request_seen(self, spider, request, dont_record=False): + def request_seen(self, request, dont_record=False): fp = request_fingerprint(request) - if fp in self.fingerprints[spider]: + if fp in self.fingerprints: return True if not dont_record: - self.fingerprints[spider].add(fp) + self.fingerprints.add(fp) return False diff --git a/scrapy/core/engine.py b/scrapy/core/engine.py index 5d674700b..9d2e89f7b 100644 --- a/scrapy/core/engine.py +++ b/scrapy/core/engine.py @@ -17,18 +17,18 @@ from scrapy.exceptions import DontCloseSpider from scrapy.http import Response, Request from scrapy.utils.misc import load_object from scrapy.utils.signal import send_catch_log, send_catch_log_deferred -from scrapy.utils.defer import mustbe_deferred from scrapy.utils.reactor import CallLaterOnce class Slot(object): - def __init__(self, start_requests, close_if_idle, nextcall): + def __init__(self, start_requests, close_if_idle, nextcall, scheduler): self.closing = False self.inprogress = set() # requests in progress self.start_requests = iter(start_requests) self.close_if_idle = close_if_idle self.nextcall = nextcall + self.scheduler = scheduler def add_request(self, request): self.inprogress.add(request) @@ -56,7 +56,7 @@ class ExecutionEngine(object): self.slots = {} self.running = False self.paused = False - self.scheduler = load_object(settings['SCHEDULER'])() + self.scheduler_cls = load_object(settings['SCHEDULER']) self.downloader = Downloader(self.settings) self.scraper = Scraper(self, self.settings) self._concurrent_spiders = settings.getint('CONCURRENT_SPIDERS') @@ -85,10 +85,6 @@ class ExecutionEngine(object): """Resume the execution engine""" self.paused = False - def is_idle(self): - return self.scheduler.is_idle() and self.downloader.is_idle() and \ - self.scraper.is_idle() - def _next_request(self, spider): try: slot = self.slots[spider] @@ -121,13 +117,13 @@ class ExecutionEngine(object): or self.scraper.slots[spider].needs_backout() def _next_request_from_scheduler(self, spider): - request = self.scheduler.next_request(spider) + slot = self.slots[spider] + request = slot.scheduler.next_request() if not request: return d = self._download(request, spider) d.addBoth(self._handle_downloader_output, request, spider) d.addErrback(log.msg, spider=spider) - slot = self.slots[spider] d.addBoth(lambda _: slot.remove_request(request)) d.addErrback(log.msg, spider=spider) d.addBoth(lambda _: slot.nextcall.schedule()) @@ -148,7 +144,7 @@ class ExecutionEngine(object): def spider_is_idle(self, spider): scraper_idle = spider in self.scraper.slots \ and self.scraper.slots[spider].is_idle() - pending = self.scheduler.spider_has_pending_requests(spider) + pending = self.slots[spider].scheduler.has_pending_requests() downloading = bool(self.downloader.slots) idle = scraper_idle and not (pending or downloading) return idle @@ -168,7 +164,7 @@ class ExecutionEngine(object): self.slots[spider].nextcall.schedule() def schedule(self, request, spider): - return self.scheduler.enqueue_request(spider, request) + return self.slots[spider].scheduler.enqueue_request(request) def download(self, request, spider): slot = self.slots[spider] @@ -210,9 +206,10 @@ class ExecutionEngine(object): spider.name log.msg("Spider opened", spider=spider) nextcall = CallLaterOnce(self._next_request, spider) - slot = Slot(start_requests or (), close_if_idle, nextcall) + scheduler = self.scheduler_cls.from_settings(self.settings) + slot = Slot(start_requests or (), close_if_idle, nextcall, scheduler) self.slots[spider] = slot - yield self.scheduler.open_spider(spider) + yield scheduler.open(spider) yield self.scraper.open_spider(spider) stats.open_spider(spider) yield send_catch_log_deferred(signals.spider_opened, spider=spider) @@ -244,14 +241,12 @@ class ExecutionEngine(object): return slot.closing log.msg("Closing spider (%s)" % reason, spider=spider) - self.scheduler.clear_pending_requests(spider) - dfd = slot.close() dfd.addBoth(lambda _: self.scraper.close_spider(spider)) dfd.addErrback(log.err, spider=spider) - dfd.addBoth(lambda _: self.scheduler.close_spider(spider)) + dfd.addBoth(lambda _: slot.scheduler.close(reason)) dfd.addErrback(log.err, spider=spider) dfd.addBoth(lambda _: send_catch_log_deferred(signal=signals.spider_closed, \ diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 72c026fdc..1e2741e8c 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -1,45 +1,33 @@ from scrapy.utils.datatypes import PriorityQueue, PriorityStack from scrapy.utils.misc import load_object -from scrapy.conf import settings class Scheduler(object): - def __init__(self): - self.pending_requests = {} - self.dfo = settings['SCHEDULER_ORDER'].upper() == 'DFO' - self.dupefilter = load_object(settings['DUPEFILTER_CLASS'])() + def __init__(self, dupefilter, dfo=False): + self.dupefilter = dupefilter + Queue = PriorityStack if dfo else PriorityQueue + self.pending_requests = Queue() - def spider_has_pending_requests(self, spider): - if spider in self.pending_requests: - return bool(self.pending_requests[spider]) + @classmethod + def from_settings(cls, settings): + dfo = settings['SCHEDULER_ORDER'].upper() == 'DFO' + dupefilter_cls = load_object(settings['DUPEFILTER_CLASS']) + dupefilter = dupefilter_cls.from_settings(settings) + return cls(dupefilter, dfo=dfo) - def open_spider(self, spider): - if spider in self.pending_requests: - raise RuntimeError('Scheduler spider already opened: %s' % spider) + def has_pending_requests(self): + return bool(self.pending_requests) - Priority = PriorityStack if self.dfo else PriorityQueue - self.pending_requests[spider] = Priority() - return self.dupefilter.open_spider(spider) + def enqueue_request(self, request): + if request.dont_filter or not self.dupefilter.request_seen(request): + self.pending_requests.push(request, -request.priority) - def close_spider(self, spider): - if spider not in self.pending_requests: - raise RuntimeError('Scheduler spider is not open: %s' % spider) - self.pending_requests.pop(spider, None) - return self.dupefilter.close_spider(spider) + def next_request(self): + if self.pending_requests: + return self.pending_requests.pop()[0] - def enqueue_request(self, spider, request): - if request.dont_filter or not self.dupefilter.request_seen(spider, request): - self.pending_requests[spider].push(request, -request.priority) - - def clear_pending_requests(self, spider): - # TODO: flush queue here or discard enqueued requests, depending on how - # the spider is being closed. + def open(self, spider): pass - def next_request(self, spider): - q = self.pending_requests[spider] - if q: - return q.pop()[0] - - def is_idle(self): - return not self.pending_requests + def close(self, reason): + pass diff --git a/scrapy/utils/engine.py b/scrapy/utils/engine.py index 9e1638d90..56d7147b9 100644 --- a/scrapy/utils/engine.py +++ b/scrapy/utils/engine.py @@ -11,7 +11,6 @@ def get_engine_status(engine=None): global_tests = [ "time()-engine.start_time", - "engine.is_idle()", "engine.has_capacity()", "engine.scheduler.is_idle()", "len(engine.scheduler.pending_requests)",