Some changes to persistent scheduler after some initial usage feedback:

* added LIFO queues, in addition to the original FIFO queues
* use LIFO queues (instead of FIFO queues) by default, since they resemble DFO
  better which is a more convenient crawling order for most cases
* do not adjust the priority based on depth by default (DEPTH_PRIORITY = 0)

If someone does need to use strict BFO order, it can be by done by setting:

    DEPTH_PRIORITY = 1
    SCHEDULER_DISK_QUEUE = 'scrapy.squeue.PickleFifoDiskQueue'
    SCHEDULER_MEMORY_QUEUE = 'scrapy.squeue.FifoMemoryQueue'
This commit is contained in:
Pablo Hoffman 2011-09-23 13:03:07 -03:00
parent cfddc314ce
commit f850a44784
7 changed files with 361 additions and 77 deletions

View File

@ -1,8 +1,8 @@
from __future__ import with_statement
import os
from os.path import join, exists
from scrapy.utils.queue import MemoryQueue
from scrapy.utils.pqueue import PriorityQueue
from scrapy.utils.reqser import request_to_dict, request_from_dict
from scrapy.utils.misc import load_object
@ -13,10 +13,11 @@ from scrapy import log
class Scheduler(object):
def __init__(self, dupefilter, jobdir=None, dqclass=None, logunser=False):
def __init__(self, dupefilter, jobdir=None, dqclass=None, mqclass=None, logunser=False):
self.df = dupefilter
self.dqdir = join(jobdir, 'requests.queue') if jobdir else None
self.dqdir = self._dqdir(jobdir)
self.dqclass = dqclass
self.mqclass = mqclass
self.logunser = logunser
@classmethod
@ -24,8 +25,9 @@ class Scheduler(object):
dupefilter_cls = load_object(settings['DUPEFILTER_CLASS'])
dupefilter = dupefilter_cls.from_settings(settings)
dqclass = load_object(settings['SCHEDULER_DISK_QUEUE'])
mqclass = load_object(settings['SCHEDULER_MEMORY_QUEUE'])
logunser = settings.getbool('LOG_UNSERIALIZABLE_REQUESTS')
return cls(dupefilter, job_dir(settings), dqclass, logunser)
return cls(dupefilter, job_dir(settings), dqclass, mqclass, logunser)
def has_pending_requests(self):
return len(self) > 0
@ -81,7 +83,7 @@ class Scheduler(object):
return request_from_dict(d, self.spider)
def _newmq(self, priority):
return MemoryQueue()
return self.mqclass()
def _newdq(self, priority):
return self.dqclass(join(self.dqdir, 'p%s' % priority))
@ -98,3 +100,10 @@ class Scheduler(object):
log.msg("Resuming crawl (%d requests scheduled)" % len(q), \
spider=self.spider)
return q
def _dqdir(self, jobdir):
if jobdir:
dqdir = join(jobdir, 'requests.queue')
if not exists(dqdir):
os.makedirs(dqdir)
return dqdir

View File

@ -45,7 +45,7 @@ DEFAULT_RESPONSE_ENCODING = 'ascii'
DEPTH_LIMIT = 0
DEPTH_STATS = True
DEPTH_PRIORITY = 1
DEPTH_PRIORITY = 0
DNSCACHE_ENABLED = True
@ -224,7 +224,8 @@ RETRY_PRIORITY_ADJUST = -1
ROBOTSTXT_OBEY = False
SCHEDULER = 'scrapy.core.scheduler.Scheduler'
SCHEDULER_DISK_QUEUE = 'scrapy.squeue.PickleDiskQueue'
SCHEDULER_DISK_QUEUE = 'scrapy.squeue.PickleLifoDiskQueue'
SCHEDULER_MEMORY_QUEUE = 'scrapy.squeue.LifoMemoryQueue'
SELECTORS_BACKEND = None # possible values: libxml2, lxml

View File

