mirror of https://github.com/scrapy/scrapy.git
Refactor chainDeferred() usages (#7008)
This commit is contained in:
parent
d3e15a10cf
commit
3c546bdb82
|
|
@ -8,6 +8,7 @@ from time import time
|
|||
from typing import TYPE_CHECKING, Any, cast
|
||||
|
||||
from twisted.internet.defer import Deferred, inlineCallbacks
|
||||
from twisted.python.failure import Failure
|
||||
|
||||
from scrapy import Request, Spider, signals
|
||||
from scrapy.core.downloader.handlers import DownloadHandlers
|
||||
|
|
@ -22,7 +23,7 @@ from scrapy.utils.asyncio import (
|
|||
)
|
||||
from scrapy.utils.defer import (
|
||||
_defer_sleep_async,
|
||||
deferred_from_coro,
|
||||
_schedule_coro,
|
||||
maybe_deferred_to_future,
|
||||
)
|
||||
from scrapy.utils.httpobj import urlparse_cached
|
||||
|
|
@ -200,6 +201,7 @@ class Downloader:
|
|||
)
|
||||
return self.get_slot_key(request)
|
||||
|
||||
# passed as download_func into self.middleware.download() in self.fetch()
|
||||
@inlineCallbacks
|
||||
def _enqueue_request(
|
||||
self, request: Request
|
||||
|
|
@ -216,7 +218,7 @@ class Downloader:
|
|||
slot.queue.append((request, d))
|
||||
self._process_queue(slot)
|
||||
try:
|
||||
return (yield d)
|
||||
return (yield d) # fired in _wait_for_download()
|
||||
finally:
|
||||
slot.active.remove(request)
|
||||
|
||||
|
|
@ -237,9 +239,8 @@ class Downloader:
|
|||
# Process enqueued requests if there are free slots to transfer for this slot
|
||||
while slot.queue and slot.free_transfer_slots() > 0:
|
||||
slot.lastseen = now
|
||||
request, deferred = slot.queue.popleft()
|
||||
dfd = deferred_from_coro(self._download(slot, request))
|
||||
dfd.chainDeferred(deferred)
|
||||
request, queue_dfd = slot.queue.popleft()
|
||||
_schedule_coro(self._wait_for_download(slot, request, queue_dfd))
|
||||
# prevent burst if inter-request delays were configured
|
||||
if delay:
|
||||
self._process_queue(slot)
|
||||
|
|
@ -282,6 +283,16 @@ class Downloader:
|
|||
spider=self.crawler.spider,
|
||||
)
|
||||
|
||||
async def _wait_for_download(
|
||||
self, slot: Slot, request: Request, queue_dfd: Deferred[Response]
|
||||
) -> None:
|
||||
try:
|
||||
response = await self._download(slot, request)
|
||||
except Exception:
|
||||
queue_dfd.errback(Failure())
|
||||
else:
|
||||
queue_dfd.callback(response) # awaited in _enqueue_request()
|
||||
|
||||
def close(self) -> None:
|
||||
self._slot_gc_loop.stop()
|
||||
for slot in self.slots.values():
|
||||
|
|
|
|||
|
|
@ -27,7 +27,6 @@ from scrapy.utils.defer import (
|
|||
_defer_sleep_async,
|
||||
_schedule_coro,
|
||||
aiter_errback,
|
||||
deferred_f_from_coro_f,
|
||||
deferred_from_coro,
|
||||
ensure_awaitable,
|
||||
iter_errback,
|
||||
|
|
@ -235,7 +234,7 @@ class Scraper:
|
|||
dfd = self.slot.add_response_request(result, request)
|
||||
self._scrape_next()
|
||||
try:
|
||||
yield dfd
|
||||
yield dfd # fired in _wait_for_processing()
|
||||
except Exception:
|
||||
logger.error(
|
||||
"Scraper bug processing %(request)s",
|
||||
|
|
@ -251,10 +250,9 @@ class Scraper:
|
|||
def _scrape_next(self) -> None:
|
||||
assert self.slot is not None # typing
|
||||
while self.slot.queue:
|
||||
result, request, deferred = self.slot.next_response_request_deferred()
|
||||
self._scrape(result, request).chainDeferred(deferred)
|
||||
result, request, queue_dfd = self.slot.next_response_request_deferred()
|
||||
_schedule_coro(self._wait_for_processing(result, request, queue_dfd))
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def _scrape(self, result: Response | Failure, request: Request) -> None:
|
||||
"""Handle the downloaded response or failure through the spider callback/errback."""
|
||||
if not isinstance(result, (Response, Failure)):
|
||||
|
|
@ -296,6 +294,16 @@ class Scraper:
|
|||
else:
|
||||
await self.handle_spider_output_async(output, request, result)
|
||||
|
||||
async def _wait_for_processing(
|
||||
self, result: Response | Failure, request: Request, queue_dfd: Deferred[None]
|
||||
) -> None:
|
||||
try:
|
||||
await self._scrape(result, request)
|
||||
except Exception:
|
||||
queue_dfd.errback(Failure())
|
||||
else:
|
||||
queue_dfd.callback(None) # awaited in enqueue_scrape()
|
||||
|
||||
def call_spider(
|
||||
self, result: Response | Failure, request: Request, spider: Spider | None = None
|
||||
) -> Deferred[Iterable[Any] | AsyncIterator[Any]]:
|
||||
|
|
|
|||
|
|
@ -0,0 +1,27 @@
|
|||
from __future__ import annotations
|
||||
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from scrapy.utils.defer import deferred_f_from_coro_f, maybe_deferred_to_future
|
||||
from scrapy.utils.test import get_crawler
|
||||
from tests.spiders import SimpleSpider
|
||||
|
||||
if TYPE_CHECKING:
|
||||
import pytest
|
||||
|
||||
from tests.mockserver.http import MockServer
|
||||
|
||||
|
||||
@deferred_f_from_coro_f
|
||||
async def test_scraper_exception(
|
||||
mockserver: MockServer,
|
||||
caplog: pytest.LogCaptureFixture,
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
crawler = get_crawler(SimpleSpider)
|
||||
monkeypatch.setattr(
|
||||
"scrapy.core.engine.Scraper.handle_spider_output_async",
|
||||
lambda *args, **kwargs: 1 / 0,
|
||||
)
|
||||
await maybe_deferred_to_future(crawler.crawl(url=mockserver.url("/")))
|
||||
assert "Scraper bug processing" in caplog.text
|
||||
Loading…
Reference in New Issue