This commit is contained in:
2026-07-27 10:34:53 +08:00
commit 3527975794
74 changed files with 16405 additions and 0 deletions
View File
+185
View File
@@ -0,0 +1,185 @@
"""浏览器兜底:纯 HTTP 被 Akamai 拦截时,用 Playwright 取回一套可用 cookie
为什么只是「兜底」而不是主路径:乐天的反爬是 Akamai Bot Manager,正常携带
完整浏览器请求头 + 复用其下发的 cookie 就能通过,不像骏河屋的 Cloudflare
Turnstile 必须真浏览器交互。因此这里的浏览器是冷路径——按需惰性启动,
长时间不用会被回收,未安装 playwright 时整体降级为「不兜底」而非报错。
导航结果里的 HTML 会一并返回:既然浏览器已经把页面取到了,就没必要让调用方
再发一次 HTTP 请求。
"""
from __future__ import annotations
import asyncio
import logging
from dataclasses import dataclass, field
from typing import Any
from app.core.config import Settings
from app.core import site
logger = logging.getLogger(__name__)
@dataclass(slots=True)
class BrowserVisit:
"""一次浏览器导航的产出:cookie 快照 + 页面 HTML"""
cookies: list[dict[str, Any]] = field(default_factory=list)
html: str = ""
class BrowserFallback:
"""按需启动的 Playwright 实例,用于取回通过反爬校验的 cookie
仅在 HTTP 抓取被判定为拦截后才会被调用;调用之间浏览器保持存活,
由 close() 在应用关闭时释放。
"""
def __init__(self, settings: Settings):
self._settings = settings
self._playwright: Any | None = None
self._browser: Any | None = None
self._lock = asyncio.Lock()
self._unavailable_reason: str | None = None
@property
def enabled(self) -> bool:
"""配置是否允许使用浏览器兜底"""
return self._settings.browser_fallback_enabled
@property
def unavailable_reason(self) -> str | None:
"""浏览器不可用的原因(未启用 / 未安装 / 启动失败)"""
if not self.enabled:
return "browser fallback is disabled"
return self._unavailable_reason
@property
def ready(self) -> bool:
"""浏览器当前是否已启动且连接正常"""
browser = self._browser
if browser is None:
return False
is_connected = getattr(browser, "is_connected", None)
if callable(is_connected):
try:
return bool(is_connected())
except Exception:
return False
return True
async def _ensure_browser(self) -> Any | None:
"""惰性启动浏览器;不可用时返回 None 并记录原因,不抛异常。"""
if not self.enabled:
return None
if self.ready:
return self._browser
async with self._lock:
if self.ready:
return self._browser
# 上一轮已经启动失败过,直接沿用结论,避免每次请求都吃一次启动超时
if self._unavailable_reason is not None and self._playwright is None:
return None
try:
from playwright.async_api import async_playwright
except ImportError as exc:
self._unavailable_reason = (
f"playwright is not installed: {exc}; "
'install it with pip install ".[browser]" && python -m playwright install chromium'
)
logger.warning("浏览器兜底不可用:%s", self._unavailable_reason)
return None
await self._close_locked()
try:
self._playwright = await async_playwright().start()
self._browser = await self._playwright.chromium.launch(
headless=self._settings.browser_headless_effective,
channel=self._settings.browser_channel or None,
proxy=self._settings.playwright_proxy,
timeout=self._settings.browser_launch_timeout_seconds * 1000,
args=["--no-first-run", "--disable-blink-features=AutomationControlled"],
)
self._unavailable_reason = None
logger.info(
"浏览器兜底已启动:headless=%s channel=%s proxy=%s",
self._settings.browser_headless_effective,
self._settings.browser_channel,
bool(self._settings.proxy_server),
)
except Exception as exc:
self._unavailable_reason = str(exc)
logger.exception("浏览器兜底启动失败:%s", self._unavailable_reason)
await self._close_locked()
return None
return self._browser
async def visit(self, url: str, *, mobile: bool) -> BrowserVisit | None:
"""用浏览器访问 url,返回 cookie 与页面 HTML;不可用时返回 None。
UA 与请求头必须和后续 HTTP 客户端使用的那一套保持一致——Akamai 会把
cookie 与指纹绑定,用 PC 指纹拿到的 cookie 拿去发手机 UA 请求可能失效。
"""
browser = await self._ensure_browser()
if browser is None:
return None
headers = site.default_headers(mobile=mobile)
context = None
try:
context = await browser.new_context(
user_agent=headers["User-Agent"],
locale="ja-JP",
timezone_id="Asia/Tokyo",
viewport={"width": 390, "height": 844} if mobile else {"width": 1440, "height": 900},
is_mobile=mobile,
has_touch=mobile,
extra_http_headers={"Accept-Language": site.ACCEPT_LANGUAGE},
)
page = await context.new_page()
await page.goto(
url,
wait_until="domcontentloaded",
timeout=self._settings.browser_nav_timeout_seconds * 1000,
)
html = await page.content()
cookies = await context.cookies()
logger.info("浏览器兜底导航完成:url=%s cookies=%s html_len=%s", url, len(cookies), len(html))
return BrowserVisit(cookies=cookies, html=html)
except Exception as exc:
logger.warning("浏览器兜底导航失败:url=%s err=%s", url, exc)
# 导航失败常常意味着浏览器已掉线,标记后由下次调用重建
if not self.ready:
self._unavailable_reason = str(exc)
return None
finally:
if context is not None:
try:
await context.close()
except Exception:
logger.debug("关闭兜底浏览器上下文失败", exc_info=True)
async def _close_locked(self) -> None:
"""释放浏览器与 Playwright 运行时(调用方需已持锁或处于关闭流程)"""
if self._browser is not None:
try:
await self._browser.close()
except Exception:
logger.debug("关闭兜底浏览器失败", exc_info=True)
self._browser = None
if self._playwright is not None:
try:
await self._playwright.stop()
except Exception:
logger.debug("停止兜底 Playwright 运行时失败", exc_info=True)
self._playwright = None
async def close(self) -> None:
"""应用关闭时释放浏览器资源"""
async with self._lock:
await self._close_locked()
+130
View File
@@ -0,0 +1,130 @@
"""ラクマ 抓取客户端:把请求参数翻译成站点 URL,抓取后解析为结构化数据
四个接口都走同一条 HTTP 通道(ラクマ 无 Akamai 限速,不需要分指纹通道):
搜索、商品详情、卖家详情、卖家商品列表。
"""
from __future__ import annotations
import asyncio
import logging
from app.core.config import Settings
from app.models.scrape import (
RakumaItemDetailData,
RakumaItemDetailRequest,
RakumaSearchRequest,
RakumaSearchResultData,
RakumaShopDetailData,
RakumaShopDetailRequest,
RakumaShopItemsData,
RakumaShopItemsRequest,
)
from app.parsers.rakuma.item import parse_item_detail
from app.parsers.rakuma.search import parse_search
from app.parsers.rakuma.shop import parse_shop_detail, parse_shop_items
from app.services.rakuma_session import RakumaSession
from app.utils.rakuma_urls import (
build_item_url,
build_search_url,
build_shop_url,
normalize_search_url,
split_item_url,
split_shop_url,
)
logger = logging.getLogger(__name__)
class RakumaClient:
"""ラクマ(fril.jp)抓取客户端"""
def __init__(self, settings: Settings, session: RakumaSession):
self._settings = settings
self._session = session
async def search(self, payload: RakumaSearchRequest) -> RakumaSearchResultData:
"""抓取搜索结果列表"""
if payload.search_url is not None:
url = normalize_search_url(str(payload.search_url), payload.page)
else:
url = build_search_url(payload)
logger.info("抓取 ラクマ 搜索页:url=%s", url)
html = await self._session.fetch_html(url)
result = parse_search(
html,
request_url=url,
page=payload.page,
keyword=payload.keyword.strip(),
)
logger.info(
"ラクマ 搜索完成:url=%s items=%s total=%s",
url, len(result.items), result.total_count,
)
return result
async def item_detail(self, payload: RakumaItemDetailRequest) -> RakumaItemDetailData:
"""抓取商品详情"""
if payload.item_url is not None:
item_id = split_item_url(str(payload.item_url))
else:
item_id = (payload.item_id or "").strip()
url = build_item_url(item_id)
logger.info("抓取 ラクマ 商品详情:url=%s", url)
html = await self._session.fetch_html(url)
detail = parse_item_detail(html, item_id=item_id, item_url=url)
logger.info(
"ラクマ 详情完成:url=%s name=%s price=%s sold_out=%s",
url, detail.item_name[:40], detail.price, detail.is_sold_out,
)
return detail
async def shop_detail(self, payload: RakumaShopDetailRequest) -> RakumaShopDetailData:
"""抓取卖家详情
评价明细在单独的 /review 页上,仅在 include_reviews=true 时并发多取一次。
"""
shop_id = self._resolve_shop_id(payload.shop_url, payload.shop_id)
url = build_shop_url(shop_id)
logger.info("抓取 ラクマ 卖家详情:url=%s reviews=%s", url, payload.include_reviews)
if payload.include_reviews:
html, review_html = await asyncio.gather(
self._session.fetch_html(url),
self._session.fetch_html(build_shop_url(shop_id, review=True)),
)
else:
html, review_html = await self._session.fetch_html(url), None
detail = parse_shop_detail(
html, shop_id=shop_id, shop_url=url, review_html=review_html
)
logger.info(
"ラクマ 卖家详情完成:shop_id=%s name=%s items=%s reviews=%s",
shop_id, detail.shop_name, detail.item_count, detail.review_count,
)
return detail
async def shop_items(self, payload: RakumaShopItemsRequest) -> RakumaShopItemsData:
"""抓取卖家的商品列表"""
shop_id = self._resolve_shop_id(payload.shop_url, payload.shop_id)
url = build_shop_url(shop_id, page=payload.page)
logger.info("抓取 ラクマ 卖家商品:url=%s", url)
html = await self._session.fetch_html(url)
result = parse_shop_items(
html, shop_id=shop_id, request_url=url, page=payload.page
)
logger.info(
"ラクマ 卖家商品完成:shop_id=%s items=%s total=%s",
shop_id, len(result.items), result.total_count,
)
return result
@staticmethod
def _resolve_shop_id(shop_url: object, shop_id: str | None) -> str:
"""从 shop_url 或 shop_id 里取出店铺 hash(模型校验已保证两者不同时为空)"""
if shop_url is not None:
return split_shop_url(str(shop_url))
return (shop_id or "").strip()
+135
View File
@@ -0,0 +1,135 @@
"""ラクマ(fril.jp)站点会话:单通道 HTTP 抓取
与乐天市场那条链路(app/services/site_session.py)分开维护,因为两站的
抓取前提完全不同:
- 乐天前置 Akamai Bot Manager,无 cookie 时每个响应被拖到 ~11s,必须先访问
首页预热再复用 cookie;且搜索页与详情页要用不同 UA,需要两条指纹通道。
- ラクマ 实测没有这类限速:冷请求(无 cookie、无预热)即 0.5-1.1s,与预热后
持平;PC UA 在搜索页、详情页、店铺页上都能拿到完整模板。
因此这里只有一条通道、不做预热,也不接浏览器兜底——没有需要兜底的拦截行为。
若将来站点加了防护,再按乐天那套补预热与升级链路。
"""
from __future__ import annotations
import asyncio
import logging
from typing import Any
import httpx
from app.core import rakuma_site as site
from app.core.config import Settings
from app.core.errors import (
ItemNotFoundError,
ResourceBusyError,
UpstreamBlockedError,
UpstreamRequestError,
)
logger = logging.getLogger(__name__)
class RakumaSession:
"""ラクマ 站点抓取会话,管理并发限流与失败重试"""
def __init__(self, settings: Settings):
self._settings = settings
self._semaphore = asyncio.Semaphore(settings.max_site_concurrency)
self._client: httpx.AsyncClient | None = None
# ---- 生命周期 ----
async def start(self) -> None:
"""创建 HTTP 客户端"""
if self._client is not None:
return
self._client = httpx.AsyncClient(
headers=site.default_headers(),
timeout=self._settings.request_timeout_seconds,
follow_redirects=True,
proxy=self._settings.httpx_proxy,
http2=True,
)
logger.info(
"ラクマ 会话已就绪:concurrency=%s proxy=%s",
self._settings.max_site_concurrency,
bool(self._settings.proxy_server),
)
async def close(self) -> None:
"""关闭 HTTP 客户端"""
if self._client is None:
return
try:
await self._client.aclose()
except Exception:
logger.debug("关闭 ラクマ HTTP 客户端失败", exc_info=True)
self._client = None
# ---- 状态 ----
def status(self) -> dict[str, Any]:
"""会话状态,供健康检查展示"""
return {"ready": self._client is not None}
# ---- 抓取 ----
async def fetch_html(self, url: str) -> str:
"""抓取页面 HTML
Raises:
ItemNotFoundError: 目标页面 404(商品已下架或 ID 不存在)
UpstreamRequestError: 网络异常或上游 5xx
UpstreamBlockedError: 反复取不到正常页面
ResourceBusyError: 等待并发槽位超时
"""
client = self._client
if client is None:
raise UpstreamRequestError("ラクマ 会话尚未初始化")
max_attempts = max(1, self._settings.http_max_attempts)
last_error = "unknown error"
try:
await asyncio.wait_for(
self._semaphore.acquire(),
timeout=self._settings.request_timeout_seconds,
)
except TimeoutError as exc:
raise ResourceBusyError() from exc
try:
for attempt in range(1, max_attempts + 1):
try:
response = await client.get(url)
except httpx.HTTPError as exc:
last_error = f"{type(exc).__name__}: {exc}"
logger.warning(
"ラクマ 抓取请求异常:url=%s attempt=%s/%s err=%s",
url, attempt, max_attempts, last_error,
)
continue
if response.status_code == 404:
raise ItemNotFoundError(f"Page not found: {url}")
if response.status_code < 400:
return response.text
last_error = (
f"upstream status {response.status_code}"
if response.status_code >= 500
else f"status {response.status_code}"
)
logger.warning(
"ラクマ 抓取结果异常:url=%s attempt=%s/%s %s",
url, attempt, max_attempts, last_error,
)
if last_error.startswith("upstream status") or ":" in last_error:
raise UpstreamRequestError(f"Upstream request failed: {last_error}")
raise UpstreamBlockedError(f"Failed to fetch {url}: {last_error}")
finally:
self._semaphore.release()
+169
View File
@@ -0,0 +1,169 @@
"""乐天抓取客户端:把请求参数翻译成站点 URL,抓取后解析为结构化数据
两个接口分别走不同的指纹通道:
- 搜索:PC 通道(search.rakuten.co.jp 的搜索页)
- 详情:手机通道(item.rakuten.co.jp 只有手机 UA 才返回统一模板)
"""
from __future__ import annotations
import logging
from app.core import site
from app.core.config import Settings
from app.core.errors import ScrapeParseError
from app.models.scrape import (
GenreData,
GenreRequest,
ItemDetailData,
ItemDetailRequest,
SearchRequest,
SearchResultData,
ShopDetailData,
ShopDetailRequest,
ShopItemsRequest,
)
from app.parsers.genre import parse_genres
from app.parsers.item import parse_item_detail
from app.parsers.search import parse_search
from app.parsers.shop import parse_shop_detail
from app.parsers.state import extract_initial_state
from app.parsers.subsites import (
SubsitePage,
build_item_page_validator,
host_of,
parse_subsite_item,
)
from app.services.site_session import SiteSession
from app.utils.urls import (
build_genre_url,
build_item_url,
build_search_url,
build_shop_url,
normalize_search_url,
split_item_url,
split_shop_url,
)
logger = logging.getLogger(__name__)
class RakutenClient:
"""乐天市场抓取客户端"""
def __init__(self, settings: Settings, session: SiteSession):
self._settings = settings
self._session = session
async def search(self, payload: SearchRequest) -> SearchResultData:
"""抓取搜索结果列表"""
if payload.search_url is not None:
url = normalize_search_url(str(payload.search_url), payload.page)
else:
url = build_search_url(payload)
logger.info("抓取搜索页:url=%s", url)
html = await self._session.fetch_html(url, mobile=False)
state = extract_initial_state(html)
result = parse_search(
state,
request_url=url,
page=payload.page,
exclude_ads=payload.exclude_ads,
)
logger.info(
"搜索完成:url=%s items=%s ads=%s total=%s",
url, len(result.items), result.ad_count, result.total_count,
)
return result
async def genres(self, payload: GenreRequest) -> GenreData:
"""抓取分类树:顶层分类列表,或指定分类的信息与直接子分类"""
genre_id = (payload.genre_id or "").strip() or None
url = build_genre_url(genre_id)
logger.info("抓取分类树:url=%s genre_id=%s", url, genre_id)
html = await self._session.fetch_html(url, mobile=False)
state = extract_initial_state(html)
result = parse_genres(state, genre_id=genre_id)
logger.info(
"分类树完成:genre_id=%s name=%s children=%s",
result.genre_id or "root", result.name, len(result.children),
)
return result
async def shop_detail(self, payload: ShopDetailRequest) -> ShopDetailData:
"""抓取商家详情"""
if payload.shop_url is not None:
shop_code = split_shop_url(str(payload.shop_url))
else:
shop_code = (payload.shop_code or "").strip()
url = build_shop_url(shop_code)
logger.info("抓取店铺页:url=%s", url)
html = await self._session.fetch_html(url, mobile=False)
state = extract_initial_state(html)
result = parse_shop_detail(state, shop_code=shop_code, html=html)
logger.info(
"店铺详情完成:shop_code=%s shop_id=%s name=%s reviews=%s",
result.shop_code, result.shop_id, result.shop_name, result.review_count,
)
return result
async def shop_items(self, payload: ShopItemsRequest) -> SearchResultData:
"""抓取商家名下的商品列表
站点没有可分页抓取的「店铺内商品」页,店铺商品实际就是搜索结果按
`sid=` 限定店铺后的产物,因此这里转成一次搜索。只给 shop_code 时
需要先取一次店铺详情换出 shop_id。
"""
shop_id = payload.shop_id
if shop_id is None:
detail = await self.shop_detail(ShopDetailRequest(shop_code=payload.shop_code))
if detail.shop_id is None:
raise ScrapeParseError(f"未能取得店铺 {payload.shop_code} 的 shop_id")
shop_id = detail.shop_id
return await self.search(payload.to_search_request(shop_id))
async def item_detail(self, payload: ItemDetailRequest) -> ItemDetailData:
"""抓取商品详情"""
if payload.item_url is not None:
shop_code, item_code = split_item_url(str(payload.item_url))
else:
# 模型校验已保证此分支下两者都非空
shop_code, item_code = payload.shop_code or "", payload.item_code or ""
# 统一收敛到规范地址,去掉 variantId 等查询参数——所有 SKU 组合都会随详情返回
url = build_item_url(shop_code, item_code)
logger.info("抓取商品详情:url=%s", url)
# 部分官方旗舰店的商品页会跳出市场域名,落到各自的独立站点;
# 校验与解析都按最终落地域名分派,落到未登记站点时由校验器抛错。
page = await self._session.fetch(
url, mobile=True, validator=build_item_page_validator(url)
)
if host_of(page.url) == site.ITEM_HOST:
detail = parse_item_detail(
extract_initial_state(page.html),
item_url=url,
shop_code=shop_code,
include_sku_variants=payload.include_sku_variants,
)
else:
detail = parse_subsite_item(
SubsitePage(
html=page.html,
requested_url=url,
final_url=page.url,
shop_code=shop_code,
item_code=item_code,
include_sku_variants=payload.include_sku_variants,
)
)
logger.info(
"详情完成:url=%s source=%s name=%s price=%s skus=%s",
url, detail.source, detail.item_name[:40], detail.price, detail.sku.variant_count,
)
return detail
+310
View File
@@ -0,0 +1,310 @@
"""站点会话:带 Akamai cookie 复用的 HTTP 抓取通道,附带浏览器兜底
乐天前置 Akamai Bot Manager,行为特征(实测):
- 请求头不完整、无 cookie 时不封禁,而是把每个响应拖到 ~11s(与响应体大小无关)
- 补齐浏览器导航请求头并复用 Akamai 下发的 cookie 后,稳定在 ~0.6-0.9s
因此这里按「指纹画像」维护两条独立通道:搜索页用 PC 画像,商品详情页用手机
画像(详情页只有手机 UA 才返回带 __INITIAL_STATE__ 的统一模板)。两条通道
各自持有独立 cookie 罐,避免把 PC 指纹拿到的 cookie 混用到手机请求上。
抓取失败时的升级路径:重新预热 → 浏览器兜底取 cookie → 放弃。
"""
from __future__ import annotations
import asyncio
import logging
import time
from dataclasses import dataclass, field
from typing import Any
import httpx
from app.core import site
from app.core.config import Settings
from app.core.errors import (
ItemNotFoundError,
ResourceBusyError,
UpstreamBlockedError,
UpstreamRequestError,
)
from app.parsers.state import PageValidator, require_state_marker
from app.services.browser_fallback import BrowserFallback
logger = logging.getLogger(__name__)
@dataclass(slots=True)
class FetchedPage:
"""一次成功抓取的产物"""
html: str
url: str # 最终落地 URL;发生跳转时与请求地址不同
# Akamai 拦截页/挑战页的特征串
_BLOCK_MARKERS = (
"access denied",
"pardon our interruption",
"reference #",
"errors.edgesuite.net",
"/_sec/cp_challenge/",
)
@dataclass(slots=True)
class _Profile:
"""一条指纹通道:独立的 httpx 客户端、cookie 罐与预热状态"""
name: str
mobile: bool
client: httpx.AsyncClient
lock: asyncio.Lock = field(default_factory=asyncio.Lock)
warmed_at: float = 0.0
@property
def cookie_names(self) -> set[str]:
return {cookie.name for cookie in self.client.cookies.jar}
class SiteSession:
"""乐天站点抓取会话,管理 cookie 预热、并发限流与失败升级"""
def __init__(self, settings: Settings, browser_fallback: BrowserFallback):
self._settings = settings
self._browser = browser_fallback
self._semaphore = asyncio.Semaphore(settings.max_site_concurrency)
self._profiles: dict[str, _Profile] = {}
# ---- 生命周期 ----
async def start(self) -> None:
"""创建两条指纹通道的 HTTP 客户端"""
for name, mobile in (("pc", False), ("sp", True)):
if name in self._profiles:
continue
self._profiles[name] = _Profile(
name=name,
mobile=mobile,
client=httpx.AsyncClient(
headers=site.default_headers(mobile=mobile),
timeout=self._settings.request_timeout_seconds,
follow_redirects=True,
proxy=self._settings.httpx_proxy,
http2=True,
),
)
logger.info(
"站点会话已就绪:profiles=%s concurrency=%s proxy=%s",
list(self._profiles),
self._settings.max_site_concurrency,
bool(self._settings.proxy_server),
)
async def close(self) -> None:
"""关闭所有 HTTP 客户端"""
for profile in self._profiles.values():
try:
await profile.client.aclose()
except Exception:
logger.debug("关闭 HTTP 客户端失败:profile=%s", profile.name, exc_info=True)
self._profiles.clear()
# ---- 状态 ----
def profile_status(self) -> dict[str, dict[str, Any]]:
"""各通道的预热状态,供健康检查展示"""
now = time.monotonic()
return {
name: {
"warmed": profile.warmed_at > 0,
"age_seconds": round(now - profile.warmed_at, 1) if profile.warmed_at else None,
"cookies": sorted(profile.cookie_names & set(site.AKAMAI_COOKIE_NAMES)),
}
for name, profile in self._profiles.items()
}
# ---- 抓取 ----
async def fetch_html(self, url: str, *, mobile: bool) -> str:
"""抓取要求含 __INITIAL_STATE__ 的页面,只取 HTML"""
page = await self.fetch(url, mobile=mobile)
return page.html
async def fetch(
self,
url: str,
*,
mobile: bool,
validator: PageValidator | None = None,
) -> FetchedPage:
"""抓取页面,返回 HTML 与最终落地地址
失败时按 重新预热 → 浏览器兜底 的顺序逐级升级重试。
Args:
validator: 页面校验器,默认要求页面含 __INITIAL_STATE__。跨站抓取时
由调用方注入按落地域名分派的校验逻辑。
Raises:
ItemNotFoundError: 目标页面 404
UpstreamBlockedError: 反复被反爬阻断
UpstreamRequestError: 网络异常或上游 5xx
ResourceBusyError: 等待并发槽位超时
AppError: 校验器判定为不可重试的失败(如落到不支持的站点)
"""
profile = self._profiles["sp" if mobile else "pc"]
max_attempts = max(1, self._settings.http_max_attempts)
validate = validator or require_state_marker
last_error: str = "unknown error"
try:
await asyncio.wait_for(
self._semaphore.acquire(),
timeout=self._settings.request_timeout_seconds,
)
except TimeoutError as exc:
raise ResourceBusyError() from exc
try:
for attempt in range(1, max_attempts + 1):
await self._ensure_warm(profile)
try:
response = await profile.client.get(url)
except httpx.HTTPError as exc:
last_error = f"{type(exc).__name__}: {exc}"
logger.warning(
"抓取请求异常:url=%s profile=%s attempt=%s/%s err=%s",
url, profile.name, attempt, max_attempts, last_error,
)
continue
if response.status_code == 404:
raise ItemNotFoundError(f"Page not found: {url}")
text = response.text
final_url = str(response.url)
if response.status_code < 400:
# 校验器可能直接抛 AppError 表示「重试也没用」,此处不拦截
reason = validate(text, final_url)
if reason is None:
return FetchedPage(html=text, url=final_url)
last_error = self._refine_failure(reason, text)
else:
last_error = self._describe_error_status(response.status_code)
logger.warning(
"抓取结果异常:url=%s profile=%s attempt=%s/%s %s",
url, profile.name, attempt, max_attempts, last_error,
)
if attempt >= max_attempts:
break
# 第一次失败先便宜地换一套 cookie;仍失败才动用浏览器
if attempt == 1:
await self._invalidate(profile)
else:
page = await self._escalate_to_browser(profile, url, validate)
if page is not None:
return page
if self._is_server_error(last_error):
raise UpstreamRequestError(f"Upstream request failed: {last_error}")
raise UpstreamBlockedError(f"Blocked while fetching {url}: {last_error}")
finally:
self._semaphore.release()
# ---- 内部 ----
@staticmethod
def _describe_error_status(status_code: int) -> str:
return (
f"upstream status {status_code}"
if status_code >= 500
else f"status {status_code}"
)
@staticmethod
def _refine_failure(reason: str, text: str) -> str:
"""页面内容校验失败时,优先报出更具体的反爬挑战页特征"""
lowered = text[:4000].lower()
matched = next((marker for marker in _BLOCK_MARKERS if marker in lowered), None)
return f"challenge page detected ({matched})" if matched else reason
@staticmethod
def _is_server_error(reason: str) -> bool:
return reason.startswith("upstream status") or reason.startswith("httpx") or "Error:" in reason
async def _ensure_warm(self, profile: _Profile) -> None:
"""确保通道持有新鲜的 Akamai cookie;过期或缺失时访问首页预热"""
if self._is_warm(profile):
return
async with profile.lock:
if self._is_warm(profile):
return
try:
response = await profile.client.get(self._settings.home_url)
profile.warmed_at = time.monotonic()
logger.info(
"会话预热完成:profile=%s status=%s cookies=%s",
profile.name,
response.status_code,
sorted(profile.cookie_names & set(site.AKAMAI_COOKIE_NAMES)),
)
except httpx.HTTPError as exc:
# 预热失败不阻断本次抓取:直连目标页仍可能成功,只是慢
logger.warning("会话预热失败:profile=%s err=%s", profile.name, exc)
profile.warmed_at = time.monotonic()
def _is_warm(self, profile: _Profile) -> bool:
if not profile.warmed_at:
return False
if time.monotonic() - profile.warmed_at > self._settings.session_ttl_seconds:
return False
return bool(profile.cookie_names & set(site.AKAMAI_COOKIE_NAMES))
async def _invalidate(self, profile: _Profile) -> None:
"""清空通道 cookie 并强制下次重新预热"""
async with profile.lock:
profile.client.cookies.clear()
profile.warmed_at = 0.0
logger.info("已清空会话 cookie,将重新预热:profile=%s", profile.name)
async def _escalate_to_browser(
self, profile: _Profile, url: str, validate: PageValidator
) -> FetchedPage | None:
"""用浏览器访问目标页,把 cookie 回灌给 HTTP 客户端
浏览器已经拿到合格页面时直接返回,省掉一次重复请求。
"""
visit = await self._browser.visit(url, mobile=profile.mobile)
if visit is None:
logger.warning(
"浏览器兜底不可用,放弃升级:profile=%s reason=%s",
profile.name,
self._browser.unavailable_reason,
)
return None
async with profile.lock:
for cookie in visit.cookies:
name = cookie.get("name")
value = cookie.get("value")
if not name or value is None:
continue
profile.client.cookies.set(
name,
value,
domain=cookie.get("domain") or "",
path=cookie.get("path") or "/",
)
profile.warmed_at = time.monotonic()
# 浏览器不回报最终 URL,这里以请求地址为准;跨站跳转场景下 HTTP 通道已先行
# 报错,走不到这一步。
if validate(visit.html, url) is None:
logger.info("浏览器兜底直接取回页面:url=%s profile=%s", url, profile.name)
return FetchedPage(html=visit.html, url=url)
logger.warning("浏览器兜底页面仍未通过校验:url=%s profile=%s", url, profile.name)
return None