Test / app.py
ScrapyTheScrapper's picture
Rename app-4.py to app.py
585c469 verified
Raw
History Blame Contribute Delete
19.8 kB
"""
Scraper SOTA v3 — FastAPI + Scrapy + Pydantic v2.
Architecture :
- FastAPI gère les routes HTTP et la validation (Pydantic v2)
- Scrapy est le moteur de crawl primaire (retry, throttle, middlewares intégrés)
- curl_cffi / cloudscraper / httpx sont des fallbacks pour les sites protégés
- structlog + prometheus pour l'observabilité
- Cache LRU en mémoire (remplace le dict manuel)
"""
from __future__ import annotations
import asyncio
import logging
import time
from collections import OrderedDict
from contextlib import asynccontextmanager
from datetime import datetime, timezone
from typing import Any, Optional
from urllib.parse import urlparse
import httpx
import orjson
import structlog
from curl_cffi import requests as curl_requests
import cloudscraper
from fastapi import BackgroundTasks, FastAPI, HTTPException, Request
from fastapi.responses import ORJSONResponse, Response
from prometheus_client import (
CONTENT_TYPE_LATEST,
Counter,
Gauge,
Histogram,
generate_latest,
)
from tenacity import (
retry,
retry_if_exception_type,
stop_after_attempt,
wait_exponential,
)
from models import (
ContentData,
ExtractionConfig,
ExtractionMode,
HealthResponse,
ImagesData,
ImageItem,
LinksData,
LinkItem,
MetadataData,
PerformanceMetrics,
ScrapingMethod,
ScrapeOptions,
ScrapeRequest,
ScrapeResponse,
settings,
)
from scraper import PageItem, ScrapyRunner, scrapy_runner
from utils import ContentCleaner, URLValidator
# ---------------------------------------------------------------------------
# Logging structuré
# ---------------------------------------------------------------------------
structlog.configure(
processors=[
structlog.stdlib.filter_by_level,
structlog.processors.TimeStamper(fmt="iso"),
structlog.stdlib.add_logger_name,
structlog.stdlib.add_log_level,
structlog.processors.StackInfoRenderer(),
structlog.processors.format_exc_info,
structlog.processors.JSONRenderer(serializer=orjson.dumps),
],
wrapper_class=structlog.stdlib.BoundLogger,
logger_factory=structlog.stdlib.LoggerFactory(),
cache_logger_on_first_use=True,
)
log = structlog.get_logger()
# ---------------------------------------------------------------------------
# Métriques Prometheus
# ---------------------------------------------------------------------------
REQUESTS_TOTAL = Counter("scraper_requests_total", "Total requests", ["method", "status"])
REQUESTS_DURATION = Histogram("scraper_duration_seconds", "Request duration", ["method"])
ACTIVE_REQUESTS = Gauge("scraper_active_requests", "Active requests")
CACHE_HITS = Counter("scraper_cache_hits_total", "Cache hits")
CACHE_MISSES = Counter("scraper_cache_misses_total", "Cache misses")
ERRORS_TOTAL = Counter("scraper_errors_total", "Errors", ["error_type"])
# ---------------------------------------------------------------------------
# Cache LRU en mémoire
# ---------------------------------------------------------------------------
class LRUCache:
"""Cache LRU thread-safe (asyncio) avec TTL par entrée."""
def __init__(self, max_size: int, default_ttl: int) -> None:
self._store: OrderedDict[str, tuple[Any, float]] = OrderedDict()
self.max_size = max_size
self.default_ttl = default_ttl
self._hits = 0
self._misses = 0
def get(self, key: str) -> Optional[Any]:
if key not in self._store:
self._misses += 1
CACHE_MISSES.inc()
return None
value, expires_at = self._store[key]
if time.monotonic() > expires_at:
del self._store[key]
self._misses += 1
CACHE_MISSES.inc()
return None
self._store.move_to_end(key)
self._hits += 1
CACHE_HITS.inc()
return value
def set(self, key: str, value: Any, ttl: Optional[int] = None) -> None:
effective_ttl = ttl if ttl is not None else self.default_ttl
if key in self._store:
self._store.move_to_end(key)
elif len(self._store) >= self.max_size:
self._store.popitem(last=False) # évicte le plus ancien
self._store[key] = (value, time.monotonic() + effective_ttl)
@property
def size(self) -> int:
return len(self._store)
@property
def hit_rate(self) -> float:
total = self._hits + self._misses
return self._hits / total if total else 0.0
# ---------------------------------------------------------------------------
# État global du worker
# ---------------------------------------------------------------------------
class WorkerState:
def __init__(self) -> None:
self.start_time = time.monotonic()
self.total_requests = 0
self.active_requests = 0
self.total_errors = 0
self.cache = LRUCache(
max_size=settings.cache_max_size,
default_ttl=settings.cache_ttl,
)
self._response_times: list[float] = []
self._rt_max = 1000 # fenêtre glissante
# HTTP clients alternatifs (fallback)
self._cloudscraper = cloudscraper.create_scraper(
browser={"browser": "chrome", "platform": "windows", "mobile": False},
delay=10,
)
self._httpx_client: Optional[httpx.AsyncClient] = None
async def get_httpx_client(self) -> httpx.AsyncClient:
if self._httpx_client is None:
limits = httpx.Limits(
max_connections=settings.max_concurrent_requests,
max_keepalive_connections=settings.max_concurrent_requests // 2,
keepalive_expiry=30,
)
self._httpx_client = httpx.AsyncClient(
timeout=httpx.Timeout(settings.request_timeout),
limits=limits,
follow_redirects=settings.follow_redirects,
http2=True,
verify=settings.verify_ssl,
)
return self._httpx_client
def record_response_time(self, duration: float) -> None:
self._response_times.append(duration)
if len(self._response_times) > self._rt_max:
self._response_times.pop(0)
@property
def avg_response_time(self) -> float:
if not self._response_times:
return 0.0
return sum(self._response_times) / len(self._response_times)
@property
def error_rate(self) -> float:
if self.total_requests == 0:
return 0.0
return self.total_errors / self.total_requests
async def close(self) -> None:
if self._httpx_client:
await self._httpx_client.aclose()
await scrapy_runner.shutdown()
state = WorkerState()
# ---------------------------------------------------------------------------
# Lifecycle FastAPI
# ---------------------------------------------------------------------------
@asynccontextmanager
async def lifespan(app: FastAPI):
log.info(
"worker_startup",
worker_id=settings.worker_id,
environment=settings.environment,
)
yield
log.info("worker_shutdown", worker_id=settings.worker_id)
await state.close()
# ---------------------------------------------------------------------------
# Application
# ---------------------------------------------------------------------------
app = FastAPI(
title="Scraper SOTA v3",
description="Worker de scraping haute performance — FastAPI + Scrapy + Pydantic v2",
version="3.0.0",
lifespan=lifespan,
default_response_class=ORJSONResponse,
)
# ---------------------------------------------------------------------------
# Middleware
# ---------------------------------------------------------------------------
@app.middleware("http")
async def timing_middleware(request: Request, call_next):
t0 = time.perf_counter()
state.active_requests += 1
ACTIVE_REQUESTS.set(state.active_requests)
try:
response = await call_next(request)
elapsed = time.perf_counter() - t0
response.headers["X-Process-Time"] = f"{elapsed:.4f}"
response.headers["X-Worker-ID"] = settings.worker_id
state.record_response_time(elapsed)
return response
finally:
state.active_requests -= 1
ACTIVE_REQUESTS.set(state.active_requests)
# ---------------------------------------------------------------------------
# Routes
# ---------------------------------------------------------------------------
@app.get("/")
async def root():
return {
"service": "Scraper SOTA",
"version": "3.0.0",
"worker_id": settings.worker_id,
"status": "operational",
"engine": "Scrapy + FastAPI + Pydantic v2",
"features": [
"scrapy-primary-crawler",
"curl_cffi / cloudscraper fallback",
"pydantic-v2-strict-models",
"lru-cache",
"prometheus-metrics",
"structured-logging",
"multi-method-extraction",
],
}
@app.get("/health", response_model=HealthResponse)
async def health() -> HealthResponse:
return HealthResponse(
status="healthy" if state.error_rate < 0.3 else "degraded",
worker_id=settings.worker_id,
uptime_seconds=time.monotonic() - state.start_time,
total_requests=state.total_requests,
active_requests=state.active_requests,
cache_size=state.cache.size,
cache_hit_rate=state.cache.hit_rate,
avg_response_time=state.avg_response_time,
error_rate=state.error_rate,
)
@app.get("/metrics")
async def metrics():
return Response(content=generate_latest(), media_type=CONTENT_TYPE_LATEST)
@app.post("/scrape", response_model=ScrapeResponse)
async def scrape_url(
request: ScrapeRequest,
background_tasks: BackgroundTasks,
) -> ScrapeResponse:
t0 = time.perf_counter()
url_str = str(request.url)
state.total_requests += 1
log.info("scrape_start", url=url_str, method=request.options.method.value)
# --- Cache ---
if request.options.cache.enabled and not request.options.cache.force_refresh:
cache_key = ContentCleaner.compute_content_hash(url_str)
cached = state.cache.get(cache_key)
if cached is not None:
log.info("cache_hit", url=url_str)
cached["performance"]["cache_hit"] = True
cached["performance"]["total_time"] = time.perf_counter() - t0
return ScrapeResponse(**cached)
try:
# --- Sélection de méthode ---
method = _select_method(request)
# --- Téléchargement ---
dl_start = time.perf_counter()
html, status_code, final_url = await _download(url_str, method, request.options)
dl_time = time.perf_counter() - dl_start
REQUESTS_DURATION.labels(method=method.value).observe(dl_time)
# --- Extraction ---
parse_start = time.perf_counter()
content_data = _extract_content(html, final_url, request.options.extraction)
parse_time = time.perf_counter() - parse_start
ex_start = time.perf_counter()
metadata_data = (
_extract_metadata(html) if request.options.extraction.include_metadata else None
)
links_data = (
_extract_links(html, final_url) if request.options.extraction.include_links else None
)
images_data = (
_extract_images(html, final_url) if request.options.extraction.include_images else None
)
ex_time = time.perf_counter() - ex_start
total_time = time.perf_counter() - t0
performance = PerformanceMetrics(
total_time=round(total_time, 4),
download_time=round(dl_time, 4),
parsing_time=round(parse_time, 4),
extraction_time=round(ex_time, 4),
content_size=len(html.encode("utf-8", errors="replace")),
cache_hit=False,
scrapy_used=(method == ScrapingMethod.SCRAPY),
)
response_dict: dict[str, Any] = {
"success": True,
"worker_id": settings.worker_id,
"url": url_str,
"final_url": final_url,
"status_code": status_code,
"method_used": method,
"content": content_data.model_dump(),
"metadata": metadata_data.model_dump() if metadata_data else None,
"links": links_data.model_dump() if links_data else None,
"images": images_data.model_dump() if images_data else None,
"performance": performance.model_dump(),
"timestamp": datetime.now(timezone.utc).isoformat(),
}
# --- Mise en cache ---
if request.options.cache.enabled:
cache_key = ContentCleaner.compute_content_hash(url_str)
state.cache.set(response_dict, cache_key, request.options.cache.ttl)
REQUESTS_TOTAL.labels(method=method.value, status="success").inc()
log.info("scrape_success", url=url_str, duration=total_time, method=method.value)
return ScrapeResponse(**response_dict)
except Exception as exc:
state.total_errors += 1
ERRORS_TOTAL.labels(error_type=type(exc).__name__).inc()
REQUESTS_TOTAL.labels(method="unknown", status="error").inc()
log.error("scrape_error", url=url_str, error=str(exc))
return ScrapeResponse(
success=False,
worker_id=settings.worker_id,
url=url_str,
error=str(exc),
performance=PerformanceMetrics(
total_time=round(time.perf_counter() - t0, 4),
download_time=0.0,
parsing_time=0.0,
extraction_time=0.0,
content_size=0,
cache_hit=False,
),
)
# ---------------------------------------------------------------------------
# Logique métier
# ---------------------------------------------------------------------------
def _select_method(request: ScrapeRequest) -> ScrapingMethod:
if request.options.method != ScrapingMethod.AUTO:
return request.options.method
url_lower = str(request.url).lower()
if any(x in url_lower for x in ("cloudflare", "cf-", "captcha")):
return ScrapingMethod.CLOUDSCRAPER
# Scrapy est le moteur par défaut
return ScrapingMethod.SCRAPY
async def _download(
url: str,
method: ScrapingMethod,
options: ScrapeOptions,
) -> tuple[str, Optional[int], str]:
"""Délègue le téléchargement à Scrapy ou aux clients HTTP alternatifs."""
timeout = options.timeout or settings.request_timeout
verify = options.verify_ssl if options.verify_ssl is not None else settings.verify_ssl
headers = dict(options.headers or {})
if method == ScrapingMethod.SCRAPY:
item: PageItem = await scrapy_runner.fetch(
url,
timeout=timeout,
verify_ssl=verify,
custom_headers=headers,
)
if item.get("error"):
raise RuntimeError(f"Scrapy error: {item['error']}")
return item["html"], item.get("status_code"), item.get("final_url", url)
# Fallbacks
if "User-Agent" not in headers:
headers["User-Agent"] = settings.user_agent
headers.update(
{
"Accept": "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8",
"Accept-Language": "en-US,en;q=0.9",
"Accept-Encoding": "gzip, deflate, br",
}
)
if method == ScrapingMethod.CURL_CFFI:
loop = asyncio.get_running_loop()
resp = await loop.run_in_executor(
None,
lambda: curl_requests.get(
url,
headers=headers,
timeout=timeout,
impersonate="chrome120",
verify=verify,
allow_redirects=options.follow_redirects if options.follow_redirects is not None else settings.follow_redirects,
),
)
return resp.text, resp.status_code, str(resp.url)
if method == ScrapingMethod.CLOUDSCRAPER:
loop = asyncio.get_running_loop()
resp = await loop.run_in_executor(
None,
lambda: state._cloudscraper.get(
url,
headers=headers,
timeout=timeout,
verify=verify,
),
)
resp.raise_for_status()
return resp.text, resp.status_code, resp.url
if method == ScrapingMethod.HTTPX:
client = await state.get_httpx_client()
resp = await client.get(url, headers=headers, timeout=timeout)
resp.raise_for_status()
return resp.text, resp.status_code, str(resp.url)
raise ValueError(f"Méthode non supportée : {method}")
def _extract_content(html: str, url: str, config: ExtractionConfig) -> ContentData:
if config.mode == ExtractionMode.RAW:
return ContentData(raw_html=html[: settings.max_content_size])
raw_extracted: dict[str, Any] = {}
if config.mode in (ExtractionMode.CLEAN, ExtractionMode.FULL):
clean_html = ContentCleaner.clean_html_fast(html)
if config.mode in (ExtractionMode.MAIN_CONTENT, ExtractionMode.FULL):
raw_extracted = ContentCleaner.extract_main_content(html, url)
text = raw_extracted.get("text", "")
if config.normalize_text:
text = ContentCleaner.normalize_text(text)
return ContentData(
clean_html=(ContentCleaner.clean_html_fast(html)[: settings.max_content_size]
if config.mode == ExtractionMode.FULL else None),
text=text[: settings.max_content_size],
title=raw_extracted.get("title"),
author=raw_extracted.get("author"),
date=raw_extracted.get("date"),
description=raw_extracted.get("description"),
language=raw_extracted.get("language"),
word_count=len(text.split()),
)
# CLEAN only
return ContentData(clean_html=ContentCleaner.clean_html_fast(html)[: settings.max_content_size])
def _extract_metadata(html: str) -> MetadataData:
md = ContentCleaner.extract_metadata(html)
return MetadataData(
og_data={k.replace("og_", ""): v for k, v in md.items() if k.startswith("og_")} or None,
twitter_data={k.replace("twitter_", ""): v for k, v in md.items() if k.startswith("twitter_")} or None,
meta_tags={k: v for k, v in md.items() if not k.startswith(("og_", "twitter_"))} or None,
canonical_url=md.get("canonical"),
)
def _extract_links(html: str, base_url: str) -> LinksData:
base_domain = urlparse(base_url).netloc
all_links = ContentCleaner.extract_links(html, base_url)[: settings.max_links]
internal, external = [], []
for lk in all_links:
item = LinkItem(**lk)
if urlparse(lk["url"]).netloc == base_domain:
internal.append(item)
else:
external.append(item)
return LinksData(internal=internal, external=external)
def _extract_images(html: str, base_url: str) -> ImagesData:
imgs = ContentCleaner.extract_images(html, base_url)[: settings.max_images]
return ImagesData(images=[ImageItem(**i) for i in imgs])
# ---------------------------------------------------------------------------
# Entrypoint
# ---------------------------------------------------------------------------
if __name__ == "__main__":
import uvicorn
uvicorn.run(
"app:app",
host="0.0.0.0",
port=settings.port,
log_level="warning",
access_log=False,
loop="uvloop",
http="httptools",
limit_concurrency=settings.max_concurrent_requests,
)