@ -1,33 +1,39 @@
"""
Scheduler disk-based queues
Scheduler queues
"""
import marshal, cPickle as pickle
from scrapy.utils.queue import DiskQueue
from scrapy.utils import queue
def _serializable_queue(queue_class, serialize, deserialize):
class PickleDiskQueue(DiskQueue):
class SerializableQueue(queue_class):
def push(self, obj):
try:
s = pickle.dumps(obj, protocol=2)
except pickle.PicklingError, e:
raise ValueError(str(e))
super(PickleDiskQueue, self).push(s)
def push(self, obj):
s = serialize(obj)
super(SerializableQueue, self).push(s)
def pop(self):
s = super(PickleDiskQueue, self).pop()
if s:
return pickle.loads(s)
def pop(self):
s = super(SerializableQueue, self).pop()
if s:
return deserialize(s)
return SerializableQueue
class MarshalDiskQueue(DiskQueue):
def _pickle_serialize(obj):
try:
return pickle.dumps(obj, protocol=2)
except pickle.PicklingError, e:
raise ValueError(str(e))
def push(self, obj):
super(MarshalDiskQueue, self).push(marshal.dumps(obj))
def pop(self):
s = super(MarshalDiskQueue, self).pop()
if s:
return marshal.loads(s)
PickleFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, \
_pickle_serialize, pickle.loads)
PickleLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, \
_pickle_serialize, pickle.loads)
MarshalFifoDiskQueue = _serializable_queue(queue.FifoDiskQueue, \
marshal.dumps, marshal.loads)
MarshalLifoDiskQueue = _serializable_queue(queue.LifoDiskQueue, \
marshal.dumps, marshal.loads)
FifoMemoryQueue = queue.FifoMemoryQueue
LifoMemoryQueue = queue.LifoMemoryQueue

View File

@ -1,5 +1,5 @@
from scrapy.tests import test_utils_queue as t
from scrapy.squeue import MarshalDiskQueue, PickleDiskQueue
from scrapy.squeue import MarshalFifoDiskQueue, MarshalLifoDiskQueue, PickleFifoDiskQueue, PickleLifoDiskQueue
from scrapy.item import Item, Field
from scrapy.http import Request
from scrapy.contrib.loader import ItemLoader
@ -14,12 +14,12 @@ class TestLoader(ItemLoader):
default_item_class = TestItem
name_out = staticmethod(test_processor)
class MarshalDiskQueueTest(t.DiskQueueTest):
class MarshalFifoDiskQueueTest(t.FifoDiskQueueTest):
chunksize = 100000
def queue(self):
return MarshalDiskQueue(self.qdir, chunksize=self.chunksize)
return MarshalFifoDiskQueue(self.qdir, chunksize=self.chunksize)
def test_serialize(self):
q = self.queue()
@ -34,25 +34,25 @@ class MarshalDiskQueueTest(t.DiskQueueTest):
q = self.queue()
self.assertRaises(ValueError, q.push, lambda x: x)
class ChunkSize1MarshalDiskQueueTest(MarshalDiskQueueTest):
class ChunkSize1MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest):
chunksize = 1
class ChunkSize2MarshalDiskQueueTest(MarshalDiskQueueTest):
class ChunkSize2MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest):
chunksize = 2
class ChunkSize3MarshalDiskQueueTest(MarshalDiskQueueTest):
class ChunkSize3MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest):
chunksize = 3
class ChunkSize4MarshalDiskQueueTest(MarshalDiskQueueTest):
class ChunkSize4MarshalFifoDiskQueueTest(MarshalFifoDiskQueueTest):
chunksize = 4
class PickleDiskQueueTest(MarshalDiskQueueTest):
class PickleFifoDiskQueueTest(MarshalFifoDiskQueueTest):
chunksize = 100000
def queue(self):
return PickleDiskQueue(self.qdir, chunksize=self.chunksize)
return PickleFifoDiskQueue(self.qdir, chunksize=self.chunksize)
def test_serialize_item(self):
q = self.queue()
@ -81,15 +81,66 @@ class PickleDiskQueueTest(MarshalDiskQueueTest):
self.assertEqual(r.url, r2.url)
assert r2.meta['request'] is r2
class ChunkSize1PickleDiskQueueTest(PickleDiskQueueTest):
class ChunkSize1PickleFifoDiskQueueTest(PickleFifoDiskQueueTest):
chunksize = 1
class ChunkSize2PickleDiskQueueTest(PickleDiskQueueTest):
class ChunkSize2PickleFifoDiskQueueTest(PickleFifoDiskQueueTest):
chunksize = 2
class ChunkSize3PickleDiskQueueTest(PickleDiskQueueTest):
class ChunkSize3PickleFifoDiskQueueTest(PickleFifoDiskQueueTest):
chunksize = 3
class ChunkSize4PickleDiskQueueTest(PickleDiskQueueTest):
class ChunkSize4PickleFifoDiskQueueTest(PickleFifoDiskQueueTest):
chunksize = 4
class MarshalLifoDiskQueueTest(t.LifoDiskQueueTest):
def queue(self):
return MarshalLifoDiskQueue(self.path)
def test_serialize(self):
q = self.queue()
q.push('a')
q.push(123)
q.push({'a': 'dict'})
self.assertEqual(q.pop(), {'a': 'dict'})
self.assertEqual(q.pop(), 123)
self.assertEqual(q.pop(), 'a')
def test_nonserializable_object(self):
q = self.queue()
self.assertRaises(ValueError, q.push, lambda x: x)
class PickleLifoDiskQueueTest(MarshalLifoDiskQueueTest):
def queue(self):
return PickleLifoDiskQueue(self.path)
def test_serialize_item(self):
q = self.queue()
i = TestItem(name='foo')
q.push(i)
i2 = q.pop()
assert isinstance(i2, TestItem)
self.assertEqual(i, i2)
def test_serialize_loader(self):
q = self.queue()
l = TestLoader()
q.push(l)
l2 = q.pop()
assert isinstance(l2, TestLoader)
assert l2.default_item_class is TestItem
self.assertEqual(l2.name_out('x'), 'xx')
def test_serialize_request_recursive(self):
q = self.queue()
r = Request('http://www.example.com')
r.meta['request'] = r
q.push(r)
r2 = q.pop()
assert isinstance(r2, Request)
self.assertEqual(r.url, r2.url)
assert r2.meta['request'] is r2

