Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 40 additions & 2 deletions estela_scrapy/extensions.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,12 @@
from scrapy.exporters import PythonItemExporter
from twisted.internet import task

from estela_scrapy.utils import json_serializer, producer, update_job
from estela_scrapy.utils import (
get_obj_byte_size_id_padding_bytes,
json_serializer,
producer,
update_job,
)

import logging

Expand Down Expand Up @@ -45,7 +50,6 @@
DEFAULT_MAX_DUPLICATE_TRACKING = 100_000
DEFAULT_TIMELINE_BUCKETS = 30


class BaseExtension:
def __init__(self, stats, *args, **kwargs):
self.stats = stats
Expand Down Expand Up @@ -90,6 +94,10 @@ def __init__(self, stats, schema=None, unique_field=None, max_buckets=30):
raise NotConfigured("REDIS_URL not found in the settings")
self.redis_conn = redis.from_url(redis_url)

self._item_size_exporter = PythonItemExporter(
binary=False, dont_fail=True
)

self.stats_key = os.getenv("REDIS_STATS_KEY")
self.interval = float(os.getenv("REDIS_STATS_INTERVAL"))

Expand All @@ -101,6 +109,10 @@ def __init__(self, stats, schema=None, unique_field=None, max_buckets=30):
"false").lower() == "true"
)

self._item_obj_byte_size_id_padding_bytes = (
get_obj_byte_size_id_padding_bytes()
)

# Always initialize all metrics tracking (no conditional)
self._init_metrics_tracking(schema, unique_field, max_buckets)

Expand Down Expand Up @@ -158,6 +170,10 @@ def from_crawler(cls, crawler):
max_buckets=max_buckets
)

from estela_scrapy.log import set_scrapy_stats

set_scrapy_stats(crawler.stats)

crawler.signals.connect(
ext.spider_opened, signal=signals.spider_opened)
crawler.signals.connect(
Expand All @@ -169,6 +185,9 @@ def from_crawler(cls, crawler):
def spider_opened(self, spider):
# Initialize duplicate counter
self.stats.set_value("advanced_metrics/items_duplicates", 0)
self.stats.set_value("item_obj_byte_size", 0)
self.stats.set_value("request_obj_byte_size", 0)
self.stats.set_value("log_obj_byte_size", 0)

update_job(self.job_url, self.auth_token, status=RUNNING_STATUS)
self.task = task.LoopingCall(self.store_stats, spider)
Expand All @@ -178,9 +197,28 @@ def item_scraped(self, item, spider):
# Track metrics for this item
self._track_item_metrics(item, spider)

def _item_json_payload_byte_size(self, exported_dict):
"""UTF-8 byte length of JSON (same encoding as job payloads)."""
return len(
json.dumps(exported_dict, default=json_serializer).encode("utf-8")
)

def _item_stat_byte_size(self, item):
"""JSON UTF-8 payload size (job shape) plus stored-id padding bytes.

Padding is ``ESTELA_ITEM_OBJ_BYTE_SIZE_ID_PADDING_BYTES`` (default 17).
"""
exported = dict(self._item_size_exporter.export_item(item))
return (
self._item_json_payload_byte_size(exported)
+ self._item_obj_byte_size_id_padding_bytes
)

def _track_item_metrics(self, item, spider):
"""Track timeline and duplicate metrics for each scraped item"""

self.stats.inc_value("item_obj_byte_size", self._item_stat_byte_size(item))

# Timeline tracking with cached time calculation (performance optimization)
self.items_since_time_update += 1
if self.items_since_time_update >= TIME_CACHE_UPDATE_INTERVAL:
Expand Down
21 changes: 20 additions & 1 deletion estela_scrapy/log.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import json
import logging
import os
import sys
Expand All @@ -7,16 +8,34 @@
from estela_queue_adapter import queue_noisy_libraries
from twisted.python import log as txlog

from estela_scrapy.utils import producer, to_standard_str
from estela_scrapy.utils import (
get_obj_byte_size_id_padding_bytes,
json_serializer,
producer,
to_standard_str,
)

_stderr = sys.stderr

# Set from RedisStatsCollector.from_crawler; logging may run before the crawler exists.
_scrapy_stats = None


def set_scrapy_stats(stats):
global _scrapy_stats
_scrapy_stats = stats


def _logfn(level, message, parent="none"):
data = {
"jid": os.getenv("ESTELA_SPIDER_JOB"),
"payload": {"log": str(message), "datetime": float(time.time())},
}
if _scrapy_stats is not None:
payload_byte_size = len(
json.dumps(data["payload"], default=json_serializer).encode("utf-8")
) + get_obj_byte_size_id_padding_bytes()
_scrapy_stats.inc_value("log_obj_byte_size", payload_byte_size)
producer.send("job_logs", data)


Expand Down
39 changes: 29 additions & 10 deletions estela_scrapy/middlewares.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import json
import logging
import os
import re
Expand All @@ -7,7 +8,12 @@
from scrapy.utils.request import fingerprint as request_fingerprint
from twisted.web import http

from estela_scrapy.utils import parse_time, producer
from estela_scrapy.utils import (
get_obj_byte_size_id_padding_bytes,
json_serializer,
parse_time,
producer,
)

proxy_logger = logging.getLogger("proxy_mw")

Expand All @@ -26,18 +32,31 @@ def get_status_size(response_status):


class StorageDownloaderMiddleware:
def __init__(self, stats=None):
self.stats = stats

@classmethod
def from_crawler(cls, crawler):
return cls(stats=crawler.stats)

def process_response(self, request, response, spider):
payload = {
"url": response.url,
"status": int(response.status),
"method": request.method,
"duration": int(request.meta.get("download_latency", 0) * 1000),
"time": parse_time(),
"response_size": len(response.body),
"fingerprint": request_fingerprint(request),
}
if self.stats is not None:
payload_byte_size = len(
json.dumps(payload, default=json_serializer).encode("utf-8")
) + get_obj_byte_size_id_padding_bytes()
self.stats.inc_value("request_obj_byte_size", payload_byte_size)
data = {
"jid": os.getenv("ESTELA_SPIDER_JOB"),
"payload": {
"url": response.url,
"status": int(response.status),
"method": request.method,
"duration": int(request.meta.get("download_latency", 0) * 1000),
"time": parse_time(),
"response_size": len(response.body),
"fingerprint": request_fingerprint(request),
},
"payload": payload,
}
producer.send("job_requests", data)
return response
Expand Down
19 changes: 19 additions & 0 deletions estela_scrapy/utils.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,28 @@
import json
import os
from datetime import date, datetime, timedelta

import requests
from estela_queue_adapter import get_producer_interface

# Default padding for a stored document id (e.g. BSON ObjectId-style ``_id``) in
# ``*_obj_byte_size`` stats. Override with ``ESTELA_ITEM_OBJ_BYTE_SIZE_ID_PADDING_BYTES``.
DEFAULT_OBJ_BYTE_SIZE_ID_PADDING_BYTES = 17

_obj_byte_size_id_padding_cache = None


def get_obj_byte_size_id_padding_bytes():
global _obj_byte_size_id_padding_cache
if _obj_byte_size_id_padding_cache is None:
_obj_byte_size_id_padding_cache = int(
os.getenv(
"ESTELA_ITEM_OBJ_BYTE_SIZE_ID_PADDING_BYTES",
str(DEFAULT_OBJ_BYTE_SIZE_ID_PADDING_BYTES),
)
)
return _obj_byte_size_id_padding_cache


def parse_time(date=None):
if date is None:
Expand Down