接入 OTel 链路追踪 + 镜像默认启用
抓取-解析链路原本只有日志,出现"抓回内容但解析不出预期字段"时定位慢。 接入 OpenTelemetry traces(FastAPI/httpx 自动 + 手写 fetch/parse span), 解析失败时把页面 HTML 作为 span event 上报,便于事后复现。 Dockerfile 默认开启,镜像一启动即导出到自建 OTLP endpoint。 Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -27,6 +27,7 @@ from app.scraping.services.site_session import SiteSession
|
||||
from app.shared.api import register_exception_handlers
|
||||
from app.shared.config import get_settings
|
||||
from app.shared.logging_setup import configure_logging
|
||||
from app.shared.telemetry import instrument_app, setup_telemetry, shutdown_telemetry
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -56,6 +57,9 @@ async def lifespan(app: FastAPI):
|
||||
app.state.container = container
|
||||
|
||||
configure_logging(container.settings)
|
||||
# telemetry 必须在创建任何 httpx.AsyncClient 之前完成 setup,否则 httpx
|
||||
# 自动 instrumentation 抓不到现有客户端的请求。
|
||||
setup_telemetry(container.settings, service_name="rakuten-scraping")
|
||||
logger.info("抓取服务启动:%s:%s", container.settings.app_host, container.settings.app_port)
|
||||
logger.info("日志级别:%s", container.settings.log_level)
|
||||
logger.info("当前环境:%s", container.settings.app_env)
|
||||
@@ -67,6 +71,7 @@ async def lifespan(app: FastAPI):
|
||||
await container.rakuma_session.close()
|
||||
await container.site_session.close()
|
||||
await container.browser_fallback.close()
|
||||
shutdown_telemetry()
|
||||
|
||||
|
||||
def create_app() -> FastAPI:
|
||||
@@ -76,6 +81,7 @@ def create_app() -> FastAPI:
|
||||
app.include_router(scrape_router)
|
||||
app.include_router(rakuma_router)
|
||||
register_exception_handlers(app)
|
||||
instrument_app(app)
|
||||
return app
|
||||
|
||||
|
||||
|
||||
@@ -8,6 +8,8 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
from opentelemetry import trace
|
||||
|
||||
from app.scraping.core import rakuma_site as site
|
||||
from app.shared.config import Settings
|
||||
from app.scraping.models.scrape import (
|
||||
@@ -35,8 +37,10 @@ from app.scraping.utils.rakuma_urls import (
|
||||
split_item_url,
|
||||
split_shop_url,
|
||||
)
|
||||
from app.shared.telemetry import snapshot
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
tracer = trace.get_tracer(__name__)
|
||||
|
||||
|
||||
class RakumaClient:
|
||||
@@ -54,18 +58,30 @@ class RakumaClient:
|
||||
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
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuma.search") as span:
|
||||
span.set_attribute("parse.site", "rakuma")
|
||||
span.set_attribute("parse.op", "search")
|
||||
try:
|
||||
html = await self._session.fetch_html(url)
|
||||
result = parse_search(
|
||||
html,
|
||||
request_url=url,
|
||||
page=payload.page,
|
||||
keyword=payload.keyword.strip(),
|
||||
)
|
||||
span.set_attribute("parse.items", len(result.items))
|
||||
span.set_attribute("parse.total", result.total_count)
|
||||
logger.info(
|
||||
"ラクマ 搜索完成:url=%s items=%s total=%s",
|
||||
url, len(result.items), result.total_count,
|
||||
)
|
||||
return result
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
async def categories(self, payload: RakumaCategoryRequest) -> RakumaCategoryData:
|
||||
"""抓取分类树
|
||||
@@ -76,17 +92,29 @@ class RakumaClient:
|
||||
url = site.CATEGORY_LIST_URL
|
||||
|
||||
logger.info("抓取 ラクマ 分类页:category_id=%s", payload.category_id)
|
||||
html = await self._session.fetch_html(url)
|
||||
data = parse_categories(
|
||||
html,
|
||||
category_id=payload.category_id,
|
||||
include_descendants=payload.include_descendants,
|
||||
)
|
||||
logger.info(
|
||||
"ラクマ 分类完成:category_id=%s name=%s children=%s tree=%s",
|
||||
data.category_id, data.name, len(data.children), data.total_count,
|
||||
)
|
||||
return data
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuma.categories") as span:
|
||||
span.set_attribute("parse.site", "rakuma")
|
||||
span.set_attribute("parse.op", "categories")
|
||||
try:
|
||||
html = await self._session.fetch_html(url)
|
||||
data = parse_categories(
|
||||
html,
|
||||
category_id=payload.category_id,
|
||||
include_descendants=payload.include_descendants,
|
||||
)
|
||||
span.set_attribute("parse.children", len(data.children))
|
||||
span.set_attribute("parse.tree", data.total_count)
|
||||
logger.info(
|
||||
"ラクマ 分类完成:category_id=%s name=%s children=%s tree=%s",
|
||||
data.category_id, data.name, len(data.children), data.total_count,
|
||||
)
|
||||
return data
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
async def item_detail(self, payload: RakumaItemDetailRequest) -> RakumaItemDetailData:
|
||||
"""抓取商品详情"""
|
||||
@@ -97,13 +125,25 @@ class RakumaClient:
|
||||
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
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuma.item_detail") as span:
|
||||
span.set_attribute("parse.site", "rakuma")
|
||||
span.set_attribute("parse.op", "item_detail")
|
||||
try:
|
||||
html = await self._session.fetch_html(url)
|
||||
detail = parse_item_detail(html, item_id=item_id, item_url=url)
|
||||
span.set_attribute("parse.price", detail.price)
|
||||
span.set_attribute("parse.sold_out", detail.is_sold_out)
|
||||
logger.info(
|
||||
"ラクマ 详情完成:url=%s name=%s price=%s sold_out=%s",
|
||||
url, detail.item_name[:40], detail.price, detail.is_sold_out,
|
||||
)
|
||||
return detail
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
async def shop_detail(self, payload: RakumaShopDetailRequest) -> RakumaShopDetailData:
|
||||
"""抓取卖家详情
|
||||
@@ -114,22 +154,34 @@ class RakumaClient:
|
||||
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
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuma.shop_detail") as span:
|
||||
span.set_attribute("parse.site", "rakuma")
|
||||
span.set_attribute("parse.op", "shop_detail")
|
||||
try:
|
||||
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
|
||||
detail = parse_shop_detail(
|
||||
html, shop_id=shop_id, shop_url=url, review_html=review_html
|
||||
)
|
||||
span.set_attribute("parse.items", detail.item_count)
|
||||
span.set_attribute("parse.reviews", detail.review_count)
|
||||
logger.info(
|
||||
"ラクマ 卖家详情完成:shop_id=%s name=%s items=%s reviews=%s",
|
||||
shop_id, detail.shop_name, detail.item_count, detail.review_count,
|
||||
)
|
||||
return detail
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
async def shop_items(self, payload: RakumaShopItemsRequest) -> RakumaShopItemsData:
|
||||
"""抓取卖家的商品列表"""
|
||||
@@ -137,15 +189,27 @@ class RakumaClient:
|
||||
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
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuma.shop_items") as span:
|
||||
span.set_attribute("parse.site", "rakuma")
|
||||
span.set_attribute("parse.op", "shop_items")
|
||||
try:
|
||||
html = await self._session.fetch_html(url)
|
||||
result = parse_shop_items(
|
||||
html, shop_id=shop_id, request_url=url, page=payload.page
|
||||
)
|
||||
span.set_attribute("parse.items", len(result.items))
|
||||
span.set_attribute("parse.total", result.total_count)
|
||||
logger.info(
|
||||
"ラクマ 卖家商品完成:shop_id=%s items=%s total=%s",
|
||||
shop_id, len(result.items), result.total_count,
|
||||
)
|
||||
return result
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
@staticmethod
|
||||
def _resolve_shop_id(shop_url: object, shop_id: str | None) -> str:
|
||||
|
||||
@@ -18,6 +18,7 @@ import logging
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
from opentelemetry import trace
|
||||
|
||||
from app.scraping.core import rakuma_site as site
|
||||
from app.shared.config import Settings
|
||||
@@ -27,6 +28,7 @@ from app.shared.errors import (
|
||||
UpstreamBlockedError,
|
||||
UpstreamRequestError,
|
||||
)
|
||||
from app.shared.telemetry import snapshot
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -89,47 +91,71 @@ class RakumaSession:
|
||||
if client is None:
|
||||
raise UpstreamRequestError("ラクマ 会话尚未初始化")
|
||||
|
||||
tracer = trace.get_tracer(__name__)
|
||||
max_attempts = max(1, self._settings.http_max_attempts)
|
||||
last_error = "unknown error"
|
||||
last_text: str | None = None
|
||||
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
self._semaphore.acquire(),
|
||||
timeout=self._settings.request_timeout_seconds,
|
||||
)
|
||||
except TimeoutError as exc:
|
||||
raise ResourceBusyError() from exc
|
||||
with tracer.start_as_current_span("scrape.fetch") as span:
|
||||
span.set_attribute("scrape.url", url)
|
||||
span.set_attribute("scrape.profile", "rakuma")
|
||||
|
||||
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}"
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
self._semaphore.acquire(),
|
||||
timeout=self._settings.request_timeout_seconds,
|
||||
)
|
||||
except TimeoutError as exc:
|
||||
err = ResourceBusyError()
|
||||
span.record_exception(err)
|
||||
span.set_attribute("scrape.fail_reason", "ResourceBusyError")
|
||||
raise err from exc
|
||||
|
||||
try:
|
||||
for attempt in range(1, max_attempts + 1):
|
||||
span.set_attribute("scrape.attempts", attempt)
|
||||
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:
|
||||
err = ItemNotFoundError(f"Page not found: {url}")
|
||||
span.record_exception(err)
|
||||
span.set_attribute("scrape.fail_reason", "ItemNotFoundError")
|
||||
raise err
|
||||
|
||||
if response.status_code < 400:
|
||||
span.set_attribute("scrape.last_status_code", response.status_code)
|
||||
span.set_attribute("scrape.final_url", str(response.url))
|
||||
span.set_attribute("scrape.html_bytes", len(response.text))
|
||||
return response.text
|
||||
|
||||
last_error = (
|
||||
f"upstream status {response.status_code}"
|
||||
if response.status_code >= 500
|
||||
else f"status {response.status_code}"
|
||||
)
|
||||
last_text = response.text
|
||||
span.set_attribute("scrape.last_status_code", response.status_code)
|
||||
span.set_attribute("scrape.html_bytes", len(response.text))
|
||||
logger.warning(
|
||||
"ラクマ 抓取请求异常:url=%s attempt=%s/%s err=%s",
|
||||
"ラクマ 抓取结果异常:url=%s attempt=%s/%s %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()
|
||||
if last_error.startswith("upstream status") or ":" in last_error:
|
||||
err = UpstreamRequestError(f"Upstream request failed: {last_error}")
|
||||
else:
|
||||
err = UpstreamBlockedError(f"Failed to fetch {url}: {last_error}")
|
||||
span.record_exception(err)
|
||||
span.set_attribute("scrape.fail_reason", type(err).__name__)
|
||||
snapshot(span, "scrape.failed_html", last_text, self._settings.otel_snapshot_max_bytes)
|
||||
raise err
|
||||
finally:
|
||||
self._semaphore.release()
|
||||
|
||||
@@ -8,6 +8,8 @@ from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from opentelemetry import trace
|
||||
|
||||
from app.scraping.core import site
|
||||
from app.shared.config import Settings
|
||||
from app.shared.errors import ScrapeParseError
|
||||
@@ -43,8 +45,10 @@ from app.scraping.utils.urls import (
|
||||
split_item_url,
|
||||
split_shop_url,
|
||||
)
|
||||
from app.shared.telemetry import snapshot
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
tracer = trace.get_tracer(__name__)
|
||||
|
||||
|
||||
class RakutenClient:
|
||||
@@ -62,19 +66,32 @@ class RakutenClient:
|
||||
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
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuten.search") as span:
|
||||
span.set_attribute("parse.site", "rakuten")
|
||||
span.set_attribute("parse.op", "search")
|
||||
try:
|
||||
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,
|
||||
)
|
||||
span.set_attribute("parse.items", len(result.items))
|
||||
span.set_attribute("parse.ads", result.ad_count)
|
||||
span.set_attribute("parse.total", result.total_count)
|
||||
logger.info(
|
||||
"搜索完成:url=%s items=%s ads=%s total=%s",
|
||||
url, len(result.items), result.ad_count, result.total_count,
|
||||
)
|
||||
return result
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
async def genres(self, payload: GenreRequest) -> GenreData:
|
||||
"""抓取分类树:顶层分类列表,或指定分类的信息与直接子分类"""
|
||||
@@ -82,14 +99,25 @@ class RakutenClient:
|
||||
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
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuten.genres") as span:
|
||||
span.set_attribute("parse.site", "rakuten")
|
||||
span.set_attribute("parse.op", "genres")
|
||||
try:
|
||||
html = await self._session.fetch_html(url, mobile=False)
|
||||
state = extract_initial_state(html)
|
||||
result = parse_genres(state, genre_id=genre_id)
|
||||
span.set_attribute("parse.children", len(result.children))
|
||||
logger.info(
|
||||
"分类树完成:genre_id=%s name=%s children=%s",
|
||||
result.genre_id or "root", result.name, len(result.children),
|
||||
)
|
||||
return result
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
async def shop_detail(self, payload: ShopDetailRequest) -> ShopDetailData:
|
||||
"""抓取商家详情"""
|
||||
@@ -100,14 +128,25 @@ class RakutenClient:
|
||||
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
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuten.shop_detail") as span:
|
||||
span.set_attribute("parse.site", "rakuten")
|
||||
span.set_attribute("parse.op", "shop_detail")
|
||||
try:
|
||||
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)
|
||||
span.set_attribute("parse.reviews", result.review_count)
|
||||
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
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
async def shop_items(self, payload: ShopItemsRequest) -> SearchResultData:
|
||||
"""抓取商家名下的商品列表
|
||||
@@ -116,14 +155,25 @@ class RakutenClient:
|
||||
`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
|
||||
with tracer.start_as_current_span("parse.rakuten.shop_items") as span:
|
||||
span.set_attribute("parse.site", "rakuten")
|
||||
span.set_attribute("parse.op", "shop_items")
|
||||
try:
|
||||
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))
|
||||
result = await self.search(payload.to_search_request(shop_id))
|
||||
span.set_attribute("parse.items", len(result.items))
|
||||
span.set_attribute("parse.total", result.total_count)
|
||||
return result
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
raise
|
||||
|
||||
async def item_detail(self, payload: ItemDetailRequest) -> ItemDetailData:
|
||||
"""抓取商品详情"""
|
||||
@@ -137,33 +187,47 @@ class RakutenClient:
|
||||
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,
|
||||
html: str | None = None
|
||||
with tracer.start_as_current_span("parse.rakuten.item_detail") as span:
|
||||
span.set_attribute("parse.site", "rakuten")
|
||||
span.set_attribute("parse.op", "item_detail")
|
||||
try:
|
||||
# 部分官方旗舰店的商品页会跳出市场域名,落到各自的独立站点;
|
||||
# 校验与解析都按最终落地域名分派,落到未登记站点时由校验器抛错。
|
||||
page = await self._session.fetch(
|
||||
url, mobile=True, validator=build_item_page_validator(url)
|
||||
)
|
||||
)
|
||||
html = page.html
|
||||
|
||||
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
|
||||
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,
|
||||
)
|
||||
)
|
||||
|
||||
span.set_attribute("parse.source", detail.source)
|
||||
span.set_attribute("parse.price", detail.price)
|
||||
span.set_attribute("parse.variant_count", detail.sku.variant_count)
|
||||
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
|
||||
except Exception:
|
||||
span.record_exception()
|
||||
span.set_attribute("parse.fail_reason", "parse_error")
|
||||
snapshot(span, "parse.failed_html", html, self._settings.otel_snapshot_max_bytes)
|
||||
raise
|
||||
|
||||
@@ -19,6 +19,7 @@ from dataclasses import dataclass, field
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
from opentelemetry import trace
|
||||
|
||||
from app.scraping.core import site
|
||||
from app.shared.config import Settings
|
||||
@@ -30,6 +31,7 @@ from app.shared.errors import (
|
||||
)
|
||||
from app.scraping.parsers.state import PageValidator, require_state_marker
|
||||
from app.scraping.services.browser_fallback import BrowserFallback
|
||||
from app.shared.telemetry import snapshot
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -152,66 +154,94 @@ class SiteSession:
|
||||
ResourceBusyError: 等待并发槽位超时
|
||||
AppError: 校验器判定为不可重试的失败(如落到不支持的站点)
|
||||
"""
|
||||
tracer = trace.get_tracer(__name__)
|
||||
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"
|
||||
# 失败时供 trace snapshot 上报:保留最后一次拿到的响应体(含挑战页 / 5xx 响应)。
|
||||
last_text: str | None = None
|
||||
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
self._semaphore.acquire(),
|
||||
timeout=self._settings.request_timeout_seconds,
|
||||
)
|
||||
except TimeoutError as exc:
|
||||
raise ResourceBusyError() from exc
|
||||
with tracer.start_as_current_span("scrape.fetch") as span:
|
||||
span.set_attribute("scrape.url", url)
|
||||
span.set_attribute("scrape.profile", profile.name)
|
||||
|
||||
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}"
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
self._semaphore.acquire(),
|
||||
timeout=self._settings.request_timeout_seconds,
|
||||
)
|
||||
except TimeoutError as exc:
|
||||
err = ResourceBusyError()
|
||||
span.record_exception(err)
|
||||
span.set_attribute("scrape.fail_reason", "ResourceBusyError")
|
||||
raise err from exc
|
||||
|
||||
try:
|
||||
for attempt in range(1, max_attempts + 1):
|
||||
span.set_attribute("scrape.attempts", attempt)
|
||||
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:
|
||||
err = ItemNotFoundError(f"Page not found: {url}")
|
||||
span.record_exception(err)
|
||||
span.set_attribute("scrape.fail_reason", "ItemNotFoundError")
|
||||
raise err
|
||||
|
||||
text = response.text
|
||||
final_url = str(response.url)
|
||||
span.set_attribute("scrape.last_status_code", response.status_code)
|
||||
span.set_attribute("scrape.final_url", final_url)
|
||||
span.set_attribute("scrape.html_bytes", len(text))
|
||||
|
||||
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)
|
||||
if last_error.startswith("challenge page detected"):
|
||||
span.set_attribute("scrape.challenge_detected", True)
|
||||
else:
|
||||
last_error = self._describe_error_status(response.status_code)
|
||||
last_text = text
|
||||
logger.warning(
|
||||
"抓取请求异常:url=%s profile=%s attempt=%s/%s err=%s",
|
||||
"抓取结果异常:url=%s profile=%s attempt=%s/%s %s",
|
||||
url, profile.name, attempt, max_attempts, last_error,
|
||||
)
|
||||
continue
|
||||
|
||||
if response.status_code == 404:
|
||||
raise ItemNotFoundError(f"Page not found: {url}")
|
||||
if attempt >= max_attempts:
|
||||
break
|
||||
|
||||
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)
|
||||
# 第一次失败先便宜地换一套 cookie;仍失败才动用浏览器
|
||||
if attempt == 1:
|
||||
await self._invalidate(profile)
|
||||
else:
|
||||
page = await self._escalate_to_browser(profile, url, validate)
|
||||
if page is not None:
|
||||
span.set_attribute("scrape.fell_back_to_browser", True)
|
||||
span.set_attribute("scrape.final_url", page.url)
|
||||
return page
|
||||
|
||||
if self._is_server_error(last_error):
|
||||
err = UpstreamRequestError(f"Upstream request failed: {last_error}")
|
||||
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()
|
||||
err = UpstreamBlockedError(f"Blocked while fetching {url}: {last_error}")
|
||||
span.record_exception(err)
|
||||
span.set_attribute("scrape.fail_reason", type(err).__name__)
|
||||
snapshot(span, "scrape.failed_html", last_text, self._settings.otel_snapshot_max_bytes)
|
||||
raise err
|
||||
finally:
|
||||
self._semaphore.release()
|
||||
|
||||
# ---- 内部 ----
|
||||
|
||||
|
||||
@@ -82,6 +82,18 @@ class Settings(BaseSettings):
|
||||
proxy_username: str | None = None
|
||||
proxy_password: str | None = None
|
||||
|
||||
# ---- OpenTelemetry traces(通用,可选;默认关闭)----
|
||||
# 启用后把抓取-解析链路以 span 导出到 OTLP/HTTP endpoint,用于排查"抓到
|
||||
# 内容能否被解析"。两侧 main.py 在 lifespan 启动时按服务名分别初始化。
|
||||
otel_enabled: bool = False
|
||||
otel_endpoint: str | None = None # OTLP/HTTP traces endpoint,例如 https://oltp.jerryyan.top/v1/traces
|
||||
otel_service_name: str = "rakuten" # 仅作兜底;实际值由两侧 main.py 显式覆盖
|
||||
otel_headers: str | None = None # OTLP 鉴权头,形如 "k=v,k=v";当前 endpoint 裸跑,留空
|
||||
otel_export_interval_ms: int = 5000
|
||||
# 解析失败时把页面 HTML 作为 span event 上报的上限字节;超出截断并标注。
|
||||
# 单个搜索页 HTML 可达 200KB-2MB,调高时同步关注 OTLP 单次请求大小限制。
|
||||
otel_snapshot_max_bytes: int = 2_000_000
|
||||
|
||||
# ---- 登录态与下单配置(仅交易服务)----
|
||||
# 人工登录一次后落盘的 Playwright storage_state 目录(相对项目根目录)。
|
||||
# 目录里是可直接冒充账号的 cookie,务必不要提交到版本库。
|
||||
|
||||
@@ -0,0 +1,145 @@
|
||||
"""OpenTelemetry traces 接入:仅在 otel_enabled=true 时初始化,否则全 noop
|
||||
|
||||
抓取与交易两个进程在 lifespan 启动时各自调用 `setup_telemetry(settings,
|
||||
service_name=...)`:注册 TracerProvider + OTLP/HTTP exporter + 自动
|
||||
instrumentation(FastAPI、httpx)。失败时(如 endpoint 不可达)不阻断主流程,
|
||||
仅打日志;traces 是辅助观测,不应让进程起不来。
|
||||
|
||||
`shutdown_telemetry` 在 lifespan 关闭时 force_flush 后再 shutdown,确保缓冲区
|
||||
里的 span 都已上报。
|
||||
|
||||
OTel SDK 默认的 ProxyTracerProvider 在 setup 之前就能用(noop span),所以
|
||||
其它代码里直接 `trace.get_tracer(__name__)` + `start_as_current_span` 即可,
|
||||
不必关心 telemetry 是否启用——禁用时 span 不会真正产生与上报。
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import TYPE_CHECKING
|
||||
|
||||
from opentelemetry import trace
|
||||
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
|
||||
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor
|
||||
from opentelemetry.instrumentation.httpx import HTTPXClientInstrumentor
|
||||
from opentelemetry.sdk.resources import SERVICE_NAME, Resource
|
||||
from opentelemetry.sdk.trace import TracerProvider
|
||||
from opentelemetry.sdk.trace.export import BatchSpanProcessor
|
||||
from opentelemetry.sdk.trace.sampling import ALWAYS_ON
|
||||
from opentelemetry.trace import Span
|
||||
from opentelemetry.util.types import AttributeValue
|
||||
|
||||
from app.shared.config import Settings
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from fastapi import FastAPI
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# 全局 provider 引用,用于 instrument_app / shutdown 时判断当前是否已初始化。
|
||||
# 显式持有比依赖 trace.get_tracer_provider() 的类型判断更稳——后者在测试场景
|
||||
# 下可能被其它用例改动全局状态。
|
||||
_provider: TracerProvider | None = None
|
||||
|
||||
|
||||
def setup_telemetry(settings: Settings, *, service_name: str) -> None:
|
||||
"""初始化 OTel:TracerProvider + OTLP exporter + httpx 自动 instrumentation。
|
||||
|
||||
- `otel_enabled=False` 或 endpoint 未配时仅打日志,不做任何事。
|
||||
- 必须在创建任何 httpx.AsyncClient 之前调用,否则 httpx 不会被打桩。
|
||||
两侧 main.py 的 lifespan 已把 setup 放在 site_session.start() 之前。
|
||||
- 重复调用安全(_provider 已设时直接返回)。
|
||||
"""
|
||||
global _provider
|
||||
if _provider is not None:
|
||||
return
|
||||
if not settings.otel_enabled or not settings.otel_endpoint:
|
||||
logger.info("OpenTelemetry 未启用(endpoint 或 otel_enabled 未配置)")
|
||||
return
|
||||
|
||||
resource = Resource.create({SERVICE_NAME: service_name})
|
||||
provider = TracerProvider(resource=resource, sampler=ALWAYS_ON)
|
||||
|
||||
exporter = OTLPSpanExporter(
|
||||
endpoint=settings.otel_endpoint,
|
||||
headers=_parse_headers(settings.otel_headers),
|
||||
timeout=10,
|
||||
)
|
||||
provider.add_span_processor(
|
||||
BatchSpanProcessor(
|
||||
exporter,
|
||||
schedule_delay_millis=settings.otel_export_interval_ms,
|
||||
)
|
||||
)
|
||||
trace.set_tracer_provider(provider)
|
||||
_provider = provider
|
||||
|
||||
# httpx 是抓取/交易两侧唯一的外部 HTTP 客户端,打桩后所有 AsyncClient 请求
|
||||
# 自动产生 CLIENT span。失败不影响主链路:进程仍可运行,只是看不到 span。
|
||||
try:
|
||||
HTTPXClientInstrumentor().instrument()
|
||||
except Exception:
|
||||
logger.warning("httpx 自动 instrumentation 失败", exc_info=True)
|
||||
|
||||
logger.info(
|
||||
"OpenTelemetry 已启用:endpoint=%s service=%s",
|
||||
settings.otel_endpoint,
|
||||
service_name,
|
||||
)
|
||||
|
||||
|
||||
def instrument_app(app: "FastAPI") -> None:
|
||||
"""FastAPI 应用打桩;未初始化时 noop,调用顺序无要求。"""
|
||||
if _provider is None:
|
||||
return
|
||||
FastAPIInstrumentor.instrument_app(app)
|
||||
|
||||
|
||||
def shutdown_telemetry() -> None:
|
||||
"""flush + shutdown;幂等,未初始化时直接返回。"""
|
||||
global _provider
|
||||
if _provider is None:
|
||||
return
|
||||
try:
|
||||
_provider.force_flush()
|
||||
_provider.shutdown()
|
||||
except Exception:
|
||||
logger.debug("关闭 OpenTelemetry provider 失败", exc_info=True)
|
||||
_provider = None
|
||||
|
||||
|
||||
def snapshot(span: Span, name: str, html: str | None, max_bytes: int) -> None:
|
||||
"""把 HTML 作为 span event 上报,超 max_bytes 截断并标注。
|
||||
|
||||
用于解析失败时复现页面:span 自身只放结构化指标(items 数、source 等),
|
||||
完整 HTML 体量大、含商品/价格内容,仅在失败分支通过 event 携带。
|
||||
"""
|
||||
if html is None or not html:
|
||||
return
|
||||
original_bytes = len(html)
|
||||
truncated = original_bytes > max_bytes
|
||||
payload = html if not truncated else html[:max_bytes]
|
||||
attributes: dict[str, AttributeValue] = {
|
||||
"snapshot.html": payload,
|
||||
"snapshot.original_bytes": original_bytes,
|
||||
}
|
||||
if truncated:
|
||||
attributes["snapshot.truncated"] = True
|
||||
span.add_event(name, attributes=attributes)
|
||||
|
||||
|
||||
def _parse_headers(raw: str | None) -> list[tuple[str, str]] | None:
|
||||
"""解析 "k1=v1,k2=v2" 形式的 header 配置;空输入返回 None。"""
|
||||
if not raw or not raw.strip():
|
||||
return None
|
||||
pairs: list[tuple[str, str]] = []
|
||||
for chunk in raw.split(","):
|
||||
key, sep, value = chunk.partition("=")
|
||||
if not sep or not key.strip():
|
||||
continue
|
||||
pairs.append((key.strip(), value.strip()))
|
||||
return pairs or None
|
||||
|
||||
|
||||
def is_initialized() -> bool:
|
||||
"""供测试断言使用。"""
|
||||
return _provider is not None
|
||||
@@ -22,6 +22,7 @@ from fastapi import FastAPI
|
||||
from app.shared.api import register_exception_handlers
|
||||
from app.shared.config import get_settings
|
||||
from app.shared.logging_setup import configure_logging
|
||||
from app.shared.telemetry import instrument_app, setup_telemetry, shutdown_telemetry
|
||||
from app.trading.api.routes.auth import router as auth_router
|
||||
from app.trading.api.routes.health import router as health_router
|
||||
from app.trading.container import TradingContainer
|
||||
@@ -43,6 +44,8 @@ async def lifespan(app: FastAPI):
|
||||
app.state.container = container
|
||||
|
||||
configure_logging(container.settings)
|
||||
# 与抓取侧同样:必须在创建 httpx 客户端之前 setup。
|
||||
setup_telemetry(container.settings, service_name="rakuten-trading")
|
||||
logger.info(
|
||||
"交易服务启动:%s:%s",
|
||||
container.settings.trading_host,
|
||||
@@ -54,6 +57,7 @@ async def lifespan(app: FastAPI):
|
||||
yield
|
||||
finally:
|
||||
await container.auth_session.close()
|
||||
shutdown_telemetry()
|
||||
|
||||
|
||||
def create_app() -> FastAPI:
|
||||
@@ -62,6 +66,7 @@ def create_app() -> FastAPI:
|
||||
app.include_router(health_router)
|
||||
app.include_router(auth_router)
|
||||
register_exception_handlers(app)
|
||||
instrument_app(app)
|
||||
return app
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user