mirror of https://github.com/scrapy/scrapy.git
Added support for returning deferreds from (some) signal handlers. Closes #193
This commit is contained in:
parent
5755bdcddc
commit
3e3a66620b
10
docs/faq.rst
10
docs/faq.rst
|
|
@ -29,8 +29,8 @@ application framework).
|
|||
Does Scrapy work with Python 3.0?
|
||||
---------------------------------
|
||||
|
||||
No, and there are no plans to port Scrapy to Python 3.0 yet. At the moment
|
||||
Scrapy works with Python 2.5 or 2.6.
|
||||
No, and there are no plans to port Scrapy to Python 3.0 yet. At the moment,
|
||||
Scrapy works with Python 2.5, 2.6 and 2.7.
|
||||
|
||||
Did Scrapy "steal" X from Django?
|
||||
---------------------------------
|
||||
|
|
@ -167,3 +167,9 @@ Can I use JSON for large exports?
|
|||
It'll depend on how large your output is. See :ref:`this warning
|
||||
<json-with-large-data>` in :class:`~scrapy.contrib.exporter.JsonItemExporter`
|
||||
documentation.
|
||||
|
||||
Can I return (Twisted) deferreds from signal handlers?
|
||||
------------------------------------------------------
|
||||
|
||||
Some signals support returning deferreds form their handlers, others don't. See
|
||||
the :ref:`topics-signals-ref` to know which ones.
|
||||
|
|
|
|||
|
|
@ -4,19 +4,28 @@
|
|||
Signals
|
||||
=======
|
||||
|
||||
Scrapy uses signals extensively to notify when certain actions occur. You can
|
||||
catch some of those signals in your Scrapy project or extension to perform
|
||||
additional tasks or extend Scrapy to add functionality not provided out of the
|
||||
box.
|
||||
Scrapy uses signals extensively to notify when certain events occur. You can
|
||||
catch some of those signals in your Scrapy project (using an :ref:`extension
|
||||
<topics-extensions>`, for example) to perform additional tasks or extend Scrapy
|
||||
to add functionality not provided out of the box.
|
||||
|
||||
Even though signals provide several arguments, the handlers which catch them
|
||||
don't have to receive all of them.
|
||||
Even though signals provide several arguments, the handlers that catch them
|
||||
don't need to accept all of them - the signal dispatching mechanism will only
|
||||
deliver the arguments that the handler receives.
|
||||
|
||||
For more information about working when see the documentation of
|
||||
`pydispatcher`_ (library used to implement signals).
|
||||
Finally, for more detailed information about signals internals see the
|
||||
documentation of `pydispatcher`_ (the which the signal dispatching mechanism is
|
||||
based on).
|
||||
|
||||
.. _pydispatcher: http://pydispatcher.sourceforge.net/
|
||||
|
||||
Deferred signal handlers
|
||||
========================
|
||||
|
||||
Some signals support returning `Twisted deferreds`_ from their handlers, see
|
||||
the :ref:`topics-signals-ref` below to know which ones.
|
||||
|
||||
.. _Twisted deferreds: http://twistedmatrix.com/documents/current/core/howto/defer.html
|
||||
|
||||
.. _topics-signals-ref:
|
||||
|
||||
|
|
@ -26,8 +35,7 @@ Built-in signals reference
|
|||
.. module:: scrapy.signals
|
||||
:synopsis: Signals definitions
|
||||
|
||||
Here's a list of signals used in Scrapy and their meaning, in alphabetical
|
||||
order.
|
||||
Here's the list of Scrapy built-in signals and their meaning.
|
||||
|
||||
engine_started
|
||||
--------------
|
||||
|
|
@ -38,6 +46,8 @@ engine_started
|
|||
Sent when the Scrapy engine is started (for example, when a crawling
|
||||
process has started).
|
||||
|
||||
This signal supports returning deferreds from their handlers.
|
||||
|
||||
engine_stopped
|
||||
--------------
|
||||
|
||||
|
|
@ -47,6 +57,8 @@ engine_stopped
|
|||
Sent when the Scrapy engine is stopped (for example, when a crawling
|
||||
process has finished).
|
||||
|
||||
This signal supports returning deferreds from their handlers.
|
||||
|
||||
item_scraped
|
||||
------------
|
||||
|
||||
|
|
@ -56,6 +68,8 @@ item_scraped
|
|||
Sent when the engine receives a new scraped item from the spider, and right
|
||||
before the item is sent to the :ref:`topics-item-pipeline`.
|
||||
|
||||
This signal supports returning deferreds from their handlers.
|
||||
|
||||
:param item: is the item scraped
|
||||
:type item: :class:`~scrapy.item.Item` object
|
||||
|
||||
|
|
@ -74,6 +88,8 @@ item_passed
|
|||
Sent after an item has passed all the :ref:`topics-item-pipeline` stages without
|
||||
being dropped.
|
||||
|
||||
This signal supports returning deferreds from their handlers.
|
||||
|
||||
:param item: the item which passed the pipeline
|
||||
:type item: :class:`~scrapy.item.Item` object
|
||||
|
||||
|
|
@ -93,6 +109,8 @@ item_dropped
|
|||
Sent after an item has been dropped from the :ref:`topics-item-pipeline`
|
||||
when some stage raised a :exc:`~scrapy.exceptions.DropItem` exception.
|
||||
|
||||
This signal supports returning deferreds from their handlers.
|
||||
|
||||
:param item: the item dropped from the :ref:`topics-item-pipeline`
|
||||
:type item: :class:`~scrapy.item.Item` object
|
||||
|
||||
|
|
@ -113,6 +131,8 @@ spider_closed
|
|||
Sent after a spider has been closed. This can be used to release per-spider
|
||||
resources reserved on :signal:`spider_opened`.
|
||||
|
||||
This signal supports returning deferreds from their handlers.
|
||||
|
||||
:param spider: the spider which has been closed
|
||||
:type spider: :class:`~scrapy.spider.BaseSpider` object
|
||||
|
||||
|
|
@ -135,6 +155,8 @@ spider_opened
|
|||
reserve per-spider resources, but can be used for any task that needs to be
|
||||
performed when a spider is opened.
|
||||
|
||||
This signal supports returning deferreds from their handlers.
|
||||
|
||||
:param spider: the spider which has been opened
|
||||
:type spider: :class:`~scrapy.spider.BaseSpider` object
|
||||
|
||||
|
|
@ -157,6 +179,8 @@ spider_idle
|
|||
You can, for example, schedule some requests in your :signal:`spider_idle`
|
||||
handler to prevent the spider from being closed.
|
||||
|
||||
This signal does not support returning deferreds from their handlers.
|
||||
|
||||
:param spider: the spider which has gone idle
|
||||
:type spider: :class:`~scrapy.spider.BaseSpider` object
|
||||
|
||||
|
|
@ -168,6 +192,8 @@ request_received
|
|||
|
||||
Sent when the engine receives a :class:`~scrapy.http.Request` from a spider.
|
||||
|
||||
This signal does not support returning deferreds from their handlers.
|
||||
|
||||
:param request: the request received
|
||||
:type request: :class:`~scrapy.http.Request` object
|
||||
|
||||
|
|
@ -182,6 +208,8 @@ request_uploaded
|
|||
|
||||
Sent right after the download has sent a :class:`~scrapy.http.Request`.
|
||||
|
||||
This signal does not support returning deferreds from their handlers.
|
||||
|
||||
:param request: the request uploaded/sent
|
||||
:type request: :class:`~scrapy.http.Request` object
|
||||
|
||||
|
|
@ -194,15 +222,17 @@ response_received
|
|||
.. signal:: response_received
|
||||
.. function:: response_received(response, spider)
|
||||
|
||||
Sent when the engine receives a new :class:`~scrapy.http.Response` from the
|
||||
downloader.
|
||||
|
||||
This signal does not support returning deferreds from their handlers.
|
||||
|
||||
:param response: the response received
|
||||
:type response: :class:`~scrapy.http.Response` object
|
||||
|
||||
:param spider: the spider for which the response is intended
|
||||
:type spider: :class:`~scrapy.spider.BaseSpider` object
|
||||
|
||||
Sent when the engine receives a new :class:`~scrapy.http.Response` from the
|
||||
downloader.
|
||||
|
||||
response_downloaded
|
||||
-------------------
|
||||
|
||||
|
|
@ -211,6 +241,8 @@ response_downloaded
|
|||
|
||||
Sent by the downloader right after a ``HTTPResponse`` is downloaded.
|
||||
|
||||
This signal does not support returning deferreds from their handlers.
|
||||
|
||||
:param response: the response downloaded
|
||||
:type response: :class:`~scrapy.http.Response` object
|
||||
|
||||
|
|
|
|||
|
|
@ -19,7 +19,7 @@ from scrapy.exceptions import IgnoreRequest, DontCloseSpider
|
|||
from scrapy.http import Response, Request
|
||||
from scrapy.spider import spiders
|
||||
from scrapy.utils.misc import load_object
|
||||
from scrapy.utils.signal import send_catch_log
|
||||
from scrapy.utils.signal import send_catch_log, send_catch_log_deferred
|
||||
from scrapy.utils.defer import mustbe_deferred
|
||||
|
||||
class ExecutionEngine(object):
|
||||
|
|
@ -44,11 +44,12 @@ class ExecutionEngine(object):
|
|||
self.configured = True
|
||||
self._spider_closed_callback = spider_closed_callback
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def start(self):
|
||||
"""Start the execution engine"""
|
||||
assert not self.running, "Engine already running"
|
||||
self.start_time = time()
|
||||
send_catch_log(signal=signals.engine_started)
|
||||
yield send_catch_log_deferred(signal=signals.engine_started)
|
||||
self.running = True
|
||||
|
||||
def stop(self):
|
||||
|
|
@ -218,7 +219,7 @@ class ExecutionEngine(object):
|
|||
self.downloader.open_spider(spider)
|
||||
yield self.scraper.open_spider(spider)
|
||||
stats.open_spider(spider)
|
||||
send_catch_log(signals.spider_opened, spider=spider)
|
||||
yield send_catch_log_deferred(signals.spider_opened, spider=spider)
|
||||
self.next_request(spider)
|
||||
|
||||
def _spider_idle(self, spider):
|
||||
|
|
@ -268,11 +269,12 @@ class ExecutionEngine(object):
|
|||
def _finish_closing_spider(self, spider):
|
||||
"""This function is called after the spider has been closed"""
|
||||
reason = self.closing.pop(spider, 'finished')
|
||||
send_catch_log(signal=signals.spider_closed, spider=spider, reason=reason)
|
||||
call = self._next_request_calls.pop(spider, None)
|
||||
if call and call.active():
|
||||
call.cancel()
|
||||
dfd = defer.maybeDeferred(stats.close_spider, spider, reason=reason)
|
||||
dfd = send_catch_log_deferred(signal=signals.spider_closed, \
|
||||
spider=spider, reason=reason)
|
||||
dfd.addBoth(lambda _: stats.close_spider(spider, reason=reason))
|
||||
dfd.addErrback(log.err, "Unhandled error in stats.close_spider()",
|
||||
spider=spider)
|
||||
dfd.addBoth(lambda _: spiders.close_spider(spider))
|
||||
|
|
@ -283,5 +285,6 @@ class ExecutionEngine(object):
|
|||
dfd.addBoth(lambda _: self._spider_closed_callback(spider))
|
||||
return dfd
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def _finish_stopping_engine(self):
|
||||
send_catch_log(signal=signals.engine_stopped)
|
||||
yield send_catch_log_deferred(signal=signals.engine_stopped)
|
||||
|
|
|
|||
|
|
@ -7,7 +7,7 @@ from twisted.internet import defer
|
|||
from scrapy.utils.defer import defer_result, defer_succeed, parallel, iter_errback
|
||||
from scrapy.utils.spider import iterate_spider_output
|
||||
from scrapy.utils.misc import load_object
|
||||
from scrapy.utils.signal import send_catch_log
|
||||
from scrapy.utils.signal import send_catch_log, send_catch_log_deferred
|
||||
from scrapy.exceptions import IgnoreRequest, DropItem
|
||||
from scrapy import signals
|
||||
from scrapy.http import Request, Response
|
||||
|
|
@ -171,14 +171,10 @@ class Scraper(object):
|
|||
elif isinstance(output, BaseItem):
|
||||
log.msg("Scraped %s in <%s>" % (output, request.url), level=log.DEBUG, \
|
||||
spider=spider)
|
||||
send_catch_log(signal=signals.item_scraped, \
|
||||
item=output, spider=spider, response=response)
|
||||
self.sites[spider].itemproc_size += 1
|
||||
# FIXME: this can't be called here because the stats spider may be
|
||||
# already closed
|
||||
#stats.max_value('scraper/max_itemproc_size', \
|
||||
# self.sites[spider].itemproc_size, spider=spider)
|
||||
dfd = self.itemproc.process_item(output, spider)
|
||||
dfd = send_catch_log_deferred(signal=signals.item_scraped, \
|
||||
item=output, spider=spider, response=response)
|
||||
dfd.addBoth(lambda _: self.itemproc.process_item(output, spider))
|
||||
dfd.addBoth(self._itemproc_finished, output, spider)
|
||||
return dfd
|
||||
elif output is None:
|
||||
|
|
@ -210,12 +206,12 @@ class Scraper(object):
|
|||
ex = output.value
|
||||
if isinstance(ex, DropItem):
|
||||
log.msg("Dropped %s - %s" % (item, str(ex)), level=log.WARNING, spider=spider)
|
||||
send_catch_log(signal=signals.item_dropped, \
|
||||
return send_catch_log_deferred(signal=signals.item_dropped, \
|
||||
item=item, spider=spider, exception=output.value)
|
||||
else:
|
||||
log.err(output, 'Error processing %s' % item, spider=spider)
|
||||
else:
|
||||
log.msg("Passed %s" % item, log.INFO, spider=spider)
|
||||
send_catch_log(signal=signals.item_passed, \
|
||||
return send_catch_log_deferred(signal=signals.item_passed, \
|
||||
item=item, spider=spider, output=output)
|
||||
|
||||
|
|
|
|||
|
|
@ -1,45 +1,84 @@
|
|||
from twisted.trial import unittest
|
||||
from twisted.python import log as txlog
|
||||
from twisted.python.failure import Failure
|
||||
from twisted.internet import defer, reactor
|
||||
|
||||
from scrapy.xlib.pydispatch import dispatcher
|
||||
from scrapy.utils.signal import send_catch_log
|
||||
from scrapy.utils.signal import send_catch_log, send_catch_log_deferred
|
||||
from scrapy import log
|
||||
|
||||
test_signal = object()
|
||||
|
||||
class SignalUtilsTest(unittest.TestCase):
|
||||
class SendCatchLogTest(unittest.TestCase):
|
||||
|
||||
@defer.inlineCallbacks
|
||||
def test_send_catch_log(self):
|
||||
test_signal = object()
|
||||
handlers_called = set()
|
||||
|
||||
def test_handler_error(arg):
|
||||
handlers_called.add(test_handler_error)
|
||||
a = 1/0
|
||||
|
||||
def test_handler_check(arg):
|
||||
handlers_called.add(test_handler_check)
|
||||
assert arg == 'test'
|
||||
return "OK"
|
||||
|
||||
def log_received(event):
|
||||
handlers_called.add(log_received)
|
||||
assert "test_handler_error" in event['message'][0]
|
||||
assert "error_handler" in event['message'][0]
|
||||
assert event['logLevel'] == log.ERROR
|
||||
|
||||
txlog.addObserver(log_received)
|
||||
dispatcher.connect(test_handler_error, signal=test_signal)
|
||||
dispatcher.connect(test_handler_check, signal=test_signal)
|
||||
result = send_catch_log(test_signal, arg='test')
|
||||
dispatcher.connect(self.error_handler, signal=test_signal)
|
||||
dispatcher.connect(self.ok_handler, signal=test_signal)
|
||||
result = yield defer.maybeDeferred(self._get_result, test_signal, arg='test', \
|
||||
handlers_called=handlers_called)
|
||||
|
||||
assert test_handler_error in handlers_called
|
||||
assert test_handler_check in handlers_called
|
||||
assert self.error_handler in handlers_called
|
||||
assert self.ok_handler in handlers_called
|
||||
assert log_received in handlers_called
|
||||
self.assertEqual(result[0][0], test_handler_error)
|
||||
self.assertEqual(result[0][0], self.error_handler)
|
||||
self.assert_(isinstance(result[0][1], Failure))
|
||||
self.assertEqual(result[1], (test_handler_check, "OK"))
|
||||
self.assertEqual(result[1], (self.ok_handler, "OK"))
|
||||
|
||||
txlog.removeObserver(log_received)
|
||||
self.flushLoggedErrors()
|
||||
dispatcher.disconnect(test_handler_error, signal=test_signal)
|
||||
dispatcher.disconnect(test_handler_check, signal=test_signal)
|
||||
dispatcher.disconnect(self.error_handler, signal=test_signal)
|
||||
dispatcher.disconnect(self.ok_handler, signal=test_signal)
|
||||
|
||||
def _get_result(self, signal, *a, **kw):
|
||||
return send_catch_log(signal, *a, **kw)
|
||||
|
||||
def error_handler(self, arg, handlers_called):
|
||||
handlers_called.add(self.error_handler)
|
||||
a = 1/0
|
||||
|
||||
def ok_handler(self, arg, handlers_called):
|
||||
handlers_called.add(self.ok_handler)
|
||||
assert arg == 'test'
|
||||
return "OK"
|
||||
|
||||
|
||||
class SendCatchLogDeferredTest(SendCatchLogTest):
|
||||
|
||||
def _get_result(self, signal, *a, **kw):
|
||||
return send_catch_log_deferred(signal, *a, **kw)
|
||||
|
||||
|
||||
class SendCatchLogDeferredTest2(SendCatchLogTest):
|
||||
|
||||
def ok_handler(self, arg, handlers_called):
|
||||
handlers_called.add(self.ok_handler)
|
||||
assert arg == 'test'
|
||||
d = defer.Deferred()
|
||||
reactor.callLater(0, d.callback, "OK")
|
||||
return d
|
||||
|
||||
def _get_result(self, signal, *a, **kw):
|
||||
return send_catch_log_deferred(signal, *a, **kw)
|
||||
|
||||
class SendCatchLogTest2(unittest.TestCase):
|
||||
|
||||
def test_error_logged_if_deferred_not_supported(self):
|
||||
test_signal = object()
|
||||
test_handler = lambda: defer.Deferred()
|
||||
log_events = []
|
||||
txlog.addObserver(log_events.append)
|
||||
dispatcher.connect(test_handler, test_signal)
|
||||
send_catch_log(test_signal)
|
||||
self.failUnless(log_events)
|
||||
self.failUnless("Cannot return deferreds from signal handler" in str(log_events))
|
||||
txlog.removeObserver(log_events.append)
|
||||
self.flushLoggedErrors()
|
||||
dispatcher.disconnect(test_handler, test_signal)
|
||||
|
|
|
|||
|
|
@ -1,5 +1,6 @@
|
|||
"""Helper functinos for working with signals"""
|
||||
|
||||
from twisted.internet.defer import maybeDeferred, DeferredList, Deferred
|
||||
from twisted.python.failure import Failure
|
||||
|
||||
from scrapy.xlib.pydispatch.dispatcher import Any, Anonymous, liveReceivers, \
|
||||
|
|
@ -19,6 +20,9 @@ def send_catch_log(signal=Any, sender=Anonymous, *arguments, **named):
|
|||
try:
|
||||
response = robustApply(receiver, signal=signal, sender=sender,
|
||||
*arguments, **named)
|
||||
if isinstance(response, Deferred):
|
||||
log.msg("Cannot return deferreds from signal handler: %s" % \
|
||||
receiver, log.ERROR, spider=spider)
|
||||
except dont_log:
|
||||
result = Failure()
|
||||
except Exception:
|
||||
|
|
@ -30,6 +34,30 @@ def send_catch_log(signal=Any, sender=Anonymous, *arguments, **named):
|
|||
responses.append((receiver, result))
|
||||
return responses
|
||||
|
||||
def send_catch_log_deferred(signal=Any, sender=Anonymous, *arguments, **named):
|
||||
"""Like send_catch_log but supports returning deferreds on signal handlers.
|
||||
Returns a deferred that gets fired once all signal handlers deferreds were
|
||||
fired.
|
||||
"""
|
||||
def logerror(failure, recv):
|
||||
if dont_log is None or not isinstance(failure.value, dont_log):
|
||||
log.err(failure, "Error caught on signal handler: %s" % recv, \
|
||||
spider=spider)
|
||||
return failure
|
||||
|
||||
dont_log = named.pop('dont_log', None)
|
||||
spider = named.get('spider', None)
|
||||
dfds = []
|
||||
for receiver in liveReceivers(getAllReceivers(sender, signal)):
|
||||
d = maybeDeferred(robustApply, receiver, signal=signal, sender=sender,
|
||||
*arguments, **named)
|
||||
d.addErrback(logerror, receiver)
|
||||
d.addBoth(lambda result: (receiver, result))
|
||||
dfds.append(d)
|
||||
d = DeferredList(dfds)
|
||||
d.addCallback(lambda out: [x[1] for x in out])
|
||||
return d
|
||||
|
||||
def disconnect_all(signal=Any, sender=Any):
|
||||
"""Disconnect all signal handlers. Useful for cleaning up after running
|
||||
tests
|
||||
|
|
|
|||
Loading…
Reference in New Issue