mirror of https://github.com/scrapy/scrapy.git
core: Get rid of duplicate filtering as a scheduler builtin feature. closes #49.
Implements a DuplicatesFilterMiddleware as spidermiddleware, a wraper using a minimal defined API of a filtering class configurable by settings. Enabling this middleware doesn't gives us same functionality compared to scheduler duplicate filter builtin, but it filter the most important source for duplicate requests, the spiders. What requests aren't filtered by new middleware? The ones originated from any part of scrapy outside of spiders, like S3 images requests or any other request manually schedule using ``scrapyengine.schedule()`` method. Previously, we usually added dont_filter=True to requests created outside of spiders to avoid collisions downloading same pages than spider. Now, this is not required anymore because new middleware filters just the spider generated requests. There is a caveat, as usual downloadmiddlewares can returns a Request object at any point of the chain, and that request is scheduled and downloaded as usual too. One of the downloadmiddlewares using this feature is RedirectMiddleware that counts on scheduler filtering builtin to avoid redirection loops. I think we can implement a request time to live decreasing counter and add it to request's ``meta`` attribute with a default value if not present, and decrement each time the request is redirected. --HG-- extra : convert_revision : svn%3Ab85faa78-f9eb-468e-a121-7cced6da292c%40846
This commit is contained in:
parent
0aba276a64
commit
c9f2865c83
|
|
@ -338,6 +338,18 @@ Default: ``180``
|
|||
|
||||
The amount of time (in secs) that the downloader will wait before timing out.
|
||||
|
||||
.. setting:: DUPLICATESFILTER_FILTERCLASS
|
||||
|
||||
DUPLICATESFILTER_FILTERCLASS
|
||||
----------------------------
|
||||
|
||||
Default: ``scrapy.contrib.spidermiddleware.SimplePerDomainFilter``
|
||||
|
||||
The class used to detect and filter duplicated requests.
|
||||
|
||||
Default ``SimplePerDomainFilter`` filter based on request fingerprint and
|
||||
grouping per domain.
|
||||
|
||||
.. setting:: ENGINE_DEBUG
|
||||
|
||||
ENGINE_DEBUG
|
||||
|
|
|
|||
|
|
@ -76,6 +76,8 @@ DOWNLOADER_MIDDLEWARES = [
|
|||
|
||||
DOWNLOADER_STATS = True
|
||||
|
||||
DUPLICATESFILTER_FILTERCLASS = 'scrapy.contrib.spidermiddleware.duplicatesfilter.SimplePerDomainFilter'
|
||||
|
||||
ENABLED_SPIDERS_FILE = ''
|
||||
|
||||
ENGINE_DEBUG = False
|
||||
|
|
|
|||
|
|
@ -14,7 +14,20 @@ from scrapy import log
|
|||
|
||||
|
||||
class DuplicatesFilterMiddleware(object):
|
||||
"""Filter out duplicate requests to avoid visiting same page more than once"""
|
||||
"""Filter out duplicate requests to avoid visiting same page more than once.
|
||||
|
||||
filter class (defined by DUPLICATESFILTER_FILTERCLASS setting) must
|
||||
inplement a simple API:
|
||||
|
||||
* open(domain) called when a new domain starts
|
||||
|
||||
* close(domain) called when a domain is going to be closed.
|
||||
|
||||
* add(domain, request) called each time a new request needs to be tested
|
||||
looking if it is was already seen or is a new one. This is the most
|
||||
important method for a filtering class.
|
||||
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
clspath = settings.get('DUPLICATESFILTER_FILTERCLASS')
|
||||
|
|
@ -57,9 +70,3 @@ class SimplePerDomainFilter(dict):
|
|||
self[domain].add(fp)
|
||||
return True
|
||||
return False
|
||||
|
||||
def has(self, domain, request):
|
||||
"""Check if a request was already seen for a domain"""
|
||||
return request_fingerprint(request) in self[domain]
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -4,74 +4,72 @@ The Scrapy Scheduler
|
|||
|
||||
from twisted.internet import defer
|
||||
|
||||
from scrapy.core.scheduler.filter import GroupFilter
|
||||
from scrapy import log
|
||||
from scrapy.core.exceptions import IgnoreRequest
|
||||
from scrapy.utils.datatypes import PriorityQueue, PriorityStack
|
||||
from scrapy.utils.request import request_fingerprint
|
||||
from scrapy.utils.defer import defer_fail
|
||||
from scrapy.conf import settings
|
||||
|
||||
|
||||
class Scheduler(object) :
|
||||
"""
|
||||
The scheduler decides what to scrape next. In other words, it defines the
|
||||
crawling order. The scheduler schedules websites and requests to be
|
||||
scraped. Individual web pages that are to be scraped are batched up into a
|
||||
"run" for a website. As the domain is being scraped, pages that are
|
||||
discovered are added to the scheduler.
|
||||
"""The scheduler decides what to scrape next. In other words, it defines the
|
||||
crawling order.
|
||||
|
||||
The scheduler schedules websites and requests to be scraped. Individual
|
||||
web pages that are to be scraped are batched up into a "run" for a website.
|
||||
|
||||
As the domain is being scraped, pages that are discovered are added to the
|
||||
scheduler.
|
||||
|
||||
Typical usage:
|
||||
|
||||
* next_available_domain() called to find out when there is something to do
|
||||
* open_domain() called to commence scraping a website
|
||||
* enqueue_request() called multiple times when new links found
|
||||
* next_request() called multiple times when there is capacity to process urls
|
||||
* close_domain() called when there are no more pages or upon error
|
||||
* next_domain() called each time a domain slot is freed, and return
|
||||
next domain to be scraped.
|
||||
|
||||
Note a couple things:
|
||||
1) The order in which you get back the list of pages to scrape is not
|
||||
necesarily the order you put them in.
|
||||
2) A canonical URL is calculated for each url for each domain to check that
|
||||
it is unique, however the actual url passed in is returned when
|
||||
get_next_page is called.
|
||||
This is for two main reasons:
|
||||
* To be nice to the screen scraped site just incase the specific
|
||||
format of the url is significant.
|
||||
* To take advantage of any caching that uses the URL/URI as a key
|
||||
* open_domain() called to commence scraping a website
|
||||
|
||||
all_domains contains the names of all domains that are to be scheduled.
|
||||
* enqueue_request() called multiple times to enqueue new requests to be downloaded
|
||||
|
||||
* next_request() called multiple times when there is capacity to download requests
|
||||
|
||||
* close_domain() called when there are no more pages for a website
|
||||
|
||||
Notes:
|
||||
|
||||
1. The order in which you get back the list of pages to scrape is not
|
||||
necesarily the order you put them in.
|
||||
|
||||
``pending_domains_count`` contains the names of all domains that are to be scheduled.
|
||||
|
||||
Two crawling orders are available by default, which can be set with the
|
||||
SCHEDULER_ORDER settings:
|
||||
|
||||
* BFO - breath-first order (default). Consumes more memory than DFO but reaches
|
||||
most relevant pages faster.
|
||||
* DFO - depth-first order. Consumes less memory than BFO but usually takes
|
||||
longer to reach the most relevant pages.
|
||||
* BFO - breath-first order (default). Consumes more memory than DFO but reaches
|
||||
most relevant pages faster.
|
||||
|
||||
* DFO - depth-first order. Consumes less memory than BFO but usually takes
|
||||
longer to reach the most relevant pages.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self.pending_domains_count = {}
|
||||
self.domains_queue = PriorityQueue()
|
||||
self.pending_requests = {}
|
||||
self.groupfilter = GroupFilter()
|
||||
self.dfo = settings.get('SCHEDULER_ORDER', '').upper() == 'DFO'
|
||||
|
||||
def domain_is_open(self, domain):
|
||||
"""Check if scheduler's resources were allocated for a domain"""
|
||||
return domain in self.pending_requests
|
||||
|
||||
def is_pending(self, domain):
|
||||
"""Check if a domain is waiting to be scraped in domain's queue."""
|
||||
return domain in self.pending_domains_count
|
||||
|
||||
def domain_has_pending(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
|
||||
"""
|
||||
"""Return next domain available to scrape and remove it from available domains queue"""
|
||||
if self.pending_domains_count:
|
||||
domain, priority = self.domains_queue.pop()
|
||||
if self.pending_domains_count[domain] == 1:
|
||||
|
|
@ -82,11 +80,13 @@ class Scheduler(object) :
|
|||
return None
|
||||
|
||||
def add_domain(self, domain, priority=1):
|
||||
"""
|
||||
This functions schedules a new domain to be scraped, with the given
|
||||
priority. It doesn't check if the domain is already scheduled. A
|
||||
domain can be scheduled twice, either with the same or with different
|
||||
"""This functions schedules a new domain to be scraped, with the given priority.
|
||||
|
||||
It doesn't check if the domain is already scheduled.
|
||||
|
||||
A domain can be scheduled twice, either with the same or with different
|
||||
priority.
|
||||
|
||||
"""
|
||||
self.domains_queue.push(domain, priority)
|
||||
if domain not in self.pending_domains_count:
|
||||
|
|
@ -95,53 +95,37 @@ class Scheduler(object) :
|
|||
self.pending_domains_count[domain] += 1
|
||||
|
||||
def open_domain(self, domain):
|
||||
"""
|
||||
Allocates resources for maintaining a schedule for domain.
|
||||
"""
|
||||
if self.dfo:
|
||||
self.pending_requests[domain] = PriorityStack()
|
||||
else:
|
||||
self.pending_requests[domain] = PriorityQueue()
|
||||
|
||||
self.groupfilter.open(domain)
|
||||
"""Allocates resources for maintaining a schedule for domain."""
|
||||
Priority = PriorityStack if self.dfo else PriorityQueue
|
||||
self.pending_requests[domain] = Priority()
|
||||
|
||||
def enqueue_request(self, domain, request, priority=1):
|
||||
"""
|
||||
Add a page to be scraped for a domain that is currently being scraped.
|
||||
"""
|
||||
requestid = request_fingerprint(request)
|
||||
added = self.groupfilter.add(domain, requestid)
|
||||
|
||||
if request.dont_filter or added:
|
||||
deferred = defer.Deferred()
|
||||
self.pending_requests[domain].push((request, deferred), priority)
|
||||
return deferred
|
||||
else:
|
||||
return defer_fail(IgnoreRequest('Skipped (already visited): %s' % request))
|
||||
|
||||
def request_seen(self, domain, request):
|
||||
"""
|
||||
Returns True if the given Request was scheduled before for the given
|
||||
domain
|
||||
"""
|
||||
return self.groupfilter.has(domain, request_fingerprint(request))
|
||||
"""Enqueue a request to be downloaded for a domain that is currently being scraped."""
|
||||
dfd = defer.Deferred()
|
||||
self.pending_requests[domain].push((request, dfd), priority)
|
||||
return dfd
|
||||
|
||||
def next_request(self, domain):
|
||||
"""
|
||||
Get the next request to be scraped.
|
||||
"""Return the next available request to be downloaded for a domain.
|
||||
|
||||
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 should be returned if there are no more request pending for the domain passed.
|
||||
"""
|
||||
pending_list = self.pending_requests.get(domain)
|
||||
if pending_list:
|
||||
return pending_list.pop()[0]
|
||||
else:
|
||||
try:
|
||||
# The second value is the request scheduled priority, returns the first one.
|
||||
return self.pending_requests[domain].pop()[0]
|
||||
except (KeyError, IndexError), ex:
|
||||
return (None, None)
|
||||
|
||||
def close_domain(self, domain) :
|
||||
"""
|
||||
Called once we are finished scraping a domain. The scheduler will
|
||||
free any resources associated with the domain.
|
||||
"""Called once we are finished scraping a domain.
|
||||
|
||||
The scheduler will free any resources associated with the domain.
|
||||
|
||||
"""
|
||||
try :
|
||||
del self.pending_requests[domain]
|
||||
|
|
@ -149,12 +133,6 @@ class Scheduler(object) :
|
|||
msg = "Could not clear pending pages for domain %s, %s" % (domain, inst)
|
||||
log.msg(msg, level=log.WARNING)
|
||||
|
||||
try :
|
||||
self.groupfilter.close(domain)
|
||||
except Exception, inst:
|
||||
msg = "Could not clear url filter for domain %s, %s" % (domain, inst)
|
||||
log.msg(msg, level=log.WARNING)
|
||||
|
||||
def remove_pending_domain(self, domain):
|
||||
"""
|
||||
Remove a pending domain not yet started. If the domain was enqueued
|
||||
|
|
@ -162,12 +140,14 @@ class Scheduler(object) :
|
|||
|
||||
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 not self.domain_is_open(domain):
|
||||
return self.pending_domains_count.pop(domain, 0)
|
||||
|
||||
def is_idle(self):
|
||||
"""Checks if the schedulers has any request pendings"""
|
||||
return not self.pending_requests
|
||||
|
|
|
|||
|
|
@ -104,5 +104,3 @@ URLLENGTH_LIMIT = 2083
|
|||
WS_ENABLED = 0
|
||||
|
||||
SPIDERPROFILER_ENABLED = 0
|
||||
|
||||
#DUPLICATESFILTER_FILTERCLASS = 'scrapy.contrib.spidermiddleware.duplicatesfilter.SimplePerDomainFilter'
|
||||
|
|
|
|||
|
|
@ -0,0 +1,50 @@
|
|||
import unittest
|
||||
|
||||
from scrapy.spider import spiders
|
||||
from scrapy.http import Request, Response
|
||||
from scrapy.contrib.spidermiddleware.duplicatesfilter import DuplicatesFilterMiddleware, SimplePerDomainFilter
|
||||
|
||||
class DuplicatesFilterMiddlewareTest(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
spiders.spider_modules = ['scrapy.tests.test_spiders']
|
||||
spiders.reload()
|
||||
self.spider = spiders.fromdomain('scrapytest.org')
|
||||
|
||||
def test_process_spider_output(self):
|
||||
mw = DuplicatesFilterMiddleware()
|
||||
mw.filter.open('scrapytest.org')
|
||||
|
||||
response = Response('')
|
||||
r1 = Request('http://scrapytest.org/1')
|
||||
r2 = Request('http://scrapytest.org/2')
|
||||
r3 = Request('http://scrapytest.org/2')
|
||||
|
||||
filtered = list(mw.process_spider_output(response, [r1, r2, r3], self.spider))
|
||||
|
||||
assert r1 in filtered
|
||||
assert r2 in filtered
|
||||
assert r3 not in filtered
|
||||
|
||||
mw.filter.close('scrapytest.org')
|
||||
|
||||
|
||||
class SimplePerDomainFilterTest(unittest.TestCase):
|
||||
|
||||
def test_filter(self):
|
||||
domain = 'scrapytest.org'
|
||||
filter = SimplePerDomainFilter()
|
||||
filter.open(domain)
|
||||
assert domain in filter
|
||||
|
||||
r1 = Request('http://scrapytest.org/1')
|
||||
r2 = Request('http://scrapytest.org/2')
|
||||
r3 = Request('http://scrapytest.org/2')
|
||||
|
||||
assert filter.add(domain, r1)
|
||||
assert filter.add(domain, r2)
|
||||
assert not filter.add(domain, r3)
|
||||
|
||||
filter.close(domain)
|
||||
assert domain not in filter
|
||||
|
||||
Loading…
Reference in New Issue