imported patch scheduler_single_spider.patch

This commit is contained in:
Pablo Hoffman 2011-07-31 03:32:25 -03:00
parent 4d1e01a4d4
commit 6d989e3fb0
4 changed files with 45 additions and 75 deletions

View File

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

View File

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

View File

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

View File

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