From 865bbe372467742cbd748ac5b68b647980aede3a Mon Sep 17 00:00:00 2001 From: Jerry Yan <792602257@qq.com> Date: Fri, 28 Aug 2026 16:18:25 +0800 Subject: [PATCH] =?UTF-8?q?feat(observability):=20=E6=8A=93=E5=8F=96?= =?UTF-8?q?=E4=BC=9A=E8=AF=9D=E9=80=90=E6=AC=A1=E5=B0=9D=E8=AF=95=E4=B8=8E?= =?UTF-8?q?=E5=8D=87=E7=BA=A7=E8=B7=AF=E5=BE=84=E8=AE=B0=E6=88=90=20event?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 重试链路的关键信息是「升级路径」,而属性表达不了过程:同名属性后写覆盖先写, 三次尝试跑完只剩最后一次的状态码,前两次为什么失败、升到哪一级全被盖掉。 - site_session / rakuma_session:每次尝试各记一条 scrape.attempt event,带 attempt / outcome / status_code / 截断后的错误串;outcome 区分 http_error、 bad_status、challenge、validate_failed - site_session 的每次升级各记一条 scrape.escalate:rewarm_on_home、浏览器兜底 recovered / failed。兜底失败也要记——不然链路里只剩「最终失败」,看不出浏览器 这一级试过没有,而「没装 playwright」和「装了也被挡」是两个查法(reason 取 BrowserFallback.unavailable_reason,为 None 时由 add_event 跳过) - 三处 span.record_exception 换成 record_error,失败的错误码与可重试位跟着落下来 - 失败页面快照带上 url 与 profile,省得回头对是哪条通道的哪个地址 测试:test_scrape_telemetry.py 补一条会话层用例,用 MockTransport 让所有请求都回 挑战页,断言三次尝试各留一条 event(而不是只剩最后一次)、升级路径依次是 rewarm_on_home → browser,以及浏览器那一级的 outcome=failed 与 reason。 Co-Authored-By: Claude Opus 5 (1M context) --- app/scraping/services/rakuma_session.py | 31 +++++++-- app/scraping/services/site_session.py | 53 ++++++++++++-- tests/test_scrape_telemetry.py | 93 +++++++++++++++++++++++++ 3 files changed, 166 insertions(+), 11 deletions(-) diff --git a/app/scraping/services/rakuma_session.py b/app/scraping/services/rakuma_session.py index b78f658..2b199ae 100644 --- a/app/scraping/services/rakuma_session.py +++ b/app/scraping/services/rakuma_session.py @@ -29,10 +29,14 @@ from app.shared.errors import ( UpstreamBlockedError, UpstreamRequestError, ) -from app.shared.telemetry import snapshot +from app.shared.telemetry import add_event, record_error, snapshot logger = logging.getLogger(__name__) +# 逐次尝试 event 里错误串的截断长度:httpx 的异常消息可能很长(带完整 URL 与 +# 底层 socket 错误),而这里只是用来区分「这次是怎么失败的」,前半句就够 +_ERROR_MAX_CHARS = 200 + class RakumaSession: """ラクマ 站点抓取会话,管理并发限流与失败重试""" @@ -108,7 +112,7 @@ class RakumaSession: ) except TimeoutError as exc: err = ResourceBusyError() - span.record_exception(err) + record_error(span, err) span.set_attribute("scrape.fail_reason", "ResourceBusyError") raise err from exc @@ -119,6 +123,13 @@ class RakumaSession: response = await client.get(url) except httpx.HTTPError as exc: last_error = f"{type(exc).__name__}: {exc}" + # 每次尝试各记一条 event:属性同名后写覆盖先写,三次尝试 + # 跑完只剩最后一次的状态,中间那两次为什么失败全被盖掉 + add_event(span, "scrape.attempt", { + "attempt": attempt, + "outcome": "http_error", + "error": last_error[:_ERROR_MAX_CHARS], + }) logger.warning( "ラクマ 抓取请求异常:url=%s attempt=%s/%s err=%s", url, attempt, max_attempts, last_error, @@ -127,7 +138,7 @@ class RakumaSession: if response.status_code == 404: err = ItemNotFoundError(f"Page not found: {url}") - span.record_exception(err) + record_error(span, err) span.set_attribute("scrape.fail_reason", "ItemNotFoundError") raise err @@ -145,6 +156,12 @@ class RakumaSession: last_text = response.text span.set_attribute("scrape.last_status_code", response.status_code) span.set_attribute("scrape.html_bytes", len(response.text)) + add_event(span, "scrape.attempt", { + "attempt": attempt, + "outcome": "bad_status", + "status_code": response.status_code, + "error": last_error[:_ERROR_MAX_CHARS], + }) logger.warning( "ラクマ 抓取结果异常:url=%s attempt=%s/%s %s", url, attempt, max_attempts, last_error, @@ -154,9 +171,13 @@ class RakumaSession: err = UpstreamRequestError(f"Upstream request failed: {last_error}") else: err = UpstreamBlockedError(f"Failed to fetch {url}: {last_error}") - span.record_exception(err) + record_error(span, err) span.set_attribute("scrape.fail_reason", type(err).__name__) - snapshot(span, "scrape.failed_html", last_text, self._settings.otel_snapshot_max_bytes) + snapshot( + span, "scrape.failed_html", last_text, + self._settings.otel_snapshot_max_bytes, + extra={"scrape.url": url, "scrape.profile": "rakuma"}, + ) raise err finally: self._semaphore.release() diff --git a/app/scraping/services/site_session.py b/app/scraping/services/site_session.py index 44c5994..90b18da 100644 --- a/app/scraping/services/site_session.py +++ b/app/scraping/services/site_session.py @@ -38,10 +38,14 @@ 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 +from app.shared.telemetry import add_event, record_error, snapshot logger = logging.getLogger(__name__) +# 逐次尝试 event 里错误串的截断长度:httpx 异常消息可能很长(带完整 URL 与底层 +# socket 错误),这里只用来区分「这次是怎么失败的」,前半句够了 +_ERROR_MAX_CHARS = 200 + @dataclass(slots=True) class FetchedPage: @@ -191,7 +195,7 @@ class SiteSession: ) except TimeoutError as exc: err = ResourceBusyError() - span.record_exception(err) + record_error(span, err) span.set_attribute("scrape.fail_reason", "ResourceBusyError") raise err from exc @@ -204,6 +208,14 @@ class SiteSession: response = await profile.client.get(url) except httpx.HTTPError as exc: last_error = f"{type(exc).__name__}: {exc}" + # 每次尝试各记一条 event。属性同名后写覆盖先写,三次尝试跑完 + # 只剩最后一次的状态码,前两次为什么失败、升级到哪一级全被 + # 盖掉——而这条链路的关键信息恰恰是「升级路径」。 + add_event(span, "scrape.attempt", { + "attempt": attempt, + "outcome": "http_error", + "error": last_error[:_ERROR_MAX_CHARS], + }) logger.warning( "抓取请求异常:url=%s profile=%s attempt=%s/%s err=%s", url, profile.name, attempt, max_attempts, last_error, @@ -216,7 +228,7 @@ class SiteSession: if response.status_code == 404: err = ItemNotFoundError(f"Page not found: {url}") - span.record_exception(err) + record_error(span, err) span.set_attribute("scrape.fail_reason", "ItemNotFoundError") raise err @@ -232,11 +244,21 @@ class SiteSession: 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"): + challenged = last_error.startswith("challenge page detected") + if challenged: span.set_attribute("scrape.challenge_detected", True) + outcome = "challenge" if challenged else "validate_failed" else: last_error = self._describe_error_status(response.status_code) + outcome = "bad_status" last_text = text + add_event(span, "scrape.attempt", { + "attempt": attempt, + "outcome": outcome, + "status_code": response.status_code, + "html_bytes": len(text), + "error": last_error[:_ERROR_MAX_CHARS], + }) logger.warning( "抓取结果异常:url=%s profile=%s attempt=%s/%s %s", url, profile.name, attempt, max_attempts, last_error, @@ -250,20 +272,39 @@ class SiteSession: if attempt == 1: await self._rewarm_on_home(profile) span.set_attribute("scrape.rewarmed_on_home", True) + add_event(span, "scrape.escalate", { + "attempt": attempt, "to": "rewarm_on_home", + }) 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) + add_event(span, "scrape.escalate", { + "attempt": attempt, "to": "browser", "outcome": "recovered", + }) return page + # 兜底没救回来(浏览器不可用,或取回的页面仍不合格)。不记 + # 这条的话链路里只看得到「最终失败」,看不出浏览器这一级到底 + # 试过没有——而「没装 playwright」和「装了也被挡」要分开查。 + add_event(span, "scrape.escalate", { + "attempt": attempt, + "to": "browser", + "outcome": "failed", + "reason": self._browser.unavailable_reason, + }) if self._is_server_error(last_error): err = UpstreamRequestError(f"Upstream request failed: {last_error}") else: err = UpstreamBlockedError(f"Blocked while fetching {url}: {last_error}") - span.record_exception(err) + record_error(span, err) span.set_attribute("scrape.fail_reason", type(err).__name__) - snapshot(span, "scrape.failed_html", last_text, self._settings.otel_snapshot_max_bytes) + snapshot( + span, "scrape.failed_html", last_text, + self._settings.otel_snapshot_max_bytes, + extra={"scrape.url": url, "scrape.profile": profile.name}, + ) raise err finally: self._semaphore.release() diff --git a/tests/test_scrape_telemetry.py b/tests/test_scrape_telemetry.py index 0395956..6adc783 100644 --- a/tests/test_scrape_telemetry.py +++ b/tests/test_scrape_telemetry.py @@ -13,6 +13,7 @@ from __future__ import annotations import json +import httpx import pytest from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import SimpleSpanProcessor @@ -27,8 +28,10 @@ from app.scraping.models.scrape import ( ) from app.scraping.services import rakuma_client as rakuma_module from app.scraping.services import rakuten_client as rakuten_module +from app.scraping.services import site_session as session_module from app.scraping.services.rakuma_client import RakumaClient from app.scraping.services.rakuten_client import RakutenClient +from app.scraping.services.site_session import SiteSession from app.shared.config import Settings from app.shared.errors import AppError, ScrapeParseError, UpstreamBlockedError @@ -183,6 +186,96 @@ async def test_rakuma_parse_failure_records_reason_and_snapshot(spans): assert snapshot.attributes["snapshot.html"] == html +# ---- 会话层:逐次尝试与升级路径 ---- +# +# 这一层的关键信息是「升级路径」:换 cookie → 浏览器兜底 → 放弃。属性表达不了过程 +# (同名后写覆盖先写,三次尝试跑完只剩最后一次),所以每次尝试与每次升级各记一条 +# event。site_session / rakuma_session 都在 fetch 内部 `trace.get_tracer(__name__)`, +# 所以这里替 `trace.get_tracer` 本身(与 test_telemetry.py 的 spans 夹具同一手法)。 + + +@pytest.fixture +def session_spans(monkeypatch) -> InMemorySpanExporter: + exporter = InMemorySpanExporter() + provider = TracerProvider() + provider.add_span_processor(SimpleSpanProcessor(exporter)) + monkeypatch.setattr(session_module.trace, "get_tracer", provider.get_tracer) + return exporter + + +class _FakeBrowser: + """浏览器兜底替身:visit 恒不可用,让升级链走完整条路""" + + def __init__(self) -> None: + self.unavailable_reason = "playwright is not installed" + + async def visit(self, url: str, *, mobile: bool): + return None + + async def close(self) -> None: + pass + + +async def _blocked_session() -> SiteSession: + """所有请求都回 Akamai 挑战页的会话,用 MockTransport 拦截""" + settings = Settings( + _env_file=None, http_max_attempts=3, request_timeout_seconds=5.0, + max_site_concurrency=4, + ) + session = SiteSession(settings, _FakeBrowser()) + await session.start() + for profile in session._profiles.values(): + await profile.client.aclose() + profile.client = httpx.AsyncClient( + transport=httpx.MockTransport( + lambda _r: httpx.Response( + 200, text="Access Denied. Reference #18.abc" + ) + ), + follow_redirects=True, + ) + return session + + +async def test_retry_attempts_and_escalation_are_recorded_as_events(session_spans): + """每次尝试与每次升级各留一条 event,链路上能读出完整升级路径 + + 这些信息无法用属性表达:`scrape.attempts` 只剩最后一次的值,前两次为什么失败、 + 换过 cookie 没有、浏览器兜底试过没有全被覆盖掉。 + """ + session = await _blocked_session() + try: + with pytest.raises(UpstreamBlockedError): + await session.fetch_html( + "https://search.rakuten.co.jp/search/mall/x/", mobile=False + ) + finally: + await session.close() + + span = _span(session_spans, "scrape.fetch") + assert span.status.status_code is StatusCode.ERROR + assert span.attributes["scrape.fail_reason"] == "UpstreamBlockedError" + assert span.attributes["scrape.challenge_detected"] is True + + # 三次尝试都留下了自己的记录,而不是只剩最后一次 + attempts = [e for e in span.events if e.name == "scrape.attempt"] + assert [e.attributes["attempt"] for e in attempts] == [1, 2, 3] + assert {e.attributes["outcome"] for e in attempts} == {"challenge"} + + # 升级路径:第 1 次失败后换 cookie,第 2 次失败后动用浏览器(本例不可用) + escalations = [e for e in span.events if e.name == "scrape.escalate"] + assert [e.attributes["to"] for e in escalations] == ["rewarm_on_home", "browser"] + browser_step = escalations[-1] + assert browser_step.attributes["outcome"] == "failed" + # 「没装 playwright」与「装了也被挡」要能分开查 + assert browser_step.attributes["reason"] == "playwright is not installed" + + # 最终失败页面的快照带上来源,便于确认是哪条通道的哪个地址 + failed = next(e for e in span.events if e.name == "scrape.failed_html") + assert failed.attributes["scrape.profile"] == "pc" + assert "search.rakuten.co.jp" in failed.attributes["scrape.url"] + + async def test_successful_scrape_leaves_span_ok(spans, search_state): """成功路径不被误标 ERROR,不落失败快照,且照常记结果指标