Merge pull request #4205 from elacuesta/bytes_received_signal

Add bytes_received signal
This commit is contained in:
Mikhail Korobov 2020-05-11 15:09:55 +05:00 committed by GitHub
commit b183579564
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
5 changed files with 145 additions and 41 deletions

View File

@ -112,7 +112,7 @@ engine_started
Sent when the Scrapy engine has started crawling.
This signal supports returning deferreds from their handlers.
This signal supports returning deferreds from its handlers.
.. note:: This signal may be fired *after* the :signal:`spider_opened` signal,
depending on how the spider was started. So **don't** rely on this signal
@ -127,7 +127,7 @@ 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.
This signal supports returning deferreds from its handlers.
Item signals
------------
@ -149,7 +149,7 @@ item_scraped
Sent when an item has been scraped, after it has passed all the
:ref:`topics-item-pipeline` stages (without being dropped).
This signal supports returning deferreds from their handlers.
This signal supports returning deferreds from its handlers.
:param item: the item scraped
:type item: dict or :class:`~scrapy.item.Item` object
@ -169,7 +169,7 @@ 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.
This signal supports returning deferreds from its handlers.
:param item: the item dropped from the :ref:`topics-item-pipeline`
:type item: dict or :class:`~scrapy.item.Item` object
@ -194,7 +194,7 @@ item_error
Sent when a :ref:`topics-item-pipeline` generates an error (i.e. raises
an exception), except :exc:`~scrapy.exceptions.DropItem` exception.
This signal supports returning deferreds from their handlers.
This signal supports returning deferreds from its handlers.
:param item: the item dropped from the :ref:`topics-item-pipeline`
:type item: dict or :class:`~scrapy.item.Item` object
@ -220,7 +220,7 @@ 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.
This signal supports returning deferreds from its handlers.
:param spider: the spider which has been closed
:type spider: :class:`~scrapy.spiders.Spider` object
@ -244,7 +244,7 @@ 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.
This signal supports returning deferreds from its handlers.
:param spider: the spider which has been opened
:type spider: :class:`~scrapy.spiders.Spider` object
@ -268,7 +268,7 @@ spider_idle
You may raise a :exc:`~scrapy.exceptions.DontCloseSpider` exception to
prevent the spider from being closed.
This signal does not support returning deferreds from their handlers.
This signal does not support returning deferreds from its handlers.
:param spider: the spider which has gone idle
:type spider: :class:`~scrapy.spiders.Spider` object
@ -287,7 +287,7 @@ spider_error
Sent when a spider callback generates an error (i.e. raises an exception).
This signal does not support returning deferreds from their handlers.
This signal does not support returning deferreds from its handlers.
:param failure: the exception raised
:type failure: twisted.python.failure.Failure
@ -310,7 +310,7 @@ request_scheduled
Sent when the engine schedules a :class:`~scrapy.http.Request`, to be
downloaded later.
The signal does not support returning deferreds from their handlers.
This signal does not support returning deferreds from its handlers.
:param request: the request that reached the scheduler
:type request: :class:`~scrapy.http.Request` object
@ -327,7 +327,7 @@ request_dropped
Sent when a :class:`~scrapy.http.Request`, scheduled by the engine to be
downloaded later, is rejected by the scheduler.
The signal does not support returning deferreds from their handlers.
This signal does not support returning deferreds from its handlers.
:param request: the request that reached the scheduler
:type request: :class:`~scrapy.http.Request` object
@ -343,7 +343,7 @@ request_reached_downloader
Sent when a :class:`~scrapy.http.Request` reached downloader.
The signal does not support returning deferreds from their handlers.
This signal does not support returning deferreds from its handlers.
:param request: the request that reached downloader
:type request: :class:`~scrapy.http.Request` object
@ -370,6 +370,29 @@ request_left_downloader
:param spider: the spider that yielded the request
:type spider: :class:`~scrapy.spiders.Spider` object
bytes_received
~~~~~~~~~~~~~~
.. signal:: bytes_received
.. function:: bytes_received(data, request, spider)
Sent by the HTTP 1.1 and S3 download handlers when a group of bytes is
received for a specific request. This signal might be fired multiple
times for the same request, with partial data each time. For instance,
a possible scenario for a 25 kb response would be two signals fired
with 10 kb of data, and a final one with 5 kb of data.
This signal does not support returning deferreds from its handlers.
:param data: the data received by the download handler
:type spider: :class:`bytes` object
:param request: the request that generated the response
:type request: :class:`~scrapy.http.Request` object
:param spider: the spider associated with the response
:type spider: :class:`~scrapy.spiders.Spider` object
Response signals
----------------
@ -382,7 +405,7 @@ response_received
Sent when the engine receives a new :class:`~scrapy.http.Response` from the
downloader.
This signal does not support returning deferreds from their handlers.
This signal does not support returning deferreds from its handlers.
:param response: the response received
:type response: :class:`~scrapy.http.Response` object
@ -401,7 +424,7 @@ response_downloaded
Sent by the downloader right after a ``HTTPResponse`` is downloaded.
This signal does not support returning deferreds from their handlers.
This signal does not support returning deferreds from its handlers.
:param response: the response downloaded
:type response: :class:`~scrapy.http.Response` object