View File

@ -1,36 +1,32 @@
from twisted.trial import unittest
from scrapy.utils.pqueue import PriorityQueue
from scrapy.utils.queue import MemoryQueue, DiskQueue
from scrapy.utils.queue import FifoMemoryQueue, LifoMemoryQueue, FifoDiskQueue, LifoDiskQueue
class TestMemoryQueue(MemoryQueue):
def track_closed(cls):
"""Wraps a queue class to track down if close() method was called"""
def __init__(self, *a, **kw):
super(TestMemoryQueue, self).__init__(*a, **kw)
self.closed = False
class TrackingClosed(cls):
def close(self):
super(TestMemoryQueue, self).close()
self.closed = True
def __init__(self, *a, **kw):
super(TrackingClosed, self).__init__(*a, **kw)
self.closed = False
def close(self):
super(TrackingClosed, self).close()
self.closed = True
return TrackingClosed
class TestDiskQueue(DiskQueue):
def __init__(self, *a, **kw):
super(TestDiskQueue, self).__init__(*a, **kw)
self.closed = False
def close(self):
super(TestDiskQueue, self).close()
self.closed = True
class MemoryPriorityQueueTest(unittest.TestCase):
class FifoMemoryPriorityQueueTest(unittest.TestCase):
def setUp(self):
qfactory = lambda x: TestMemoryQueue()
self.q = PriorityQueue(qfactory)
self.q = PriorityQueue(self.qfactory)
def qfactory(self, prio):
return track_closed(FifoMemoryQueue)()
def test_push_pop_noprio(self):
self.q.push('a')
@ -94,11 +90,39 @@ class MemoryPriorityQueueTest(unittest.TestCase):
assert p1queue.closed
class DiskPriorityQueueTest(MemoryPriorityQueueTest):
class LifoMemoryPriorityQueueTest(FifoMemoryPriorityQueueTest):
def qfactory(self, prio):
return track_closed(LifoMemoryQueue)()
def test_push_pop_noprio(self):
self.q.push('a')
self.q.push('b')
self.q.push('c')
self.assertEqual(self.q.pop(), 'c')
self.assertEqual(self.q.pop(), 'b')
self.assertEqual(self.q.pop(), 'a')
self.assertEqual(self.q.pop(), None)
def test_push_pop_prio(self):
self.q.push('a', 3)
self.q.push('b', 1)
self.q.push('c', 2)
self.q.push('d', 1)
self.assertEqual(self.q.pop(), 'd')
self.assertEqual(self.q.pop(), 'b')
self.assertEqual(self.q.pop(), 'c')
self.assertEqual(self.q.pop(), 'a')
self.assertEqual(self.q.pop(), None)
class FifoDiskPriorityQueueTest(FifoMemoryPriorityQueueTest):
def setUp(self):
qfactory = lambda x: TestDiskQueue(self.mktemp())
self.q = PriorityQueue(qfactory)
self.q = PriorityQueue(self.qfactory)
def qfactory(self, prio):
return track_closed(FifoDiskQueue)(self.mktemp())
def test_nonserializable_object_one(self):
self.assertRaises(TypeError, self.q.push, lambda x: x, 0)
@ -123,3 +147,37 @@ class DiskPriorityQueueTest(MemoryPriorityQueueTest):
self.assertEqual(self.q.pop(), None)
self.assertEqual(self.q.close(), [])
class FifoDiskPriorityQueueTest(FifoMemoryPriorityQueueTest):
def qfactory(self, prio):
return track_closed(FifoDiskQueue)(self.mktemp())
def test_nonserializable_object_one(self):
self.assertRaises(TypeError, self.q.push, lambda x: x, 0)
self.assertEqual(self.q.close(), [])
def test_nonserializable_object_many_close(self):
self.q.push('a', 3)
self.q.push('b', 1)
self.assertRaises(TypeError, self.q.push, lambda x: x, 0)
self.q.push('c', 2)
self.assertEqual(self.q.pop(), 'b')
self.assertEqual(sorted(self.q.close()), [2, 3])
def test_nonserializable_object_many_pop(self):
self.q.push('a', 3)
self.q.push('b', 1)
self.assertRaises(TypeError, self.q.push, lambda x: x, 0)
self.q.push('c', 2)
self.assertEqual(self.q.pop(), 'b')
self.assertEqual(self.q.pop(), 'c')
self.assertEqual(self.q.pop(), 'a')
self.assertEqual(self.q.pop(), None)
self.assertEqual(self.q.close(), [])
class LifoDiskPriorityQueueTest(LifoMemoryPriorityQueueTest):
def qfactory(self, prio):
return track_closed(LifoDiskQueue)(self.mktemp())

