AsyncCrawlerProcess._start_asyncio() improvements. (#7366)

This commit is contained in:
Andrey Rakhmatullin 2026-03-27 16:21:25 +05:00 committed by GitHub
parent 5561aaec1d
commit fee20b7858
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
3 changed files with 199 additions and 31 deletions

View File

@ -605,22 +605,30 @@ class CrawlerProcessBase(CrawlerRunnerBase):
from twisted.internet import reactor from twisted.internet import reactor
install_shutdown_handlers(self._signal_kill) install_shutdown_handlers(self._signal_kill)
signame = signal_names[signum] self._log_shutdown(signum)
logger.info(
"Received %(signame)s, shutting down gracefully. Send again to force ",
{"signame": signame},
)
reactor.callFromThread(self._graceful_stop_reactor) reactor.callFromThread(self._graceful_stop_reactor)
def _signal_kill(self, signum: int, _: Any) -> None: def _signal_kill(self, signum: int, _: Any) -> None:
from twisted.internet import reactor from twisted.internet import reactor
install_shutdown_handlers(signal.SIG_IGN) install_shutdown_handlers(signal.SIG_IGN)
self._log_kill(signum)
reactor.callFromThread(self._stop_reactor)
@staticmethod
def _log_shutdown(signum: int) -> None:
signame = signal_names[signum]
logger.info(
"Received %(signame)s, shutting down gracefully. Send again to force ",
{"signame": signame},
)
@staticmethod
def _log_kill(signum: int) -> None:
signame = signal_names[signum] signame = signal_names[signum]
logger.info( logger.info(
"Received %(signame)s twice, forcing unclean shutdown", {"signame": signame} "Received %(signame)s twice, forcing unclean shutdown", {"signame": signame}
) )
reactor.callFromThread(self._stop_reactor)
def _setup_reactor(self, install_signal_handlers: bool) -> None: def _setup_reactor(self, install_signal_handlers: bool) -> None:
from twisted.internet import reactor from twisted.internet import reactor
@ -767,6 +775,7 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner):
): ):
super().__init__(settings, install_root_handler) super().__init__(settings, install_root_handler)
logger.debug("Using AsyncCrawlerProcess") logger.debug("Using AsyncCrawlerProcess")
self._reactorless_loop: asyncio.AbstractEventLoop | None = None
# We want the asyncio event loop to be installed early, so that it's # We want the asyncio event loop to be installed early, so that it's
# always the correct one. And as we do that, we can also install the # always the correct one. And as we do that, we can also install the
# reactor here. # reactor here.
@ -778,7 +787,7 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner):
raise RuntimeError( raise RuntimeError(
"TWISTED_ENABLED is False but a Twisted reactor is installed." "TWISTED_ENABLED is False but a Twisted reactor is installed."
) )
set_asyncio_event_loop(loop_path) self._reactorless_loop = set_asyncio_event_loop(loop_path)
install_reactor_import_hook() install_reactor_import_hook()
elif is_reactor_installed(): elif is_reactor_installed():
# The user could install a reactor before this class is instantiated. # The user could install a reactor before this class is instantiated.
@ -790,6 +799,7 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner):
else: else:
install_reactor(_asyncio_reactor_path, loop_path) install_reactor(_asyncio_reactor_path, loop_path)
self._initialized_reactor = True self._initialized_reactor = True
self._reactorless_main_task: asyncio.Future[None] | None = None
def _stop_dfd(self) -> Deferred[Any]: def _stop_dfd(self) -> Deferred[Any]:
return deferred_from_coro(self.stop()) return deferred_from_coro(self.stop())
@ -820,20 +830,149 @@ class AsyncCrawlerProcess(CrawlerProcessBase, AsyncCrawlerRunner):
def _start_asyncio( def _start_asyncio(
self, stop_after_crawl: bool, install_signal_handlers: bool self, stop_after_crawl: bool, install_signal_handlers: bool
) -> None: ) -> None:
# Very basic and will need multiple improvements. # We cannot use asyncio.run() here, because we can't let it handle the
# TODO https://docs.python.org/3/library/asyncio-runner.html#handling-keyboard-interruption # loop lifetime: _crawl() needs a loop (which we create in __init__()),
# TODO various exception handling # because crawl() returns a Task.
# TODO consider asyncio.run() # So we reproduce a part of asyncio.runners.Runner that is useful to us.
# Normal workflow:
# 1. _start_asyncio() creates a task for self.join() and calls _run_loop()
# 2. _run_loop() calls loop.run_until_complete(main_task)
# 3. Crawling tasks start and finish
# 4. join() completes, loop.run_until_complete() and thus _run_loop() return
# 5. _start_asyncio() calls _close_loop()
# 6. _close_loop() does finalization and calls loop.close()
# Normal workflow with stop_after_crawl=False:
# 1. _start_asyncio() creates a simple future and calls _run_loop()
# 2. _run_loop() calls loop.run_until_complete(main_task)
# 3. Crawling tasks start and finish
# 4. _run_loop() blocks until the loop is stopped externally or the
# future is cancelled via Ctrl-C
# 5. (after _run_loop() returns) _start_asyncio() calls _close_loop()
# 6. _close_loop() does finalization and calls loop.close()
# Workflow with Ctrl-C pressed once:
# 1. While loop.run_until_complete() blocks, _signal_shutdown_reactorless()
# is called
# 2. _signal_shutdown_reactorless() calls _shutdown_graceful_reactorless()
# (via call_soon_threadsafe())
# 3. _shutdown_graceful_reactorless() calls stop()
# 4. For stop_after_crawl=True: crawl tasks finish, join() completes,
# loop.run_until_complete() and thus _run_loop() return
# For stop_after_crawl=False: _shutdown_graceful_reactorless() waits
# for crawl tasks via join(), then cancels the main task,
# loop.run_until_complete() raises CancelledError, _run_loop() returns
# 5. _start_asyncio() calls _close_loop()
# 6. _close_loop() does finalization and calls loop.close()
# Workflow with Ctrl-C pressed twice:
# 1. While loop.run_until_complete() blocks, _signal_shutdown_reactorless()
# is called
# 2. _signal_shutdown_reactorless() calls _shutdown_graceful_reactorless()
# (via call_soon_threadsafe()) and installs _signal_kill_reactorless()
# as the next handler
# 3. Before _shutdown_graceful_reactorless() completes,
# _signal_kill_reactorless() is called
# 4. _signal_kill_reactorless() cancels the main task
# (via call_soon_threadsafe())
# 5. loop.run_until_complete() raises CancelledError, _run_loop() returns
# 6. _start_asyncio() calls _close_loop()
# 7. _close_loop() cancels all pending tasks (including
# _shutdown_graceful_reactorless()), does finalization and calls loop.close()
loop = self._reactorless_loop
assert loop
loop = asyncio.get_event_loop()
if stop_after_crawl: if stop_after_crawl:
join_task = loop.create_task(self.join()) self._reactorless_main_task = loop.create_task(self.join())
join_task.add_done_callback(lambda _: loop.stop()) else:
self._reactorless_main_task = loop.create_future()
self._stop_after_crawl = stop_after_crawl
try: try:
loop.run_forever() # blocking call self._run_loop(install_signal_handlers) # blocking call
except asyncio.CancelledError:
pass
finally: finally:
self._close_loop()
def _run_loop(self, install_signal_handlers: bool) -> None:
# similar to asyncio.runners.Runner.run()
if install_signal_handlers:
install_shutdown_handlers(self._signal_shutdown_reactorless)
assert self._reactorless_loop
assert self._reactorless_main_task
self._reactorless_loop.run_until_complete(self._reactorless_main_task)
def _close_loop(self) -> None:
# Similar to asyncio.runners.Runner.close()
loop = self._reactorless_loop
assert loop
try:
self._cancel_all_tasks(loop)
loop.run_until_complete(loop.shutdown_asyncgens()) loop.run_until_complete(loop.shutdown_asyncgens())
loop.run_until_complete(loop.shutdown_default_executor())
finally:
self._reactorless_main_task = None
asyncio.set_event_loop(None)
loop.close() loop.close()
self._reactorless_loop = None
@staticmethod
def _cancel_all_tasks(loop: asyncio.AbstractEventLoop) -> None:
# copy of asyncio.runners._cancel_all_tasks()
to_cancel = asyncio.all_tasks(loop)
if not to_cancel:
return
for task in to_cancel:
task.cancel()
loop.run_until_complete(asyncio.gather(*to_cancel, return_exceptions=True))
for task in to_cancel:
if task.cancelled():
continue
if task.exception() is not None:
loop.call_exception_handler(
{
"message": "unhandled exception during AsyncCrawlerProcess shutdown",
"exception": task.exception(),
"task": task,
}
)
def _signal_shutdown_reactorless(self, signum: int, _: Any) -> None:
install_shutdown_handlers(self._signal_kill_reactorless)
self._log_shutdown(signum)
if (loop := self._reactorless_loop) is None:
return
def _create_shutdown_task() -> None:
coro = self._shutdown_graceful_reactorless()
try:
loop.create_task(coro)
except RuntimeError:
coro.close()
loop.call_soon_threadsafe(_create_shutdown_task)
async def _shutdown_graceful_reactorless(self) -> None:
await self.stop()
if not self._stop_after_crawl:
# wait until crawl tasks finish and cancel the future
await self.join()
if self._reactorless_main_task and not self._reactorless_main_task.done():
self._reactorless_main_task.cancel()
def _signal_kill_reactorless(self, signum: int, _: Any) -> None:
install_shutdown_handlers(signal.SIG_IGN)
self._log_kill(signum)
if (loop := self._reactorless_loop) is None:
return
if (task := self._reactorless_main_task) is not None:
loop.call_soon_threadsafe(task.cancel)
def _start_twisted( def _start_twisted(
self, stop_after_crawl: bool, install_signal_handlers: bool self, stop_after_crawl: bool, install_signal_handlers: bool

View File

@ -0,0 +1,20 @@
import asyncio
import sys
import scrapy
from scrapy.crawler import AsyncCrawlerProcess
class SleepingSpider(scrapy.Spider):
name = "sleeping"
start_urls = ["data:,;"]
async def parse(self, response):
await asyncio.sleep(int(sys.argv[1]))
process = AsyncCrawlerProcess(settings={"TWISTED_ENABLED": False})
process.crawl(SleepingSpider)
process.start()

View File

@ -12,12 +12,10 @@ from typing import TYPE_CHECKING
import pytest import pytest
from packaging.version import parse as parse_version from packaging.version import parse as parse_version
from pexpect.popen_spawn import PopenSpawn from pexpect.popen_spawn import PopenSpawn
from twisted.internet.defer import Deferred
from w3lib import __version__ as w3lib_version from w3lib import __version__ as w3lib_version
from scrapy.utils.asyncio import call_later from tests.utils import async_sleep, get_script_run_env
from tests.utils import get_script_run_env from tests.utils.decorators import coroutine_test
from tests.utils.decorators import inline_callbacks_test
if TYPE_CHECKING: if TYPE_CHECKING:
from tests.mockserver.http import MockServer from tests.mockserver.http import MockServer
@ -207,33 +205,37 @@ class TestCrawlerProcessSubprocessBase(ScriptRunnerMixin):
assert "Spider closed (finished)" in log assert "Spider closed (finished)" in log
assert "The value of FOO is 42" in log assert "The value of FOO is 42" in log
def test_shutdown_graceful(self): def _test_shutdown_graceful(self, script: str = "sleeping.py") -> None:
sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK # type: ignore[attr-defined]
args = self.get_script_args("sleeping.py", "3") args = self.get_script_args(script, "3")
p = PopenSpawn(args, timeout=5, env=get_script_run_env()) p = PopenSpawn(args, timeout=5, env=get_script_run_env())
p.expect_exact("Spider opened") p.expect_exact("Spider opened")
p.expect_exact("Crawled (200)") p.expect_exact("Crawled (200)")
p.kill(sig) p.kill(sig)
p.expect_exact("shutting down gracefully") p.expect_exact("shutting down gracefully")
p.expect_exact("Spider closed (shutdown)") p.expect_exact("Spider closed (shutdown)")
p.wait() p.wait() # type: ignore[no-untyped-call]
@inline_callbacks_test def test_shutdown_graceful(self) -> None:
def test_shutdown_forced(self): self._test_shutdown_graceful()
sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK
args = self.get_script_args("sleeping.py", "10") async def _test_shutdown_forced(self, script: str = "sleeping.py") -> None:
sig = signal.SIGINT if sys.platform != "win32" else signal.SIGBREAK # type: ignore[attr-defined]
args = self.get_script_args(script, "10")
p = PopenSpawn(args, timeout=5, env=get_script_run_env()) p = PopenSpawn(args, timeout=5, env=get_script_run_env())
p.expect_exact("Spider opened") p.expect_exact("Spider opened")
p.expect_exact("Crawled (200)") p.expect_exact("Crawled (200)")
p.kill(sig) p.kill(sig)
p.expect_exact("shutting down gracefully") p.expect_exact("shutting down gracefully")
# sending the second signal too fast often causes problems # sending the second signal too fast often causes problems
d = Deferred() await async_sleep(0.01)
call_later(0.01, d.callback, None)
yield d
p.kill(sig) p.kill(sig)
p.expect_exact("forcing unclean shutdown") p.expect_exact("forcing unclean shutdown")
p.wait() p.wait() # type: ignore[no-untyped-call]
@coroutine_test
async def test_shutdown_forced(self) -> None:
await self._test_shutdown_forced()
class TestCrawlerProcessSubprocess(TestCrawlerProcessSubprocessBase): class TestCrawlerProcessSubprocess(TestCrawlerProcessSubprocessBase):
@ -397,6 +399,13 @@ class TestAsyncCrawlerProcessSubprocess(TestCrawlerProcessSubprocessBase):
in log in log
) )
def test_shutdown_graceful(self) -> None:
self._test_shutdown_graceful("reactorless_sleeping.py")
@coroutine_test
async def test_shutdown_forced(self) -> None:
await self._test_shutdown_forced("reactorless_sleeping.py")
class TestCrawlerRunnerSubprocessBase(ScriptRunnerMixin): class TestCrawlerRunnerSubprocessBase(ScriptRunnerMixin):
"""Common tests between CrawlerRunner and AsyncCrawlerRunner, """Common tests between CrawlerRunner and AsyncCrawlerRunner,