View File

@ -18,6 +18,7 @@ from twisted.web.http_headers import Headers as TxHeaders
from twisted.web.iweb import IBodyProducer, UNKNOWN_LENGTH
from zope.interface import implementer
from scrapy import signals
from scrapy.core.downloader.tls import openssl_methods
from scrapy.core.downloader.webclient import _parse
from scrapy.exceptions import ScrapyDeprecationWarning
@ -34,6 +35,8 @@ class HTTP11DownloadHandler:
lazy = False
def __init__(self, settings, crawler=None):
self._crawler = crawler
from twisted.internet import reactor
self._pool = HTTPConnectionPool(reactor, persistent=True)
self._pool.maxPersistentPerHost = settings.getint('CONCURRENT_REQUESTS_PER_DOMAIN')
@ -79,6 +82,7 @@ class HTTP11DownloadHandler:
maxsize=getattr(spider, 'download_maxsize', self._default_maxsize),
warnsize=getattr(spider, 'download_warnsize', self._default_warnsize),
fail_on_dataloss=self._fail_on_dataloss,
crawler=self._crawler,
)
return agent.download_request(request)
@ -276,7 +280,7 @@ class ScrapyAgent:
_TunnelingAgent = TunnelingAgent
def __init__(self, contextFactory=None, connectTimeout=10, bindAddress=None, pool=None,
maxsize=0, warnsize=0, fail_on_dataloss=True):
maxsize=0, warnsize=0, fail_on_dataloss=True, crawler=None):
self._contextFactory = contextFactory
self._connectTimeout = connectTimeout
self._bindAddress = bindAddress
@ -285,6 +289,7 @@ class ScrapyAgent:
self._warnsize = warnsize
self._fail_on_dataloss = fail_on_dataloss
self._txresponse = None
self._crawler = crawler
def _get_agent(self, request, timeout):
from twisted.internet import reactor
@ -407,7 +412,15 @@ class ScrapyAgent:
d = defer.Deferred(_cancel)
txresponse.deliverBody(
_ResponseReader(d, txresponse, request, maxsize, warnsize, fail_on_dataloss)
_ResponseReader(
finished=d,
txresponse=txresponse,
request=request,
maxsize=maxsize,
warnsize=warnsize,
fail_on_dataloss=fail_on_dataloss,
crawler=self._crawler,
)
)
# save response for timeouts
@ -449,7 +462,7 @@ class _RequestBodyProducer:
class _ResponseReader(protocol.Protocol):
def __init__(self, finished, txresponse, request, maxsize, warnsize, fail_on_dataloss):
def __init__(self, finished, txresponse, request, maxsize, warnsize, fail_on_dataloss, crawler):
self._finished = finished
self._txresponse = txresponse
self._request = request
@ -462,6 +475,7 @@ class _ResponseReader(protocol.Protocol):
self._bytes_received = 0
self._certificate = None
self._ip_address = None
self._crawler = crawler
def connectionMade(self):
if self._certificate is None:
@ -479,6 +493,13 @@ class _ResponseReader(protocol.Protocol):
self._bodybuf.write(bodyBytes)
self._bytes_received += len(bodyBytes)
self._crawler.signals.send_catch_log(
signal=signals.bytes_received,
data=bodyBytes,
request=self._request,
spider=self._crawler.spider,
)
if self._maxsize and self._bytes_received > self._maxsize:
logger.error("Received (%(bytes)s) bytes larger than download "
"max size (%(maxsize)s) in request %(request)s.",

View File

@ -17,6 +17,7 @@ request_reached_downloader = object()
request_left_downloader = object()
response_received = object()
response_downloaded = object()
bytes_received = object()
item_scraped = object()
item_dropped = object()
item_error = object()

View File

@ -730,6 +730,9 @@ class Http11ProxyTestCase(HttpProxyTestCase):
class HttpDownloadHandlerMock:
def __init__(self, *args, **kwargs):
pass
def download_request(self, request, spider):
return request

View File

@ -13,22 +13,24 @@ module with the ``runserver`` argument::
import os
import re
import sys
from collections import defaultdict
from urllib.parse import urlparse
from twisted.internet import reactor, defer
from twisted.web import server, static, util
from twisted.trial import unittest
from twisted.web import server, static, util
from pydispatch import dispatcher
from scrapy import signals
from scrapy.core.engine import ExecutionEngine
from scrapy.utils.test import get_crawler
from pydispatch import dispatcher
from tests import tests_datadir
from scrapy.spiders import Spider
from scrapy.http import Request
from scrapy.item import Item, Field
from scrapy.linkextractors import LinkExtractor
from scrapy.http import Request
from scrapy.spiders import Spider
from scrapy.utils.signal import disconnect_all
from scrapy.utils.test import get_crawler
from tests import tests_datadir, get_testdata
class TestItem(Item):
@ -88,6 +90,8 @@ def start_test_site(debug=False):
r = static.File(root_dir)
r.putChild(b"redirect", util.Redirect(b"/redirected"))
r.putChild(b"redirected", static.Data(b"Redirected here", "text/plain"))
numbers = [str(x).encode("utf8") for x in range(2**14)]
r.putChild(b"numbers", static.Data(b"".join(numbers), "text/plain"))
port = reactor.listenTCP(0, server.Site(r), interface="127.0.0.1")
if debug:
@ -107,15 +111,20 @@ class CrawlerRun:
self.reqreached = []
self.itemerror = []
self.itemresp = []
self.signals_catched = {}
self.bytes = defaultdict(lambda: list())
self.signals_caught = {}
self.spider_class = spider_class
def run(self):
self.port = start_test_site()
self.portno = self.port.getHost().port
start_urls = [self.geturl("/"), self.geturl("/redirect"),
self.geturl("/redirect")] # a duplicate
start_urls = [
self.geturl("/"),
self.geturl("/redirect"),
self.geturl("/redirect"), # duplicate
self.geturl("/numbers"),
]
for name, signal in vars(signals).items():
if not name.startswith('_'):
@ -124,6 +133,7 @@ class CrawlerRun:
self.crawler = get_crawler(self.spider_class)
self.crawler.signals.connect(self.item_scraped, signals.item_scraped)
self.crawler.signals.connect(self.item_error, signals.item_error)
self.crawler.signals.connect(self.bytes_received, signals.bytes_received)
self.crawler.signals.connect(self.request_scheduled, signals.request_scheduled)
self.crawler.signals.connect(self.request_dropped, signals.request_dropped)
self.crawler.signals.connect(self.request_reached, signals.request_reached_downloader)
@ -155,6 +165,9 @@ class CrawlerRun:
def item_scraped(self, item, spider, response):
self.itemresp.append((item, response))
def bytes_received(self, data, request, spider):
self.bytes[request].append(data)
def request_scheduled(self, request, spider):
self.reqplug.append((request, spider))
@ -172,7 +185,7 @@ class CrawlerRun:
signalargs = kwargs.copy()
sig = signalargs.pop('signal')
signalargs.pop('sender', None)
self.signals_catched[sig] = signalargs
self.signals_caught[sig] = signalargs
class EngineTest(unittest.TestCase):
@ -183,16 +196,17 @@ class EngineTest(unittest.TestCase):
self.run = CrawlerRun(spider)
yield self.run.run()
self._assert_visited_urls()
self._assert_scheduled_requests(urls_to_visit=8)
self._assert_scheduled_requests(urls_to_visit=9)
self._assert_downloaded_responses()
self._assert_scraped_items()
self._assert_signals_catched()
self._assert_signals_caught()
self._assert_bytes_received()
@defer.inlineCallbacks
def test_crawler_dupefilter(self):
self.run = CrawlerRun(TestDupeFilterSpider)
yield self.run.run()
self._assert_scheduled_requests(urls_to_visit=7)
self._assert_scheduled_requests(urls_to_visit=8)
self._assert_dropped_requests()
@defer.inlineCallbacks
@ -229,8 +243,8 @@ class EngineTest(unittest.TestCase):
def _assert_downloaded_responses(self):
# response tests
self.assertEqual(8, len(self.run.respplug))
self.assertEqual(8, len(self.run.reqreached))
self.assertEqual(9, len(self.run.respplug))
self.assertEqual(9, len(self.run.reqreached))
for response, _ in self.run.respplug:
if self.run.getpath(response.url) == '/item999.html':
@ -263,19 +277,61 @@ class EngineTest(unittest.TestCase):
self.assertEqual('Item 2 name', item['name'])
self.assertEqual('200', item['price'])
def _assert_signals_catched(self):
assert signals.engine_started in self.run.signals_catched
assert signals.engine_stopped in self.run.signals_catched
assert signals.spider_opened in self.run.signals_catched
assert signals.spider_idle in self.run.signals_catched
assert signals.spider_closed in self.run.signals_catched
def _assert_bytes_received(self):
self.assertEqual(9, len(self.run.bytes))
for request, data in self.run.bytes.items():
joined_data = b"".join(data)
if self.run.getpath(request.url) == "/":
self.assertEqual(joined_data, get_testdata("test_site", "index.html"))
elif self.run.getpath(request.url) == "/item1.html":
self.assertEqual(joined_data, get_testdata("test_site", "item1.html"))
elif self.run.getpath(request.url) == "/item2.html":
self.assertEqual(joined_data, get_testdata("test_site", "item2.html"))
elif self.run.getpath(request.url) == "/redirected":
self.assertEqual(joined_data, b"Redirected here")
elif self.run.getpath(request.url) == '/redirect':
self.assertEqual(
joined_data,
b"\n<html>\n"
b" <head>\n"
b" <meta http-equiv=\"refresh\" content=\"0;URL=/redirected\">\n"
b" </head>\n"
b" <body bgcolor=\"#FFFFFF\" text=\"#000000\">\n"
b" <a href=\"/redirected\">click here</a>\n"
b" </body>\n"
b"</html>\n"
)
elif self.run.getpath(request.url) == "/tem999.html":
self.assertEqual(
joined_data,
b"\n<html>\n"
b" <head><title>404 - No Such Resource</title></head>\n"
b" <body>\n"
b" <h1>No Such Resource</h1>\n"
b" <p>File not found.</p>\n"
b" </body>\n"
b"</html>\n"
)
elif self.run.getpath(request.url) == "/numbers":
# signal was fired multiple times
self.assertTrue(len(data) > 1)
# bytes were received in order
numbers = [str(x).encode("utf8") for x in range(2**14)]
self.assertEqual(joined_data, b"".join(numbers))
def _assert_signals_caught(self):
assert signals.engine_started in self.run.signals_caught
assert signals.engine_stopped in self.run.signals_caught
assert signals.spider_opened in self.run.signals_caught
assert signals.spider_idle in self.run.signals_caught
assert signals.spider_closed in self.run.signals_caught
self.assertEqual({'spider': self.run.spider},
self.run.signals_catched[signals.spider_opened])
self.run.signals_caught[signals.spider_opened])
self.assertEqual({'spider': self.run.spider},
self.run.signals_catched[signals.spider_idle])
self.run.signals_caught[signals.spider_idle])
self.assertEqual({'spider': self.run.spider, 'reason': 'finished'},
self.run.signals_catched[signals.spider_closed])
self.run.signals_caught[signals.spider_closed])
@defer.inlineCallbacks
def test_close_downloader(self):