feat(search): log site durations

This commit is contained in:
jxxghp
2026-09-06 16:25:41 +08:00
parent edc15e9575
commit dd7d6db9a1
3 changed files with 89 additions and 13 deletions
+8 -6
View File
@@ -1,6 +1,7 @@
"""搜索分页计划与异步逐页调度 owner。""" """搜索分页计划与异步逐页调度 owner。"""
import asyncio import asyncio
import time
from typing import Any, AsyncIterator, Awaitable, Callable, Dict, List, Optional, Tuple from typing import Any, AsyncIterator, Awaitable, Callable, Dict, List, Optional, Tuple
from app.application.configuration import ( from app.application.configuration import (
@@ -10,7 +11,7 @@ from app.chain.search.contract import _SearchOwnerBase
from app.runtime.stop import runtime_stop_state from app.runtime.stop import runtime_stop_state
from app.runtime.tasks import get_task_registry from app.runtime.tasks import get_task_registry
PageResult = Tuple[List[Any], Optional[Exception]] PageResult = Tuple[List[Any], Optional[Exception], float]
PageTask = asyncio.Task[PageResult] PageTask = asyncio.Task[PageResult]
SiteIndexer = Dict[str, Any] SiteIndexer = Dict[str, Any]
PendingPages = dict[PageTask, Tuple[SiteIndexer, int, int]] PendingPages = dict[PageTask, Tuple[SiteIndexer, int, int]]
@@ -24,14 +25,15 @@ async def _run_site_page(
) -> PageResult: ) -> PageResult:
"""执行一页请求,并把业务异常交回编排层统一传播。""" """执行一页请求,并把业务异常交回编排层统一传播。"""
async with semaphore: async with semaphore:
started_at = time.perf_counter()
try: try:
return await search_page(site, page_number) or [], None return await search_page(site, page_number) or [], None, time.perf_counter() - started_at
except asyncio.CancelledError: except asyncio.CancelledError:
raise raise
except Exception as error: # noqa: BLE001 except Exception as error: # noqa: BLE001
# TaskRegistry 会报告未收口的后台异常;这里由等待方同步传播, # TaskRegistry 会报告未收口的后台异常;这里由等待方同步传播,
# 避免同一 provider 故障同时被登记器和调用链重复记录。 # 避免同一 provider 故障同时被登记器和调用链重复记录。
return [], error return [], error, time.perf_counter() - started_at
def _submit_site_page( def _submit_site_page(
@@ -126,7 +128,7 @@ class SearchPaginationOwner(_SearchOwnerBase):
search_page: Callable[[SiteIndexer, int], Awaitable[Optional[List[Any]]]], search_page: Callable[[SiteIndexer, int], Awaitable[Optional[List[Any]]]],
should_continue: Callable[[SiteIndexer, List[Any]], bool], should_continue: Callable[[SiteIndexer, List[Any]], bool],
task_owner: str, task_owner: str,
) -> AsyncIterator[Tuple[SiteIndexer, int, List[Any], bool]]: ) -> AsyncIterator[Tuple[SiteIndexer, int, List[Any], bool, float]]:
"""统一调度站点逐页请求,并在调用方退出时取消、等待全部请求。""" """统一调度站点逐页请求,并在调用方退出时取消、等待全部请求。"""
total_num = len(indexer_sites) * len(search_pages) total_num = len(indexer_sites) * len(search_pages)
semaphore = asyncio.Semaphore(self.runtime_config.search_threadpool_size or max(1, total_num)) semaphore = asyncio.Semaphore(self.runtime_config.search_threadpool_size or max(1, total_num))
@@ -153,7 +155,7 @@ class SearchPaginationOwner(_SearchOwnerBase):
) )
for task in done_tasks: for task in done_tasks:
site, page_index, page_number = pending_tasks.pop(task) site, page_index, page_number = pending_tasks.pop(task)
page_results, error = await task page_results, error, elapsed = await task
if error is not None: if error is not None:
raise error raise error
continued = should_continue(site, page_results) and page_index + 1 < len(search_pages) continued = should_continue(site, page_results) and page_index + 1 < len(search_pages)
@@ -167,6 +169,6 @@ class SearchPaginationOwner(_SearchOwnerBase):
search_page=search_page, search_page=search_page,
task_owner=task_owner, task_owner=task_owner,
) )
yield site, page_number, page_results, continued yield site, page_number, page_results, continued, elapsed
finally: finally:
await _cancel_pending_pages(pending_tasks) await _cancel_pending_pages(pending_tasks)
+68 -6
View File
@@ -35,6 +35,27 @@ _site_request_schedule_lock = threading.Lock()
_site_next_request_at: Dict[str, float] = {} _site_next_request_at: Dict[str, float] = {}
def _site_duration_key(site: SiteIndexer) -> str:
"""返回站点耗时统计使用的稳定键。"""
site_id = site.get("id")
return str(site_id if site_id is not None else site.get("domain") or site.get("name") or "未知站点")
def _site_duration_label(site: SiteIndexer) -> str:
"""返回站点耗时日志中的可读名称。"""
return str(site.get("name") or site.get("domain") or site.get("id") or "未知站点")
def _format_site_durations(durations: Dict[str, tuple[str, float]]) -> str:
"""格式化各站点耗时,并标出本轮最慢站点。"""
if not durations:
return ""
ordered = sorted(durations.values(), key=lambda item: item[1], reverse=True)
details = "".join(f"{name}{seconds:.2f}" for name, seconds in ordered)
slowest_name, slowest_seconds = ordered[0]
return f";各站点耗时:{details};最慢站点:{slowest_name}{slowest_seconds:.2f} 秒)"
def _site_request_interval(site: SiteIndexer) -> float: def _site_request_interval(site: SiteIndexer) -> float:
"""读取站点管理中已有的单次访问间隔配置。""" """读取站点管理中已有的单次访问间隔配置。"""
try: try:
@@ -90,6 +111,25 @@ def _search_site_page(
) )
def _timed_sync_site_page(
owner: "_SearchProviderSyncOwner",
*,
site: SiteIndexer,
keyword: str,
mtype: Optional[MediaType],
page: int,
) -> tuple[List[TorrentInfo], float]:
"""执行同步站点搜索并返回本页耗时,不包含线程池排队时间。"""
started_at = time.perf_counter()
result = owner._search_site_torrents_with_budget(
site=site,
keyword=keyword,
mtype=mtype,
page=page,
)
return result, time.perf_counter() - started_at
@dataclass(frozen=True) @dataclass(frozen=True)
class ProviderBatch: class ProviderBatch:
"""一次 provider 页完成后发布的统一事实。""" """一次 provider 页完成后发布的统一事实。"""
@@ -101,6 +141,7 @@ class ProviderBatch:
total: int total: int
total_items: int total_items: int
continued: bool continued: bool
elapsed: float
class _SearchProviderSyncOwner(_SearchOwnerBase): class _SearchProviderSyncOwner(_SearchOwnerBase):
@@ -186,7 +227,8 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
"""向进程共享线程 owner 提交一页,并登记该站点的续页位置。""" """向进程共享线程 owner 提交一页,并登记该站点的续页位置。"""
page_number = search_pages[page_index] page_number = search_pages[page_index]
future = ThreadHelper().submit( future = ThreadHelper().submit(
self._search_site_torrents_with_budget, _timed_sync_site_page,
self,
site=site, site=site,
keyword=search_keyword, keyword=search_keyword,
mtype=media_type, mtype=media_type,
@@ -312,6 +354,7 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
self.runtime_config.search_threadpool_size or len(indexer_sites), self.runtime_config.search_threadpool_size or len(indexer_sites),
) )
finish_count = 0 finish_count = 0
site_durations: Dict[str, tuple[str, float]] = {}
for _ in range(max_workers): for _ in range(max_workers):
SearchProviderOwner._submit_next_sync_site( SearchProviderOwner._submit_next_sync_site(
@@ -332,7 +375,14 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
) )
for future in done: for future in done:
site, page_index, page_number = pending.pop(future) site, page_index, page_number = pending.pop(future)
page_results = future.result() or [] page_results, elapsed = future.result()
page_results = page_results or []
duration_key = _site_duration_key(site)
previous = site_durations.get(duration_key)
site_durations[duration_key] = (
_site_duration_label(site),
(previous[1] if previous else 0.0) + elapsed,
)
finish_count += 1 finish_count += 1
results.extend(page_results) results.extend(page_results)
continued = self._should_continue_search_pages( continued = self._should_continue_search_pages(
@@ -368,6 +418,7 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
finally: finally:
for future in pending: for future in pending:
future.cancel() future.cancel()
return site_durations
def _search_all_sites( def _search_all_sites(
self, self,
@@ -402,7 +453,7 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
value=0, value=0,
text=(f"开始搜索,共 {len(indexer_sites)} 个站点,{len(search_pages)} 页 ..."), text=(f"开始搜索,共 {len(indexer_sites)} 个站点,{len(search_pages)} 页 ..."),
) )
SearchProviderOwner._collect_sync_site_results( site_durations = SearchProviderOwner._collect_sync_site_results(
self, self,
keyword=keyword, keyword=keyword,
indexer_sites=indexer_sites, indexer_sites=indexer_sites,
@@ -423,7 +474,10 @@ class _SearchProviderSyncOwner(_SearchOwnerBase):
if isinstance(getattr(self, "_subscription_site_budget", None), SubscriptionSiteBudget) if isinstance(getattr(self, "_subscription_site_budget", None), SubscriptionSiteBudget)
else logger.info else logger.info
) )
log(f"站点搜索完成,有效资源数:{len(results)},总耗时 {elapsed}") log(
f"站点搜索完成,有效资源数:{len(results)},总耗时 {elapsed}"
f"{_format_site_durations(site_durations)}"
)
return results return results
finally: finally:
progress.end() progress.end()
@@ -457,7 +511,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
task_owner=task_owner, task_owner=task_owner,
) )
async with aclosing(page_iterator): async with aclosing(page_iterator):
async for site, page_number, page_results, continued in page_iterator: async for site, page_number, page_results, continued, elapsed in page_iterator:
finish_count += 1 finish_count += 1
total_items += len(page_results) total_items += len(page_results)
yield ProviderBatch( yield ProviderBatch(
@@ -468,6 +522,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
total=total, total=total,
total_items=total_items, total_items=total_items,
continued=continued, continued=continued,
elapsed=elapsed,
) )
@staticmethod @staticmethod
@@ -615,6 +670,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
started_at = datetime.now() started_at = datetime.now()
last_total_items = len(initial_items) last_total_items = len(initial_items)
site_durations: Dict[str, tuple[str, float]] = {}
# 既有流式协议约定插件资源先于站点进度发布;列表入口也消费同一事件流。 # 既有流式协议约定插件资源先于站点进度发布;列表入口也消费同一事件流。
if initial_items: if initial_items:
yield SearchProviderOwner._plugin_event( yield SearchProviderOwner._plugin_event(
@@ -646,6 +702,12 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
async with aclosing(batches): async with aclosing(batches):
async for batch in batches: async for batch in batches:
last_total_items = batch.total_items last_total_items = batch.total_items
duration_key = _site_duration_key(batch.site)
previous = site_durations.get(duration_key)
site_durations[duration_key] = (
_site_duration_label(batch.site),
(previous[1] if previous else 0.0) + batch.elapsed,
)
yield await SearchProviderOwner._batch_event( yield await SearchProviderOwner._batch_event(
batch=batch, batch=batch,
keyword=keyword, keyword=keyword,
@@ -658,7 +720,7 @@ class SearchProviderOwner(_SearchProviderSyncOwner):
f"站点{label}搜索完成,有效{'字幕' if subtitle else '资源'}数:{last_total_items},总耗时 {elapsed}" f"站点{label}搜索完成,有效{'字幕' if subtitle else '资源'}数:{last_total_items},总耗时 {elapsed}"
) )
await progress.update(value=100, text=done_text) await progress.update(value=100, text=done_text)
logger.info(done_text) logger.info(f"{done_text}{_format_site_durations(site_durations)}")
finally: finally:
await progress.end() await progress.end()
+13 -1
View File
@@ -32,6 +32,16 @@ def _make_chain() -> SearchChain:
return chain return chain
def test_format_site_durations_orders_sites_and_marks_slowest() -> None:
"""站点耗时日志应按耗时倒序,并明确标出最慢站点。"""
formatted = search_module._format_site_durations({
"1": ("站点一", 1.2),
"2": ("站点二", 3.4),
})
assert formatted == ";各站点耗时:站点二:3.40 秒,站点一:1.20 秒;最慢站点:站点二(3.40 秒)"
@pytest.mark.asyncio @pytest.mark.asyncio
async def test_site_page_iterator_cancels_and_waits_pending_requests() -> None: async def test_site_page_iterator_cancels_and_waits_pending_requests() -> None:
"""调用方提前关闭迭代器时,统一编排器必须收口其他站点请求。""" """调用方提前关闭迭代器时,统一编排器必须收口其他站点请求。"""
@@ -61,7 +71,9 @@ async def test_site_page_iterator_cancels_and_waits_pending_requests() -> None:
task_owner="test.search.site_page", task_owner="test.search.site_page",
) )
assert await anext(iterator) == ({"id": 1}, 0, ["ready"], False) result = await anext(iterator)
assert result[:4] == ({"id": 1}, 0, ["ready"], False)
assert result[4] >= 0
await blocked_started.wait() await blocked_started.wait()
await iterator.aclose() await iterator.aclose()