From ed5d7561f9c4087f825ac7aa03e287182d87b35d Mon Sep 17 00:00:00 2001 From: Pablo Hoffman Date: Fri, 11 Jun 2010 17:22:14 -0300 Subject: [PATCH] Added SQS Execution Queue, and example script to add spiders to the queue --- bin/scrapy-sqs.py | 30 +++++++++++++++++++ scrapy/commands/start.py | 7 +++-- scrapy/conf/default_settings.py | 7 +++++ scrapy/contrib/queue/__init__.py | 0 scrapy/contrib/queue/sqs.py | 49 ++++++++++++++++++++++++++++++++ scrapy/core/queue.py | 2 +- 6 files changed, 91 insertions(+), 4 deletions(-) create mode 100755 bin/scrapy-sqs.py create mode 100644 scrapy/contrib/queue/__init__.py create mode 100644 scrapy/contrib/queue/sqs.py diff --git a/bin/scrapy-sqs.py b/bin/scrapy-sqs.py new file mode 100755 index 000000000..e710aebb8 --- /dev/null +++ b/bin/scrapy-sqs.py @@ -0,0 +1,30 @@ +#!/usr/bin/env python + +import sys + +from boto.sqs.connection import SQSConnection +from boto.sqs.message import Message + +from scrapy.utils.py26 import json +from scrapy.conf import settings + +qname = settings['SQS_QUEUE'] + +if len(sys.argv) <= 1: + print "usage: %s [args]" % sys.argv[0] + print + print "available commands:" + print " put - append spider to queue" + print + print "SQS queue: %s" % qname + print + sys.exit() + +cmd, args = sys.argv[1], sys.argv[2:] + +if cmd == 'put': + conn = SQSConnection(settings['AWS_ACCESS_KEY_ID'], \ + settings['AWS_SECRET_ACCESS_KEY']) + q = conn.create_queue(qname) + msg = Message(body=json.dumps({'name': args[0]})) + q.write(msg) diff --git a/scrapy/commands/start.py b/scrapy/commands/start.py index 3c32abcc6..efe21bd64 100644 --- a/scrapy/commands/start.py +++ b/scrapy/commands/start.py @@ -1,6 +1,7 @@ -from scrapy.core.queue import KeepAliveExecutionQueue from scrapy.command import ScrapyCommand from scrapy.core.manager import scrapymanager +from scrapy.utils.misc import load_object +from scrapy.conf import settings class Command(ScrapyCommand): @@ -10,6 +11,6 @@ class Command(ScrapyCommand): return "Start the Scrapy manager but don't run any spider (idle mode)" def run(self, args, opts): - q = KeepAliveExecutionQueue() - scrapymanager.queue = q + queue_class = load_object(settings['SERVICE_QUEUE']) + scrapymanager.queue = queue_class() scrapymanager.start() diff --git a/scrapy/conf/default_settings.py b/scrapy/conf/default_settings.py index d2983f898..e53d8be90 100644 --- a/scrapy/conf/default_settings.py +++ b/scrapy/conf/default_settings.py @@ -185,6 +185,8 @@ SCHEDULER_MIDDLEWARES_BASE = { SCHEDULER_ORDER = 'DFO' +SERVICE_QUEUE = 'scrapy.core.queue.KeepAliveExecutionQueue' + SPIDER_MANAGER_CLASS = 'scrapy.contrib.spidermanager.TwistedPluginSpiderManager' SPIDER_MIDDLEWARES = {} @@ -205,6 +207,11 @@ SPIDER_MODULES = [] SPIDERPROFILER_ENABLED = False +SQS_QUEUE = 'scrapy' +SQS_VISIBILITY_TIMEOUT = 7200 +SQS_POLLING_DELAY = 30 +SQS_REGION = 'us-east-1' + STATS_CLASS = 'scrapy.stats.collector.MemoryStatsCollector' STATS_ENABLED = True STATS_DUMP = False diff --git a/scrapy/contrib/queue/__init__.py b/scrapy/contrib/queue/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/scrapy/contrib/queue/sqs.py b/scrapy/contrib/queue/sqs.py new file mode 100644 index 000000000..69ce2dc7c --- /dev/null +++ b/scrapy/contrib/queue/sqs.py @@ -0,0 +1,49 @@ +import threading + +from twisted.internet import threads +from boto.sqs.connection import SQSConnection +from boto.sqs import regions + +from scrapy.core.queue import ExecutionQueue +from scrapy.utils.py26 import json +from scrapy.conf import settings + +class SQSExecutionQueue(ExecutionQueue): + + polling_delay = settings.getint('SQS_POLLING_DELAY') + queue_name = settings['SQS_QUEUE'] + region_name = settings['SQS_REGION'] + visibility_timeout = settings.getint('SQS_VISIBILITY_TIMEOUT') + aws_access_key_id = settings['AWS_ACCESS_KEY_ID'] + aws_secret_access_key = settings['AWS_SECRET_ACCESS_KEY'] + + def __init__(self, *a, **kw): + super(SQSExecutionQueue, self).__init__(*a, **kw) + self.region = self._get_region() + self._tls = threading.local() + + def _append_next(self): + return threads.deferToThread(self._append_next_from_sqs) + + def _append_next_from_sqs(self): + q = self._get_sqs_queue() + msgs = q.get_messages(1, visibility_timeout=self.visibility_timeout) + if msgs: + msg = msgs[0] + msg.delete() + spargs = json.loads(msg.get_body()) + spname = spargs.pop('name') + self.append_spider_name(spname, **spargs) + + def _get_sqs_queue(self): + if not hasattr(self._tls, 'queue'): + c = SQSConnection(self.aws_access_key_id, self.aws_secret_access_key, \ + region=self.region) + self._tls.queue = c.create_queue(self.queue_name) + return self._tls.queue + + def _get_region(self, name=region_name): + return [r for r in regions() if r.name == name][0] + + def is_finished(self): + return False diff --git a/scrapy/core/queue.py b/scrapy/core/queue.py index f8dd340df..354f3a135 100644 --- a/scrapy/core/queue.py +++ b/scrapy/core/queue.py @@ -15,7 +15,7 @@ class ExecutionQueue(object): self._spiders = _spiders def _append_next(self): - """Called when there are no more itemsl left in self.spider_requests. + """Called when there are no more items left in self.spider_requests. This method is meant to be overriden in subclasses to add new (spider, requests) tuples to self.spider_requests. It can return a Deferred. """