mirror of https://github.com/scrapy/scrapy.git
Full typing for scrapy/extensions, part 1. (#6276)
This commit is contained in:
parent
8985a04bd1
commit
421e08dd4a
|
|
@ -4,20 +4,31 @@ conditions are met.
|
||||||
See documentation in docs/topics/extensions.rst
|
See documentation in docs/topics/extensions.rst
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
from collections import defaultdict
|
from collections import defaultdict
|
||||||
|
from typing import TYPE_CHECKING, Any, DefaultDict, Dict
|
||||||
|
|
||||||
from scrapy import signals
|
from twisted.python.failure import Failure
|
||||||
|
|
||||||
|
from scrapy import Request, Spider, signals
|
||||||
|
from scrapy.crawler import Crawler
|
||||||
from scrapy.exceptions import NotConfigured
|
from scrapy.exceptions import NotConfigured
|
||||||
|
from scrapy.http import Response
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
# typing.Self requires Python 3.11
|
||||||
|
from typing_extensions import Self
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class CloseSpider:
|
class CloseSpider:
|
||||||
def __init__(self, crawler):
|
def __init__(self, crawler: Crawler):
|
||||||
self.crawler = crawler
|
self.crawler: Crawler = crawler
|
||||||
|
|
||||||
self.close_on = {
|
self.close_on: Dict[str, Any] = {
|
||||||
"timeout": crawler.settings.getfloat("CLOSESPIDER_TIMEOUT"),
|
"timeout": crawler.settings.getfloat("CLOSESPIDER_TIMEOUT"),
|
||||||
"itemcount": crawler.settings.getint("CLOSESPIDER_ITEMCOUNT"),
|
"itemcount": crawler.settings.getint("CLOSESPIDER_ITEMCOUNT"),
|
||||||
"pagecount": crawler.settings.getint("CLOSESPIDER_PAGECOUNT"),
|
"pagecount": crawler.settings.getint("CLOSESPIDER_PAGECOUNT"),
|
||||||
|
|
@ -28,7 +39,7 @@ class CloseSpider:
|
||||||
if not any(self.close_on.values()):
|
if not any(self.close_on.values()):
|
||||||
raise NotConfigured
|
raise NotConfigured
|
||||||
|
|
||||||
self.counter = defaultdict(int)
|
self.counter: DefaultDict[str, int] = defaultdict(int)
|
||||||
|
|
||||||
if self.close_on.get("errorcount"):
|
if self.close_on.get("errorcount"):
|
||||||
crawler.signals.connect(self.error_count, signal=signals.spider_error)
|
crawler.signals.connect(self.error_count, signal=signals.spider_error)
|
||||||
|
|
@ -39,8 +50,8 @@ class CloseSpider:
|
||||||
if self.close_on.get("itemcount"):
|
if self.close_on.get("itemcount"):
|
||||||
crawler.signals.connect(self.item_scraped, signal=signals.item_scraped)
|
crawler.signals.connect(self.item_scraped, signal=signals.item_scraped)
|
||||||
if self.close_on.get("timeout_no_item"):
|
if self.close_on.get("timeout_no_item"):
|
||||||
self.timeout_no_item = self.close_on["timeout_no_item"]
|
self.timeout_no_item: int = self.close_on["timeout_no_item"]
|
||||||
self.items_in_period = 0
|
self.items_in_period: int = 0
|
||||||
crawler.signals.connect(
|
crawler.signals.connect(
|
||||||
self.spider_opened_no_item, signal=signals.spider_opened
|
self.spider_opened_no_item, signal=signals.spider_opened
|
||||||
)
|
)
|
||||||
|
|
@ -50,22 +61,25 @@ class CloseSpider:
|
||||||
crawler.signals.connect(self.spider_closed, signal=signals.spider_closed)
|
crawler.signals.connect(self.spider_closed, signal=signals.spider_closed)
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_crawler(cls, crawler):
|
def from_crawler(cls, crawler: Crawler) -> Self:
|
||||||
return cls(crawler)
|
return cls(crawler)
|
||||||
|
|
||||||
def error_count(self, failure, response, spider):
|
def error_count(self, failure: Failure, response: Response, spider: Spider) -> None:
|
||||||
self.counter["errorcount"] += 1
|
self.counter["errorcount"] += 1
|
||||||
if self.counter["errorcount"] == self.close_on["errorcount"]:
|
if self.counter["errorcount"] == self.close_on["errorcount"]:
|
||||||
|
assert self.crawler.engine
|
||||||
self.crawler.engine.close_spider(spider, "closespider_errorcount")
|
self.crawler.engine.close_spider(spider, "closespider_errorcount")
|
||||||
|
|
||||||
def page_count(self, response, request, spider):
|
def page_count(self, response: Response, request: Request, spider: Spider) -> None:
|
||||||
self.counter["pagecount"] += 1
|
self.counter["pagecount"] += 1
|
||||||
if self.counter["pagecount"] == self.close_on["pagecount"]:
|
if self.counter["pagecount"] == self.close_on["pagecount"]:
|
||||||
|
assert self.crawler.engine
|
||||||
self.crawler.engine.close_spider(spider, "closespider_pagecount")
|
self.crawler.engine.close_spider(spider, "closespider_pagecount")
|
||||||
|
|
||||||
def spider_opened(self, spider):
|
def spider_opened(self, spider: Spider) -> None:
|
||||||
from twisted.internet import reactor
|
from twisted.internet import reactor
|
||||||
|
|
||||||
|
assert self.crawler.engine
|
||||||
self.task = reactor.callLater(
|
self.task = reactor.callLater(
|
||||||
self.close_on["timeout"],
|
self.close_on["timeout"],
|
||||||
self.crawler.engine.close_spider,
|
self.crawler.engine.close_spider,
|
||||||
|
|
@ -73,21 +87,22 @@ class CloseSpider:
|
||||||
reason="closespider_timeout",
|
reason="closespider_timeout",
|
||||||
)
|
)
|
||||||
|
|
||||||
def item_scraped(self, item, spider):
|
def item_scraped(self, item: Any, spider: Spider) -> None:
|
||||||
self.counter["itemcount"] += 1
|
self.counter["itemcount"] += 1
|
||||||
if self.counter["itemcount"] == self.close_on["itemcount"]:
|
if self.counter["itemcount"] == self.close_on["itemcount"]:
|
||||||
|
assert self.crawler.engine
|
||||||
self.crawler.engine.close_spider(spider, "closespider_itemcount")
|
self.crawler.engine.close_spider(spider, "closespider_itemcount")
|
||||||
|
|
||||||
def spider_closed(self, spider):
|
def spider_closed(self, spider: Spider) -> None:
|
||||||
task = getattr(self, "task", False)
|
task = getattr(self, "task", None)
|
||||||
if task and task.active():
|
if task and task.active():
|
||||||
task.cancel()
|
task.cancel()
|
||||||
|
|
||||||
task_no_item = getattr(self, "task_no_item", False)
|
task_no_item = getattr(self, "task_no_item", None)
|
||||||
if task_no_item and task_no_item.running:
|
if task_no_item and task_no_item.running:
|
||||||
task_no_item.stop()
|
task_no_item.stop()
|
||||||
|
|
||||||
def spider_opened_no_item(self, spider):
|
def spider_opened_no_item(self, spider: Spider) -> None:
|
||||||
from twisted.internet import task
|
from twisted.internet import task
|
||||||
|
|
||||||
self.task_no_item = task.LoopingCall(self._count_items_produced, spider)
|
self.task_no_item = task.LoopingCall(self._count_items_produced, spider)
|
||||||
|
|
@ -98,10 +113,10 @@ class CloseSpider:
|
||||||
f"{self.timeout_no_item} seconds."
|
f"{self.timeout_no_item} seconds."
|
||||||
)
|
)
|
||||||
|
|
||||||
def item_scraped_no_item(self, item, spider):
|
def item_scraped_no_item(self, item: Any, spider: Spider) -> None:
|
||||||
self.items_in_period += 1
|
self.items_in_period += 1
|
||||||
|
|
||||||
def _count_items_produced(self, spider):
|
def _count_items_produced(self, spider: Spider) -> None:
|
||||||
if self.items_in_period >= 1:
|
if self.items_in_period >= 1:
|
||||||
self.items_in_period = 0
|
self.items_in_period = 0
|
||||||
else:
|
else:
|
||||||
|
|
@ -109,4 +124,5 @@ class CloseSpider:
|
||||||
f"Closing spider since no items were produced in the last "
|
f"Closing spider since no items were produced in the last "
|
||||||
f"{self.timeout_no_item} seconds."
|
f"{self.timeout_no_item} seconds."
|
||||||
)
|
)
|
||||||
|
assert self.crawler.engine
|
||||||
self.crawler.engine.close_spider(spider, "closespider_timeout_no_item")
|
self.crawler.engine.close_spider(spider, "closespider_timeout_no_item")
|
||||||
|
|
|
||||||
|
|
@ -2,18 +2,28 @@
|
||||||
Extension for collecting core stats like items scraped and start/finish times
|
Extension for collecting core stats like items scraped and start/finish times
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from datetime import datetime, timezone
|
from __future__ import annotations
|
||||||
|
|
||||||
from scrapy import signals
|
from datetime import datetime, timezone
|
||||||
|
from typing import TYPE_CHECKING, Any, Optional
|
||||||
|
|
||||||
|
from scrapy import Spider, signals
|
||||||
|
from scrapy.crawler import Crawler
|
||||||
|
from scrapy.statscollectors import StatsCollector
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
# typing.Self requires Python 3.11
|
||||||
|
from typing_extensions import Self
|
||||||
|
|
||||||
|
|
||||||
class CoreStats:
|
class CoreStats:
|
||||||
def __init__(self, stats):
|
def __init__(self, stats: StatsCollector):
|
||||||
self.stats = stats
|
self.stats: StatsCollector = stats
|
||||||
self.start_time = None
|
self.start_time: Optional[datetime] = None
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_crawler(cls, crawler):
|
def from_crawler(cls, crawler: Crawler) -> Self:
|
||||||
|
assert crawler.stats
|
||||||
o = cls(crawler.stats)
|
o = cls(crawler.stats)
|
||||||
crawler.signals.connect(o.spider_opened, signal=signals.spider_opened)
|
crawler.signals.connect(o.spider_opened, signal=signals.spider_opened)
|
||||||
crawler.signals.connect(o.spider_closed, signal=signals.spider_closed)
|
crawler.signals.connect(o.spider_closed, signal=signals.spider_closed)
|
||||||
|
|
@ -22,11 +32,12 @@ class CoreStats:
|
||||||
crawler.signals.connect(o.response_received, signal=signals.response_received)
|
crawler.signals.connect(o.response_received, signal=signals.response_received)
|
||||||
return o
|
return o
|
||||||
|
|
||||||
def spider_opened(self, spider):
|
def spider_opened(self, spider: Spider) -> None:
|
||||||
self.start_time = datetime.now(tz=timezone.utc)
|
self.start_time = datetime.now(tz=timezone.utc)
|
||||||
self.stats.set_value("start_time", self.start_time, spider=spider)
|
self.stats.set_value("start_time", self.start_time, spider=spider)
|
||||||
|
|
||||||
def spider_closed(self, spider, reason):
|
def spider_closed(self, spider: Spider, reason: str) -> None:
|
||||||
|
assert self.start_time is not None
|
||||||
finish_time = datetime.now(tz=timezone.utc)
|
finish_time = datetime.now(tz=timezone.utc)
|
||||||
elapsed_time = finish_time - self.start_time
|
elapsed_time = finish_time - self.start_time
|
||||||
elapsed_time_seconds = elapsed_time.total_seconds()
|
elapsed_time_seconds = elapsed_time.total_seconds()
|
||||||
|
|
@ -36,13 +47,13 @@ class CoreStats:
|
||||||
self.stats.set_value("finish_time", finish_time, spider=spider)
|
self.stats.set_value("finish_time", finish_time, spider=spider)
|
||||||
self.stats.set_value("finish_reason", reason, spider=spider)
|
self.stats.set_value("finish_reason", reason, spider=spider)
|
||||||
|
|
||||||
def item_scraped(self, item, spider):
|
def item_scraped(self, item: Any, spider: Spider) -> None:
|
||||||
self.stats.inc_value("item_scraped_count", spider=spider)
|
self.stats.inc_value("item_scraped_count", spider=spider)
|
||||||
|
|
||||||
def response_received(self, spider):
|
def response_received(self, spider: Spider) -> None:
|
||||||
self.stats.inc_value("response_received_count", spider=spider)
|
self.stats.inc_value("response_received_count", spider=spider)
|
||||||
|
|
||||||
def item_dropped(self, item, spider, exception):
|
def item_dropped(self, item: Any, spider: Spider, exception: BaseException) -> None:
|
||||||
reason = exception.__class__.__name__
|
reason = exception.__class__.__name__
|
||||||
self.stats.inc_value("item_dropped_count", spider=spider)
|
self.stats.inc_value("item_dropped_count", spider=spider)
|
||||||
self.stats.inc_value(f"item_dropped_reasons_count/{reason}", spider=spider)
|
self.stats.inc_value(f"item_dropped_reasons_count/{reason}", spider=spider)
|
||||||
|
|
|
||||||
|
|
@ -4,22 +4,31 @@ Extensions for debugging Scrapy
|
||||||
See documentation in docs/topics/extensions.rst
|
See documentation in docs/topics/extensions.rst
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
import signal
|
import signal
|
||||||
import sys
|
import sys
|
||||||
import threading
|
import threading
|
||||||
import traceback
|
import traceback
|
||||||
from pdb import Pdb
|
from pdb import Pdb
|
||||||
|
from types import FrameType
|
||||||
|
from typing import TYPE_CHECKING, Optional
|
||||||
|
|
||||||
|
from scrapy.crawler import Crawler
|
||||||
from scrapy.utils.engine import format_engine_status
|
from scrapy.utils.engine import format_engine_status
|
||||||
from scrapy.utils.trackref import format_live_refs
|
from scrapy.utils.trackref import format_live_refs
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
# typing.Self requires Python 3.11
|
||||||
|
from typing_extensions import Self
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class StackTraceDump:
|
class StackTraceDump:
|
||||||
def __init__(self, crawler=None):
|
def __init__(self, crawler: Crawler):
|
||||||
self.crawler = crawler
|
self.crawler: Crawler = crawler
|
||||||
try:
|
try:
|
||||||
signal.signal(signal.SIGUSR2, self.dump_stacktrace)
|
signal.signal(signal.SIGUSR2, self.dump_stacktrace)
|
||||||
signal.signal(signal.SIGQUIT, self.dump_stacktrace)
|
signal.signal(signal.SIGQUIT, self.dump_stacktrace)
|
||||||
|
|
@ -28,10 +37,11 @@ class StackTraceDump:
|
||||||
pass
|
pass
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_crawler(cls, crawler):
|
def from_crawler(cls, crawler: Crawler) -> Self:
|
||||||
return cls(crawler)
|
return cls(crawler)
|
||||||
|
|
||||||
def dump_stacktrace(self, signum, frame):
|
def dump_stacktrace(self, signum: int, frame: Optional[FrameType]) -> None:
|
||||||
|
assert self.crawler.engine
|
||||||
log_args = {
|
log_args = {
|
||||||
"stackdumps": self._thread_stacks(),
|
"stackdumps": self._thread_stacks(),
|
||||||
"enginestatus": format_engine_status(self.crawler.engine),
|
"enginestatus": format_engine_status(self.crawler.engine),
|
||||||
|
|
@ -44,7 +54,7 @@ class StackTraceDump:
|
||||||
extra={"crawler": self.crawler},
|
extra={"crawler": self.crawler},
|
||||||
)
|
)
|
||||||
|
|
||||||
def _thread_stacks(self):
|
def _thread_stacks(self) -> str:
|
||||||
id2name = dict((th.ident, th.name) for th in threading.enumerate())
|
id2name = dict((th.ident, th.name) for th in threading.enumerate())
|
||||||
dumps = ""
|
dumps = ""
|
||||||
for id_, frame in sys._current_frames().items():
|
for id_, frame in sys._current_frames().items():
|
||||||
|
|
@ -55,12 +65,13 @@ class StackTraceDump:
|
||||||
|
|
||||||
|
|
||||||
class Debugger:
|
class Debugger:
|
||||||
def __init__(self):
|
def __init__(self) -> None:
|
||||||
try:
|
try:
|
||||||
signal.signal(signal.SIGUSR2, self._enter_debugger)
|
signal.signal(signal.SIGUSR2, self._enter_debugger)
|
||||||
except AttributeError:
|
except AttributeError:
|
||||||
# win32 platforms don't support SIGUSR signals
|
# win32 platforms don't support SIGUSR signals
|
||||||
pass
|
pass
|
||||||
|
|
||||||
def _enter_debugger(self, signum, frame):
|
def _enter_debugger(self, signum: int, frame: Optional[FrameType]) -> None:
|
||||||
|
assert frame
|
||||||
Pdb().set_trace(frame.f_back)
|
Pdb().set_trace(frame.f_back)
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,18 @@
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
from typing import TYPE_CHECKING, Optional, Tuple, Union
|
||||||
|
|
||||||
from twisted.internet import task
|
from twisted.internet import task
|
||||||
|
|
||||||
from scrapy import signals
|
from scrapy import Spider, signals
|
||||||
|
from scrapy.crawler import Crawler
|
||||||
from scrapy.exceptions import NotConfigured
|
from scrapy.exceptions import NotConfigured
|
||||||
|
from scrapy.statscollectors import StatsCollector
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
# typing.Self requires Python 3.11
|
||||||
|
from typing_extensions import Self
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
@ -14,30 +23,31 @@ class LogStats:
|
||||||
* IPM - Items per Minute
|
* IPM - Items per Minute
|
||||||
"""
|
"""
|
||||||
|
|
||||||
def __init__(self, stats, interval=60.0):
|
def __init__(self, stats: StatsCollector, interval: float = 60.0):
|
||||||
self.stats = stats
|
self.stats: StatsCollector = stats
|
||||||
self.interval = interval
|
self.interval: float = interval
|
||||||
self.multiplier = 60.0 / self.interval
|
self.multiplier: float = 60.0 / self.interval
|
||||||
self.task = None
|
self.task: Optional[task.LoopingCall] = None
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_crawler(cls, crawler):
|
def from_crawler(cls, crawler: Crawler) -> Self:
|
||||||
interval = crawler.settings.getfloat("LOGSTATS_INTERVAL")
|
interval: float = crawler.settings.getfloat("LOGSTATS_INTERVAL")
|
||||||
if not interval:
|
if not interval:
|
||||||
raise NotConfigured
|
raise NotConfigured
|
||||||
|
assert crawler.stats
|
||||||
o = cls(crawler.stats, interval)
|
o = cls(crawler.stats, interval)
|
||||||
crawler.signals.connect(o.spider_opened, signal=signals.spider_opened)
|
crawler.signals.connect(o.spider_opened, signal=signals.spider_opened)
|
||||||
crawler.signals.connect(o.spider_closed, signal=signals.spider_closed)
|
crawler.signals.connect(o.spider_closed, signal=signals.spider_closed)
|
||||||
return o
|
return o
|
||||||
|
|
||||||
def spider_opened(self, spider):
|
def spider_opened(self, spider: Spider) -> None:
|
||||||
self.pagesprev = 0
|
self.pagesprev: int = 0
|
||||||
self.itemsprev = 0
|
self.itemsprev: int = 0
|
||||||
|
|
||||||
self.task = task.LoopingCall(self.log, spider)
|
self.task = task.LoopingCall(self.log, spider)
|
||||||
self.task.start(self.interval)
|
self.task.start(self.interval)
|
||||||
|
|
||||||
def log(self, spider):
|
def log(self, spider: Spider) -> None:
|
||||||
self.calculate_stats()
|
self.calculate_stats()
|
||||||
|
|
||||||
msg = (
|
msg = (
|
||||||
|
|
@ -52,14 +62,14 @@ class LogStats:
|
||||||
}
|
}
|
||||||
logger.info(msg, log_args, extra={"spider": spider})
|
logger.info(msg, log_args, extra={"spider": spider})
|
||||||
|
|
||||||
def calculate_stats(self):
|
def calculate_stats(self) -> None:
|
||||||
self.items = self.stats.get_value("item_scraped_count", 0)
|
self.items: int = self.stats.get_value("item_scraped_count", 0)
|
||||||
self.pages = self.stats.get_value("response_received_count", 0)
|
self.pages: int = self.stats.get_value("response_received_count", 0)
|
||||||
self.irate = (self.items - self.itemsprev) * self.multiplier
|
self.irate: float = (self.items - self.itemsprev) * self.multiplier
|
||||||
self.prate = (self.pages - self.pagesprev) * self.multiplier
|
self.prate: float = (self.pages - self.pagesprev) * self.multiplier
|
||||||
self.pagesprev, self.itemsprev = self.pages, self.items
|
self.pagesprev, self.itemsprev = self.pages, self.items
|
||||||
|
|
||||||
def spider_closed(self, spider, reason):
|
def spider_closed(self, spider: Spider, reason: str) -> None:
|
||||||
if self.task and self.task.running:
|
if self.task and self.task.running:
|
||||||
self.task.stop()
|
self.task.stop()
|
||||||
|
|
||||||
|
|
@ -67,7 +77,9 @@ class LogStats:
|
||||||
self.stats.set_value("responses_per_minute", rpm_final)
|
self.stats.set_value("responses_per_minute", rpm_final)
|
||||||
self.stats.set_value("items_per_minute", ipm_final)
|
self.stats.set_value("items_per_minute", ipm_final)
|
||||||
|
|
||||||
def calculate_final_stats(self, spider):
|
def calculate_final_stats(
|
||||||
|
self, spider: Spider
|
||||||
|
) -> Union[Tuple[None, None], Tuple[float, float]]:
|
||||||
start_time = self.stats.get_value("start_time")
|
start_time = self.stats.get_value("start_time")
|
||||||
finished_time = self.stats.get_value("finished_time")
|
finished_time = self.stats.get_value("finished_time")
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -4,26 +4,36 @@ MemoryDebugger extension
|
||||||
See documentation in docs/topics/extensions.rst
|
See documentation in docs/topics/extensions.rst
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import gc
|
from __future__ import annotations
|
||||||
|
|
||||||
from scrapy import signals
|
import gc
|
||||||
|
from typing import TYPE_CHECKING
|
||||||
|
|
||||||
|
from scrapy import Spider, signals
|
||||||
|
from scrapy.crawler import Crawler
|
||||||
from scrapy.exceptions import NotConfigured
|
from scrapy.exceptions import NotConfigured
|
||||||
|
from scrapy.statscollectors import StatsCollector
|
||||||
from scrapy.utils.trackref import live_refs
|
from scrapy.utils.trackref import live_refs
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
# typing.Self requires Python 3.11
|
||||||
|
from typing_extensions import Self
|
||||||
|
|
||||||
|
|
||||||
class MemoryDebugger:
|
class MemoryDebugger:
|
||||||
def __init__(self, stats):
|
def __init__(self, stats: StatsCollector):
|
||||||
self.stats = stats
|
self.stats: StatsCollector = stats
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_crawler(cls, crawler):
|
def from_crawler(cls, crawler: Crawler) -> Self:
|
||||||
if not crawler.settings.getbool("MEMDEBUG_ENABLED"):
|
if not crawler.settings.getbool("MEMDEBUG_ENABLED"):
|
||||||
raise NotConfigured
|
raise NotConfigured
|
||||||
|
assert crawler.stats
|
||||||
o = cls(crawler.stats)
|
o = cls(crawler.stats)
|
||||||
crawler.signals.connect(o.spider_closed, signal=signals.spider_closed)
|
crawler.signals.connect(o.spider_closed, signal=signals.spider_closed)
|
||||||
return o
|
return o
|
||||||
|
|
||||||
def spider_closed(self, spider, reason):
|
def spider_closed(self, spider: Spider, reason: str) -> None:
|
||||||
gc.collect()
|
gc.collect()
|
||||||
self.stats.set_value(
|
self.stats.set_value(
|
||||||
"memdebug/gc_garbage_count", len(gc.garbage), spider=spider
|
"memdebug/gc_garbage_count", len(gc.garbage), spider=spider
|
||||||
|
|
|
||||||
|
|
@ -4,24 +4,32 @@ MemoryUsage extension
|
||||||
See documentation in docs/topics/extensions.rst
|
See documentation in docs/topics/extensions.rst
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
import socket
|
import socket
|
||||||
import sys
|
import sys
|
||||||
from importlib import import_module
|
from importlib import import_module
|
||||||
from pprint import pformat
|
from pprint import pformat
|
||||||
|
from typing import TYPE_CHECKING, List
|
||||||
|
|
||||||
from twisted.internet import task
|
from twisted.internet import task
|
||||||
|
|
||||||
from scrapy import signals
|
from scrapy import signals
|
||||||
|
from scrapy.crawler import Crawler
|
||||||
from scrapy.exceptions import NotConfigured
|
from scrapy.exceptions import NotConfigured
|
||||||
from scrapy.mail import MailSender
|
from scrapy.mail import MailSender
|
||||||
from scrapy.utils.engine import get_engine_status
|
from scrapy.utils.engine import get_engine_status
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
# typing.Self requires Python 3.11
|
||||||
|
from typing_extensions import Self
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class MemoryUsage:
|
class MemoryUsage:
|
||||||
def __init__(self, crawler):
|
def __init__(self, crawler: Crawler):
|
||||||
if not crawler.settings.getbool("MEMUSAGE_ENABLED"):
|
if not crawler.settings.getbool("MEMUSAGE_ENABLED"):
|
||||||
raise NotConfigured
|
raise NotConfigured
|
||||||
try:
|
try:
|
||||||
|
|
@ -30,32 +38,33 @@ class MemoryUsage:
|
||||||
except ImportError:
|
except ImportError:
|
||||||
raise NotConfigured
|
raise NotConfigured
|
||||||
|
|
||||||
self.crawler = crawler
|
self.crawler: Crawler = crawler
|
||||||
self.warned = False
|
self.warned: bool = False
|
||||||
self.notify_mails = crawler.settings.getlist("MEMUSAGE_NOTIFY_MAIL")
|
self.notify_mails: List[str] = crawler.settings.getlist("MEMUSAGE_NOTIFY_MAIL")
|
||||||
self.limit = crawler.settings.getint("MEMUSAGE_LIMIT_MB") * 1024 * 1024
|
self.limit: int = crawler.settings.getint("MEMUSAGE_LIMIT_MB") * 1024 * 1024
|
||||||
self.warning = crawler.settings.getint("MEMUSAGE_WARNING_MB") * 1024 * 1024
|
self.warning: int = crawler.settings.getint("MEMUSAGE_WARNING_MB") * 1024 * 1024
|
||||||
self.check_interval = crawler.settings.getfloat(
|
self.check_interval: float = crawler.settings.getfloat(
|
||||||
"MEMUSAGE_CHECK_INTERVAL_SECONDS"
|
"MEMUSAGE_CHECK_INTERVAL_SECONDS"
|
||||||
)
|
)
|
||||||
self.mail = MailSender.from_settings(crawler.settings)
|
self.mail: MailSender = MailSender.from_settings(crawler.settings)
|
||||||
crawler.signals.connect(self.engine_started, signal=signals.engine_started)
|
crawler.signals.connect(self.engine_started, signal=signals.engine_started)
|
||||||
crawler.signals.connect(self.engine_stopped, signal=signals.engine_stopped)
|
crawler.signals.connect(self.engine_stopped, signal=signals.engine_stopped)
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_crawler(cls, crawler):
|
def from_crawler(cls, crawler: Crawler) -> Self:
|
||||||
return cls(crawler)
|
return cls(crawler)
|
||||||
|
|
||||||
def get_virtual_size(self):
|
def get_virtual_size(self) -> int:
|
||||||
size = self.resource.getrusage(self.resource.RUSAGE_SELF).ru_maxrss
|
size: int = self.resource.getrusage(self.resource.RUSAGE_SELF).ru_maxrss
|
||||||
if sys.platform != "darwin":
|
if sys.platform != "darwin":
|
||||||
# on macOS ru_maxrss is in bytes, on Linux it is in KB
|
# on macOS ru_maxrss is in bytes, on Linux it is in KB
|
||||||
size *= 1024
|
size *= 1024
|
||||||
return size
|
return size
|
||||||
|
|
||||||
def engine_started(self):
|
def engine_started(self) -> None:
|
||||||
|
assert self.crawler.stats
|
||||||
self.crawler.stats.set_value("memusage/startup", self.get_virtual_size())
|
self.crawler.stats.set_value("memusage/startup", self.get_virtual_size())
|
||||||
self.tasks = []
|
self.tasks: List[task.LoopingCall] = []
|
||||||
tsk = task.LoopingCall(self.update)
|
tsk = task.LoopingCall(self.update)
|
||||||
self.tasks.append(tsk)
|
self.tasks.append(tsk)
|
||||||
tsk.start(self.check_interval, now=True)
|
tsk.start(self.check_interval, now=True)
|
||||||
|
|
@ -68,15 +77,18 @@ class MemoryUsage:
|
||||||
self.tasks.append(tsk)
|
self.tasks.append(tsk)
|
||||||
tsk.start(self.check_interval, now=True)
|
tsk.start(self.check_interval, now=True)
|
||||||
|
|
||||||
def engine_stopped(self):
|
def engine_stopped(self) -> None:
|
||||||
for tsk in self.tasks:
|
for tsk in self.tasks:
|
||||||
if tsk.running:
|
if tsk.running:
|
||||||
tsk.stop()
|
tsk.stop()
|
||||||
|
|
||||||
def update(self):
|
def update(self) -> None:
|
||||||
|
assert self.crawler.stats
|
||||||
self.crawler.stats.max_value("memusage/max", self.get_virtual_size())
|
self.crawler.stats.max_value("memusage/max", self.get_virtual_size())
|
||||||
|
|
||||||
def _check_limit(self):
|
def _check_limit(self) -> None:
|
||||||
|
assert self.crawler.engine
|
||||||
|
assert self.crawler.stats
|
||||||
peak_mem_usage = self.get_virtual_size()
|
peak_mem_usage = self.get_virtual_size()
|
||||||
if peak_mem_usage > self.limit:
|
if peak_mem_usage > self.limit:
|
||||||
self.crawler.stats.set_value("memusage/limit_reached", 1)
|
self.crawler.stats.set_value("memusage/limit_reached", 1)
|
||||||
|
|
@ -106,9 +118,10 @@ class MemoryUsage:
|
||||||
{"virtualsize": peak_mem_usage / 1024 / 1024},
|
{"virtualsize": peak_mem_usage / 1024 / 1024},
|
||||||
)
|
)
|
||||||
|
|
||||||
def _check_warning(self):
|
def _check_warning(self) -> None:
|
||||||
if self.warned: # warn only once
|
if self.warned: # warn only once
|
||||||
return
|
return
|
||||||
|
assert self.crawler.stats
|
||||||
if self.get_virtual_size() > self.warning:
|
if self.get_virtual_size() > self.warning:
|
||||||
self.crawler.stats.set_value("memusage/warning_reached", 1)
|
self.crawler.stats.set_value("memusage/warning_reached", 1)
|
||||||
mem = self.warning / 1024 / 1024
|
mem = self.warning / 1024 / 1024
|
||||||
|
|
@ -126,8 +139,10 @@ class MemoryUsage:
|
||||||
self.crawler.stats.set_value("memusage/warning_notified", 1)
|
self.crawler.stats.set_value("memusage/warning_notified", 1)
|
||||||
self.warned = True
|
self.warned = True
|
||||||
|
|
||||||
def _send_report(self, rcpts, subject):
|
def _send_report(self, rcpts: List[str], subject: str) -> None:
|
||||||
"""send notification mail with some additional useful info"""
|
"""send notification mail with some additional useful info"""
|
||||||
|
assert self.crawler.engine
|
||||||
|
assert self.crawler.stats
|
||||||
stats = self.crawler.stats
|
stats = self.crawler.stats
|
||||||
s = f"Memory usage at engine startup : {stats.get_value('memusage/startup') / 1024 / 1024}M\r\n"
|
s = f"Memory usage at engine startup : {stats.get_value('memusage/startup') / 1024 / 1024}M\r\n"
|
||||||
s += f"Maximum memory usage : {stats.get_value('memusage/max') / 1024 / 1024}M\r\n"
|
s += f"Maximum memory usage : {stats.get_value('memusage/max') / 1024 / 1024}M\r\n"
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,8 @@ Mail sending helpers
|
||||||
See documentation in docs/topics/email.rst
|
See documentation in docs/topics/email.rst
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
from email import encoders as Encoders
|
from email import encoders as Encoders
|
||||||
from email.mime.base import MIMEBase
|
from email.mime.base import MIMEBase
|
||||||
|
|
@ -12,14 +14,20 @@ from email.mime.nonmultipart import MIMENonMultipart
|
||||||
from email.mime.text import MIMEText
|
from email.mime.text import MIMEText
|
||||||
from email.utils import formatdate
|
from email.utils import formatdate
|
||||||
from io import BytesIO
|
from io import BytesIO
|
||||||
|
from typing import TYPE_CHECKING
|
||||||
|
|
||||||
from twisted import version as twisted_version
|
from twisted import version as twisted_version
|
||||||
from twisted.internet import defer, ssl
|
from twisted.internet import defer, ssl
|
||||||
from twisted.python.versions import Version
|
from twisted.python.versions import Version
|
||||||
|
|
||||||
|
from scrapy.settings import BaseSettings
|
||||||
from scrapy.utils.misc import arg_to_iter
|
from scrapy.utils.misc import arg_to_iter
|
||||||
from scrapy.utils.python import to_bytes
|
from scrapy.utils.python import to_bytes
|
||||||
|
|
||||||
|
if TYPE_CHECKING:
|
||||||
|
# typing.Self requires Python 3.11
|
||||||
|
from typing_extensions import Self
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -56,7 +64,7 @@ class MailSender:
|
||||||
self.debug = debug
|
self.debug = debug
|
||||||
|
|
||||||
@classmethod
|
@classmethod
|
||||||
def from_settings(cls, settings):
|
def from_settings(cls, settings: BaseSettings) -> Self:
|
||||||
return cls(
|
return cls(
|
||||||
smtphost=settings["MAIL_HOST"],
|
smtphost=settings["MAIL_HOST"],
|
||||||
mailfrom=settings["MAIL_FROM"],
|
mailfrom=settings["MAIL_FROM"],
|
||||||
|
|
@ -203,7 +211,7 @@ class MailSender:
|
||||||
to_addrs,
|
to_addrs,
|
||||||
msg,
|
msg,
|
||||||
d,
|
d,
|
||||||
**factory_keywords
|
**factory_keywords,
|
||||||
)
|
)
|
||||||
factory.noisy = False
|
factory.noisy = False
|
||||||
return factory
|
return factory
|
||||||
|
|
|
||||||
|
|
@ -1,14 +1,15 @@
|
||||||
"""Some debugging functions for working with the Scrapy engine"""
|
"""Some debugging functions for working with the Scrapy engine"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
# used in global tests code
|
# used in global tests code
|
||||||
from time import time # noqa: F401
|
from time import time # noqa: F401
|
||||||
from typing import TYPE_CHECKING, Any, List, Tuple
|
from typing import Any, List, Tuple
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
from scrapy.core.engine import ExecutionEngine
|
||||||
from scrapy.core.engine import ExecutionEngine
|
|
||||||
|
|
||||||
|
|
||||||
def get_engine_status(engine: "ExecutionEngine") -> List[Tuple[str, Any]]:
|
def get_engine_status(engine: ExecutionEngine) -> List[Tuple[str, Any]]:
|
||||||
"""Return a report of the current engine status"""
|
"""Return a report of the current engine status"""
|
||||||
tests = [
|
tests = [
|
||||||
"time()-engine.start_time",
|
"time()-engine.start_time",
|
||||||
|
|
@ -37,7 +38,7 @@ def get_engine_status(engine: "ExecutionEngine") -> List[Tuple[str, Any]]:
|
||||||
return checks
|
return checks
|
||||||
|
|
||||||
|
|
||||||
def format_engine_status(engine: "ExecutionEngine") -> str:
|
def format_engine_status(engine: ExecutionEngine) -> str:
|
||||||
checks = get_engine_status(engine)
|
checks = get_engine_status(engine)
|
||||||
s = "Execution engine status\n\n"
|
s = "Execution engine status\n\n"
|
||||||
for test, result in checks:
|
for test, result in checks:
|
||||||
|
|
@ -47,5 +48,5 @@ def format_engine_status(engine: "ExecutionEngine") -> str:
|
||||||
return s
|
return s
|
||||||
|
|
||||||
|
|
||||||
def print_engine_status(engine: "ExecutionEngine") -> None:
|
def print_engine_status(engine: ExecutionEngine) -> None:
|
||||||
print(format_engine_status(engine))
|
print(format_engine_status(engine))
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue