perf: 量化商品页打开耗时并去重检查 (#109)

This commit is contained in:
chengma
2026-08-10 17:40:44 +08:00
parent 4c9bf3865b
commit ba0706dc47
13 changed files with 693 additions and 72 deletions
+68 -31
View File
@@ -11,6 +11,7 @@ import math
import re
import time
import xml.etree.ElementTree as ET
from contextlib import nullcontext
from dataclasses import dataclass
from datetime import datetime, timezone
from decimal import Decimal, InvalidOperation, ROUND_HALF_UP
@@ -19,6 +20,7 @@ from typing import Any, Callable, Iterable, Mapping, Optional, Sequence
from urllib.parse import parse_qs, urlparse
from .pdd_device_service import PDD_PACKAGE_NAME, PddDeviceError, PddDeviceService
from .performance_timing import current_performance_trace
from .util.get_size_panle_coord import get_size_panel_coord
@@ -691,8 +693,11 @@ class PddCollectService:
self._check_cancelled()
try:
with self._device_service.connect(self._device_address) as device:
self._open_goods(device, goods_url)
session = self._device_service.connect(self._device_address)
with session as device:
self._open_goods(
device, goods_url, session.initial_app_state
)
goods = self._collect_goods_details(device)
if not goods.title:
raise PddCollectError("PDD_DATA_TITLE_MISSING", "商品页没有可识别的标题")
@@ -779,14 +784,32 @@ class PddCollectService:
):
raise PddCollectError("PDD_PAGE_OVERALL_TIMEOUT", "PDD 采集超过 10 分钟")
def _open_goods(self, device: Any, goods_url: str) -> None:
def _open_goods(
self,
device: Any,
goods_url: str,
initial_app_state: Optional[Mapping[str, Any]] = None,
) -> None:
trace = current_performance_trace()
try:
current = device.app_current()
current = (
dict(initial_app_state)
if initial_app_state is not None
else device.app_current()
)
if current.get("package") != PDD_PACKAGE_NAME:
device.app_start(PDD_PACKAGE_NAME)
if not device.app_wait(PDD_PACKAGE_NAME, timeout=10):
raise PddCollectError("DEVICE_APP_START_FAILED", "PDD 应用启动失败")
device.open_url(goods_url)
stage = trace.stage("pdd_start_or_wait") if trace else nullcontext()
with stage:
device.app_start(PDD_PACKAGE_NAME)
if not device.app_wait(PDD_PACKAGE_NAME, timeout=10):
raise PddCollectError(
"DEVICE_APP_START_FAILED", "PDD 应用启动失败"
)
elif trace is not None:
trace.record("pdd_start_or_wait", 0, "already_foreground")
stage = trace.stage("open_url") if trace else nullcontext()
with stage:
device.open_url(goods_url)
except PddCollectError:
raise
except Exception as exc:
@@ -797,30 +820,44 @@ class PddCollectService:
) from exc
raise PddCollectError("DEVICE_APP_START_FAILED", f"无法打开 PDD 商品链接:{exc}") from exc
deadline = self._monotonic() + self._page_timeout
while self._monotonic() < deadline:
self._check_cancelled()
current = device.app_current()
xml_data = self._dump_hierarchy(device)
root = _parse_xml(xml_data)
last_labels = _all_labels(root)
pdd_node_count = sum(
1
for node in root.iter("node")
if node.get("package") == PDD_PACKAGE_NAME
ready_stage = trace.stage("goods_page_ready") if trace else nullcontext()
with ready_stage:
deadline = self._monotonic() + self._page_timeout
first_dump = True
while self._monotonic() < deadline:
self._check_cancelled()
current = device.app_current()
if first_dump and trace is not None:
with trace.stage("first_dump_hierarchy"):
xml_data = self._dump_hierarchy(device)
first_dump = False
else:
xml_data = self._dump_hierarchy(device)
root = _parse_xml(xml_data)
last_labels = _all_labels(root)
pdd_node_count = sum(
1
for node in root.iter("node")
if node.get("package") == PDD_PACKAGE_NAME
)
# 部分 OPPO/ColorOS 设备会一直把无线调试设置页报告为焦点,
# 即使屏幕和无障碍树已经是 PDD。此时以树中真实包名为准。
is_pdd_hierarchy = pdd_node_count >= 3
is_pdd_focused = current.get("package") == PDD_PACKAGE_NAME
if is_pdd_focused or is_pdd_hierarchy:
_raise_special_page(last_labels)
combined = " ".join(last_labels)
loading = "加载中" in combined or "正在加载" in combined
if not loading and any(
marker in combined for marker in _READY_MARKERS
):
if trace is not None:
trace.checkpoint("end_to_end_total")
return
self._sleep(0.5)
raise PddCollectError(
"PDD_PAGE_TIMEOUT", "等待 PDD 商品详情页加载超时"
)
# 部分 OPPO/ColorOS 设备会一直把无线调试设置页报告为焦点,
# 即使屏幕和无障碍树已经是 PDD。此时以树中真实包名为准。
is_pdd_hierarchy = pdd_node_count >= 3
is_pdd_focused = current.get("package") == PDD_PACKAGE_NAME
if is_pdd_focused or is_pdd_hierarchy:
_raise_special_page(last_labels)
combined = " ".join(last_labels)
loading = "加载中" in combined or "正在加载" in combined
if not loading and any(marker in combined for marker in _READY_MARKERS):
return
self._sleep(0.5)
raise PddCollectError("PDD_PAGE_TIMEOUT", "等待 PDD 商品详情页加载超时")
def _collect_goods_details(self, device: Any) -> GoodsSnapshot:
"""保留首屏结果,并向下浏览补采评价数量和店铺名。"""
+16 -2
View File
@@ -10,6 +10,8 @@ import threading
from contextlib import AbstractContextManager
from typing import Any, Callable, Optional
from .performance_timing import current_performance_trace
PDD_PACKAGE_NAME = "com.xunmeng.pinduoduo"
@@ -30,10 +32,12 @@ class PddDeviceSession(AbstractContextManager[Any]):
self,
device: Any,
serial: str,
initial_app_state: dict[str, Any],
release: Callable[[], None],
) -> None:
self._device = device
self.serial = serial
self.initial_app_state = dict(initial_app_state)
self._release = release
self._owner_thread_id = threading.get_ident()
self._closed = False
@@ -107,8 +111,17 @@ class PddDeviceService:
self._active_serial = checked_serial
try:
device = self._connector(checked_serial)
current = device.app_current()
trace = current_performance_trace()
if trace is None:
device = self._connector(checked_serial)
else:
with trace.stage("uiautomator2_connect"):
device = self._connector(checked_serial)
if trace is None:
current = device.app_current()
else:
with trace.stage("first_app_current"):
current = device.app_current()
if not isinstance(current, dict):
raise RuntimeError("uiautomator2 未返回有效的设备状态")
except PddDeviceError:
@@ -130,6 +143,7 @@ class PddDeviceService:
return PddDeviceSession(
device,
checked_serial,
current,
lambda: self._release(checked_serial),
)
+55 -30
View File
@@ -9,6 +9,7 @@ from __future__ import annotations
import re
import time
import xml.etree.ElementTree as ET
from contextlib import nullcontext
from decimal import Decimal, InvalidOperation, ROUND_HALF_UP
from typing import Any, Callable, Mapping, Optional
from urllib.parse import parse_qs, urlparse
@@ -18,6 +19,7 @@ from .pdd_device_service import (
PddDeviceError,
PddDeviceService,
)
from .performance_timing import current_performance_trace
from .pdd_purchase_adapter import (
PddLivePurchaseAdapter,
PddPurchaseAdapter,
@@ -266,17 +268,24 @@ class U2PddPurchaseAdapter(PddPurchaseAdapter):
try:
self._session = self._device_service.connect(self._device_address)
self._device = self._session.__enter__()
current = self._device.app_current()
current = self._session.initial_app_state
trace = current_performance_trace()
if current.get("package") != PDD_PACKAGE_NAME:
self._device.app_start(PDD_PACKAGE_NAME)
if not self._device.app_wait(PDD_PACKAGE_NAME, timeout=10):
raise PddPurchaseError(
"DEVICE_APP_START_FAILED",
"PDD 应用启动失败",
step="purchase_open_goods",
retryable=True,
)
self._device.open_url(goods_url)
stage = trace.stage("pdd_start_or_wait") if trace else nullcontext()
with stage:
self._device.app_start(PDD_PACKAGE_NAME)
if not self._device.app_wait(PDD_PACKAGE_NAME, timeout=10):
raise PddPurchaseError(
"DEVICE_APP_START_FAILED",
"PDD 应用启动失败",
step="purchase_open_goods",
retryable=True,
)
elif trace is not None:
trace.record("pdd_start_or_wait", 0, "already_foreground")
stage = trace.stage("open_url") if trace else nullcontext()
with stage:
self._device.open_url(goods_url)
self._wait_for_page("goods", self._page_timeout)
except PddPurchaseError:
raise
@@ -436,26 +445,42 @@ class U2PddPurchaseAdapter(PddPurchaseAdapter):
session.__exit__(None, None, None)
def _wait_for_page(self, expected: str, timeout: float) -> str:
deadline = self._monotonic() + timeout
last_kind = "unknown"
while self._monotonic() < deadline:
self._check_cancelled("purchase_open_goods")
device = self._require_device()
current = device.app_current()
xml_data = self._dump_hierarchy()
root = _parse_xml(xml_data)
last_kind = _page_kind(root, str(current.get("package") or ""))
if last_kind == expected:
return xml_data
if last_kind in {"captcha", "login_required", "risk_control", "payment"}:
self._raise_special_page(last_kind)
self._sleep(0.25)
raise PddPurchaseError(
"PDD_PAGE_TIMEOUT",
f"等待 PDD {expected} 页面超时,最后页面为 {last_kind}",
step="purchase_open_goods",
retryable=True,
)
trace = current_performance_trace()
ready_stage = trace.stage("goods_page_ready") if trace else nullcontext()
with ready_stage:
deadline = self._monotonic() + timeout
last_kind = "unknown"
first_dump = True
while self._monotonic() < deadline:
self._check_cancelled("purchase_open_goods")
device = self._require_device()
current = device.app_current()
if first_dump and trace is not None:
with trace.stage("first_dump_hierarchy"):
xml_data = self._dump_hierarchy()
first_dump = False
else:
xml_data = self._dump_hierarchy()
root = _parse_xml(xml_data)
last_kind = _page_kind(root, str(current.get("package") or ""))
if last_kind == expected:
if trace is not None:
trace.checkpoint("end_to_end_total")
return xml_data
if last_kind in {
"captcha",
"login_required",
"risk_control",
"payment",
}:
self._raise_special_page(last_kind)
self._sleep(0.25)
raise PddPurchaseError(
"PDD_PAGE_TIMEOUT",
f"等待 PDD {expected} 页面超时,最后页面为 {last_kind}",
step="purchase_open_goods",
retryable=True,
)
def _wait_for_confirmation_panel(self) -> str:
deadline = self._monotonic() + self._panel_timeout
+149
View File
@@ -0,0 +1,149 @@
"""任务领取到商品页的脱敏分阶段计时。
日志字段使用白名单,只包含任务编号、稳定操作名、毫秒耗时和结果分类。
不要把 URL、控件树、异常正文或业务数据传入本模块。
"""
from __future__ import annotations
import json
import logging
from contextlib import contextmanager
from contextvars import ContextVar
from logging.handlers import RotatingFileHandler
from pathlib import Path
import time
from typing import Callable, Iterator, Optional
from .db import data_dir
PERFORMANCE_LOG_NAME = "task_performance.jsonl"
_LOGGER_NAME = "cmautobuy.performance"
_ACTIVE_TRACE: ContextVar[Optional["TaskPerformanceTrace"]] = ContextVar(
"cmautobuy_task_performance_trace", default=None
)
def _default_sink(record: dict[str, object]) -> None:
logging.getLogger(_LOGGER_NAME).info(
json.dumps(record, ensure_ascii=False, separators=(",", ":"))
)
def configure_performance_logging(
directory: Optional[Path] = None,
) -> Path:
"""把性能事件写入可轮转 JSONL 文件,重复调用不会增加处理器。"""
log_directory = directory or data_dir() / "logs"
log_directory.mkdir(parents=True, exist_ok=True)
log_path = log_directory / PERFORMANCE_LOG_NAME
logger = logging.getLogger(_LOGGER_NAME)
if not any(
getattr(handler, "baseFilename", None) == str(log_path.resolve())
for handler in logger.handlers
):
handler = RotatingFileHandler(
log_path,
maxBytes=2 * 1024 * 1024,
backupCount=3,
encoding="utf-8",
)
handler.setFormatter(logging.Formatter("%(message)s"))
logger.addHandler(handler)
logger.setLevel(logging.INFO)
logger.propagate = False
return log_path
class TaskPerformanceTrace:
"""收集一条任务的阶段耗时,并在绑定任务编号后输出记录。"""
def __init__(
self,
*,
monotonic: Callable[[], float] = time.monotonic,
sink: Callable[[dict[str, object]], None] = _default_sink,
) -> None:
self._monotonic = monotonic
self._sink = sink
self._started_at = monotonic()
self._task_id = ""
self._pending: list[tuple[str, int, str]] = []
@property
def task_id(self) -> str:
return self._task_id
def bind_task(self, task_id: str) -> None:
"""绑定稳定任务编号,并输出绑定前缓存的领取阶段记录。"""
value = str(task_id or "").strip()
if not value:
return
if self._task_id and self._task_id != value:
raise RuntimeError("性能计时不能跨任务复用")
self._task_id = value
pending, self._pending = self._pending, []
for operation, duration_ms, result in pending:
self._emit(operation, duration_ms, result)
@contextmanager
def stage(self, operation: str) -> Iterator[None]:
"""记录一个阶段;异常时只写 failed,不记录异常正文。"""
started_at = self._monotonic()
try:
yield
except BaseException:
self.record(operation, self._monotonic() - started_at, "failed")
raise
else:
self.record(operation, self._monotonic() - started_at, "succeeded")
def record(self, operation: str, seconds: float, result: str) -> None:
"""记录白名单阶段数据,负耗时会安全归零。"""
checked_operation = str(operation or "").strip()
checked_result = str(result or "").strip()
if not checked_operation or not checked_result:
raise ValueError("性能操作名和结果分类不能为空")
duration_ms = max(0, round(float(seconds) * 1000))
if self._task_id:
self._emit(checked_operation, duration_ms, checked_result)
else:
self._pending.append(
(checked_operation, duration_ms, checked_result)
)
def checkpoint(self, operation: str, result: str = "succeeded") -> None:
"""记录从本轮调度开始到当前时刻的端到端耗时。"""
self.record(operation, self._monotonic() - self._started_at, result)
@contextmanager
def activate(self) -> Iterator["TaskPerformanceTrace"]:
"""让同一工作线程内的设备和页面适配层取得本计时对象。"""
token = _ACTIVE_TRACE.set(self)
try:
yield self
finally:
_ACTIVE_TRACE.reset(token)
def _emit(self, operation: str, duration_ms: int, result: str) -> None:
self._sink(
{
"task_id": self._task_id,
"operation": operation,
"duration_ms": duration_ms,
"result": result,
}
)
def current_performance_trace() -> Optional[TaskPerformanceTrace]:
"""返回当前工作线程的任务计时对象;普通调用返回 None。"""
return _ACTIVE_TRACE.get()
+30 -8
View File
@@ -27,6 +27,7 @@ from .purchase_reconcile_service import (
PurchaseReconcileFactory,
PurchaseReconcileService,
)
from .performance_timing import TaskPerformanceTrace
from .task_models import (
NewClaimedTask,
OutboxEventRecord,
@@ -139,6 +140,9 @@ class TaskDispatcher:
cancelled: Callable[[], bool] = lambda: False,
device_connection_checker: Optional[Callable[[str], None]] = None,
task_saved: Callable[[str], None] = lambda _task_id: None,
performance_trace_factory: Callable[[], TaskPerformanceTrace] = (
TaskPerformanceTrace
),
) -> None:
self._gateway = gateway
self._repository = repository
@@ -153,6 +157,7 @@ class TaskDispatcher:
self._reconcile_factory = purchase_reconcile_factory
self._cancelled = cancelled
self._task_saved = task_saved
self._performance_trace_factory = performance_trace_factory
self._device_connection_checker = (
device_connection_checker
or AndroidDeviceService().require_connected
@@ -193,6 +198,15 @@ class TaskDispatcher:
def execute_one(self) -> TaskDispatchOutcome:
"""不可逆采购优先核单,其余按 Outbox、本地任务、Admin 顺序处理。"""
trace = self._performance_trace_factory()
with trace.activate():
return self._execute_one_traced(trace)
def _execute_one_traced(
self, trace: TaskPerformanceTrace
) -> TaskDispatchOutcome:
"""在同一任务计时上下文中串行处理一轮任务。"""
unresolved_reader = getattr(
self._repository, "unresolved_irreversible_purchase", None
)
@@ -211,7 +225,9 @@ class TaskDispatcher:
raise AndroidDeviceSearchError(
"存在只允许核对的采购任务;请先连接并保存原 Android 设备"
)
self._device_connection_checker(self._device_address)
trace.bind_task(reconcile_task.remote_task_id)
with trace.stage("adb_device_check"):
self._device_connection_checker(self._device_address)
if self._reconcile_factory is None:
raise RuntimeError(
f"任务 {reconcile_task.remote_task_id} 只允许核对订单,"
@@ -234,7 +250,8 @@ class TaskDispatcher:
raise AndroidDeviceSearchError(
"请先在设置页选择并保存 Android 设备"
)
self._device_connection_checker(self._device_address)
with trace.stage("adb_device_check"):
self._device_connection_checker(self._device_address)
if not self.purchase_ready:
pending_purchase = self._repository.next_purchase_task()
@@ -248,16 +265,19 @@ class TaskDispatcher:
include_purchase=self.purchase_ready
)
if task is None:
remote = self._gateway.claim_next(
self._client, self.claim_capabilities()
)
with trace.stage("admin_claim_request"):
remote = self._gateway.claim_next(
self._client, self.claim_capabilities()
)
if remote is None:
return TaskDispatchOutcome("no_task", "暂无可领取的任务")
trace.bind_task(remote.task_id)
saved_now = False
try:
self._repository.add_claimed_task(
admin_task_to_new_claimed_task(remote)
)
with trace.stage("sqlite_local_save"):
self._repository.add_claimed_task(
admin_task_to_new_claimed_task(remote)
)
saved_now = True
except DuplicateTaskError:
pass
@@ -273,6 +293,8 @@ class TaskDispatcher:
if saved_now:
self._task_saved(remote.task_id)
trace.bind_task(task.remote_task_id)
if task.task_type is TaskType.COLLECT:
service = CollectTaskService(
self._gateway,