View File

@ -1,12 +1,12 @@
import os, glob
from twisted.trial import unittest
from scrapy.utils.queue import MemoryQueue, DiskQueue
from scrapy.utils.queue import FifoMemoryQueue, LifoMemoryQueue, FifoDiskQueue, LifoDiskQueue
class MemoryQueueTest(unittest.TestCase):
class FifoMemoryQueueTest(unittest.TestCase):
def queue(self):
return MemoryQueue()
return FifoMemoryQueue()
def test_empty(self):
"""Empty queue test"""
@ -52,7 +52,56 @@ class MemoryQueueTest(unittest.TestCase):
self.assertEqual(len(q), 0)
class DiskQueueTest(MemoryQueueTest):
class LifoMemoryQueueTest(unittest.TestCase):
def queue(self):
return LifoMemoryQueue()
def test_empty(self):
"""Empty queue test"""
q = self.queue()
assert q.pop() is None
def test_push_pop1(self):
"""Basic push/pop test"""
q = self.queue()
q.push('a')
q.push('b')
q.push('c')
self.assertEqual(q.pop(), 'c')
self.assertEqual(q.pop(), 'b')
self.assertEqual(q.pop(), 'a')
self.assertEqual(q.pop(), None)
def test_push_pop2(self):
"""Test interleaved push and pops"""
q = self.queue()
q.push('a')
q.push('b')
q.push('c')
q.push('d')
self.assertEqual(q.pop(), 'd')
self.assertEqual(q.pop(), 'c')
q.push('e')
self.assertEqual(q.pop(), 'e')
self.assertEqual(q.pop(), 'b')
self.assertEqual(q.pop(), 'a')
def test_len(self):
q = self.queue()
self.assertEqual(len(q), 0)
q.push('a')
self.assertEqual(len(q), 1)
q.push('b')
q.push('c')
self.assertEqual(len(q), 3)
q.pop()
q.pop()
q.pop()
self.assertEqual(len(q), 0)
class FifoDiskQueueTest(FifoMemoryQueueTest):
chunksize = 100000
@ -60,7 +109,7 @@ class DiskQueueTest(MemoryQueueTest):
self.qdir = self.mktemp()
def queue(self):
return DiskQueue(self.qdir, chunksize=self.chunksize)
return FifoDiskQueue(self.qdir, chunksize=self.chunksize)
def test_close_open(self):
"""Test closing and re-opening keeps state"""
@ -108,14 +157,68 @@ class DiskQueueTest(MemoryQueueTest):
assert not os.path.exists(self.qdir)
class ChunkSize1DiskQueueTest(DiskQueueTest):
class ChunkSize1FifoDiskQueueTest(FifoDiskQueueTest):
chunksize = 1
class ChunkSize2DiskQueueTest(DiskQueueTest):
class ChunkSize2FifoDiskQueueTest(FifoDiskQueueTest):
chunksize = 2
class ChunkSize3DiskQueueTest(DiskQueueTest):
class ChunkSize3FifoDiskQueueTest(FifoDiskQueueTest):
chunksize = 3
class ChunkSize4DiskQueueTest(DiskQueueTest):
class ChunkSize4FifoDiskQueueTest(FifoDiskQueueTest):
chunksize = 4
class LifoDiskQueueTest(LifoMemoryQueueTest):
def setUp(self):
self.path = self.mktemp()
def queue(self):
return LifoDiskQueue(self.path)
def test_close_open(self):
"""Test closing and re-opening keeps state"""
q = self.queue()
q.push('a')
q.push('b')
q.push('c')
q.push('d')
self.assertEqual(q.pop(), 'd')
self.assertEqual(q.pop(), 'c')
q.close()
del q
q = self.queue()
self.assertEqual(len(q), 2)
q.push('e')
self.assertEqual(q.pop(), 'e')
self.assertEqual(q.pop(), 'b')
q.close()
del q
q = self.queue()
self.assertEqual(q.pop(), 'a')
self.assertEqual(len(q), 0)
def test_cleanup(self):
"""Test queue file is removed if queue is empty"""
q = self.queue()
assert os.path.exists(self.path)
for x in range(5):
q.push(str(x))
for x in range(5):
q.pop()
q.close()
assert not os.path.exists(self.path)
def test_file_size_shrinks(self):
"""Test size of queue file shrinks when popping items"""
q = self.queue()
q.push('a')
q.push('b')
q.close()
size = os.path.getsize(self.path)
q = self.queue()
q.pop()
q.close()
self.assertLess(os.path.getsize(self.path), size)

