diff --git a/scrapy/core/scheduler.py b/scrapy/core/scheduler.py index 6d74085ea..8beb4e923 100644 --- a/scrapy/core/scheduler.py +++ b/scrapy/core/scheduler.py @@ -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 diff --git a/scrapy/settings/default_settings.py b/scrapy/settings/default_settings.py index 8dfca7215..c639167a5 100644 --- a/scrapy/settings/default_settings.py +++ b/scrapy/settings/default_settings.py @@ -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 diff --git a/scrapy/squeue.py b/scrapy/squeue.py index c01919fc1..0f90048af 100644 --- a/scrapy/squeue.py +++ b/scrapy/squeue.py @@ -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 diff --git a/scrapy/tests/test_squeue.py b/scrapy/tests/test_squeue.py index 03fd7da87..a6a12f4a1 100644 --- a/scrapy/tests/test_squeue.py +++ b/scrapy/tests/test_squeue.py @@ -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 diff --git a/scrapy/tests/test_utils_pqueue.py b/scrapy/tests/test_utils_pqueue.py index 22f3046a8..98b9d47e3 100644 --- a/scrapy/tests/test_utils_pqueue.py +++ b/scrapy/tests/test_utils_pqueue.py @@ -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()) diff --git a/scrapy/tests/test_utils_queue.py b/scrapy/tests/test_utils_queue.py index 1a18a9860..9a18abc72 100644 --- a/scrapy/tests/test_utils_queue.py +++ b/scrapy/tests/test_utils_queue.py @@ -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) diff --git a/scrapy/utils/queue.py b/scrapy/utils/queue.py index 09fce3c6b..1192b38ca 100644 --- a/scrapy/utils/queue.py +++ b/scrapy/utils/queue.py @@ -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