View File

@ -8,7 +8,7 @@ from collections import deque
from scrapy.utils.py26 import json
class MemoryQueue(object):
class FifoMemoryQueue(object):
"""Memory FIFO queue."""
def __init__(self):
@ -28,7 +28,14 @@ class MemoryQueue(object):
return len(self.q)
class DiskQueue(object):
class LifoMemoryQueue(FifoMemoryQueue):
"""Memory LIFO queue."""
def push(self, obj):
self.q.append(obj)
class FifoDiskQueue(object):
"""Persistent FIFO queue."""
szhdr_format = ">L"
@ -119,3 +126,52 @@ class DiskQueue(object):
os.remove(os.path.join(self.path, 'info.json'))
if not os.listdir(self.path):
os.rmdir(self.path)
class LifoDiskQueue(object):
"""Persistent LIFO queue."""
SIZE_FORMAT = ">L"
SIZE_SIZE = struct.calcsize(SIZE_FORMAT)
def __init__(self, path):
self.path = path
if os.path.exists(path):
self.f = open(path, 'rb+')
qsize = self.f.read(self.SIZE_SIZE)
self.size, = struct.unpack(self.SIZE_FORMAT, qsize)
self.f.seek(0, os.SEEK_END)
else:
self.f = open(path, 'wb+')
self.f.write(struct.pack(self.SIZE_FORMAT, 0))
self.size = 0
def push(self, string):
self.f.write(string)
ssize = struct.pack(self.SIZE_FORMAT, len(string))
self.f.write(ssize)
self.size += 1
def pop(self):
if not self.size:
return
self.f.seek(-self.SIZE_SIZE, os.SEEK_END)
size, = struct.unpack(self.SIZE_FORMAT, self.f.read())
self.f.seek(-size-self.SIZE_SIZE, os.SEEK_END)
data = self.f.read(size)
self.f.seek(-size, os.SEEK_CUR)
self.f.truncate()
self.size -= 1
return data
def close(self):
if self.size:
self.f.seek(0)
self.f.write(struct.pack(self.SIZE_FORMAT, self.size))
self.f.close()
if not self.size:
os.remove(self.path)
def __len__(self):
return self.size