From 2eaea2ee8194d14fa94a00ccb908c091e4daa54d Mon Sep 17 00:00:00 2001 From: simon Date: Thu, 10 Sep 2026 20:54:30 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20cninfo=20=E6=8A=93=E5=8F=96=E9=99=90?= =?UTF-8?q?=E5=88=B6=E6=B5=8F=E8=A7=88=E5=99=A8=E5=B9=B6=E5=8F=91,?= =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E6=95=B4=E6=9C=BA=E5=86=BB=E7=BB=93=20(06:00?= =?UTF-8?q?=20=E4=BB=BB=E5=8A=A1=E5=8E=8B=E7=A9=BF=E5=86=85=E5=AD=98)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因: _render_page 每次新建完整 headless Chromium(实测 857MB/实例), crawl_watchlist 对 15 只股票全量并发 → 峰值需求 ≈12.5GB ≫ 7.9GB RAM, 2026-09-06/08/10 三次 06:0x 整机冻结(load 65 → 看门狗复位)。 - crawler/cninfo.py: 新增 MAX_RENDER_CONCURRENCY(默认 2,env 可配) + 模块级 _render_sem 信号量,渲染全程持槽;顺带清理 3 个死导入 - tests/test_cninfo.py: 新增 8 个(并发上限/串行/异常释放/env 解析) - 实测: 峰值 1796MB / 22 进程 / load 2.75 / 耗时 441s(上限 900s) - 全量 291 passed --- continuation.md | 44 +++++++++++ crawler/cninfo.py | 45 ++++++++++-- scheduler/pipeline.py | 2 +- tests/test_cninfo.py | 167 ++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 251 insertions(+), 7 deletions(-) create mode 100644 tests/test_cninfo.py diff --git a/continuation.md b/continuation.md index 61899e1..4e1dd54 100644 --- a/continuation.md +++ b/continuation.md @@ -4,6 +4,50 @@ --- +## 本次完成 (2026-09-10) — cninfo 抓取压穿内存导致整机冻结的修复 + +**现象**:2026-09-06 / 09-08 / 09-10 连续三次早上 06:0x 整机冻结,看门狗(硬件 2min)硬复位。 + +**根因(证据链闭合)**: +- 三次冻结时刻 = `cninfo` 公告管道 06:00 定时任务:`logs/scheduler_error.log` 显示 09-10 06:00:03.180~.888 **0.7 秒内打印 15 条「抓取…公告」**(15 只股票协程同时进入渲染),随后日志全无直到 06:14:08 重启 +- `data/raw/cninfo/` 缺 20260906/08/10 三个目录(存续目录 mtime 均为 06:02,而落盘在 `_save_items()` 中、抓取全部完成后才执行 → 崩在写数据之前) +- `pcp-pmie` 09-10 06:01:20 报 **load 65**(4 核);DNS 全面超时(frpc/dockerd resolver) +- `journalctl --list-boots` 与三次重启时刻吻合 + +**代码缺陷**(`crawler/cninfo.py`): +1. `_render_page` **每次调用都 `async with AsyncWebCrawler(...)` 新建完整 Chromium**(最贵的错误) +2. `crawl_watchlist` 用 `asyncio.as_completed` 对 watchlist **全量并发**(15 只) +3. **无任何并发限制**;service 亦无资源限制(`CPUQuota=infinity`、`TasksMax=9626`) +4. 本机 `cgroup_disable=memory` → **`MemoryMax` 不可用**(只会 `cpuset cpu io pids`) + +**实测(本机 8GB)**:基线 chrome=0 → **2 并发峰值 1713MB / 20 进程 = 857MB/实例**;外推 15 并发 ≈**12.5GB** ≫ 7.9GB RAM → 必然压穿(时好时坏是 zram 与时序侥幸) + +**修复(第一步:限并发)**: +- `crawler/cninfo.py`:新增 `MAX_RENDER_CONCURRENCY`(默认 **2**,env `CNINFO_RENDER_CONCURRENCY` 可覆盖)+ 模块级 `_render_sem` 信号量,`_render_page` 全程持槽(含浏览器启停) +- 顺带清理 3 个既有死导入(`time`/`datetime`/`Any`) + +**实测验证(2026-09-10 20:43 实跑 `a-share cninfo`)**: + +| 指标 | 修复前(15 并发) | 修复后(2 并发) | +|------|----------------|---------------| +| Chromium 内存峰值 | ≈12.5GB(外推) | **1796 MB** | +| 进程峰值 | ≈150 | **22** | +| load 峰值 | **65** | **2.75** | +| available 最低 | 压穿冻结 | **3561 MB** | +| 耗时 | 崩(无 END) | **441 s**(上限 900s) | +| 结果 | 无数据 | 抓 34 条/存 23 条,补上 09-10 缺失目录 ✅ | + +**测试**:新增 `tests/test_cninfo.py` 8 个(并发上限/串行/异常释放槽位/env 解析);全量 **291 passed**(3 个 crawler 基线失败无关);ruff 干净 + +**遗留与后续优化**: +- 耗时由 127~145s 增至 441s(用时间换内存安全),`cninfo_crawl` 超时 900s 余量由 6 倍降至 2 倍;如嫌慢可 `CNINFO_RENDER_CONCURRENCY=3`(峰值约 2.6GB,仍安全) +- **第二步(未做)**:复用单个 `AsyncWebCrawler` + `arun_many()` 批量渲染(浏览器组 15→1,可同时提速降内存) +- **第三步(未做)**:`_fetch_irm_requests()` 是同步 `requests.get` 却在 async 中直接调用,**阻塞事件循环**;可改 `asyncio.to_thread` +- 长期:改用 cninfo 公告 POST JSON API(`hisAnnouncement/query`)彻底去掉浏览器依赖 +- 未执行 systemd 资源限制:`CPUQuota` 对内存型崩溃基本无效(反而延长驻留),`TasksMax` 过小会让抓取永久失败 + +--- + ## 本次完成 (2026-08-23) — MCP 新闻查询服务确认与修复 **用户需求**:实现新闻查询 MCP 服务(阅读文档步骤,确认是否已实现)。 diff --git a/crawler/cninfo.py b/crawler/cninfo.py index b4f01b4..bb19de7 100644 --- a/crawler/cninfo.py +++ b/crawler/cninfo.py @@ -19,10 +19,8 @@ import hashlib import json import os import re -import time -from datetime import date, datetime, timedelta +from datetime import date, timedelta from pathlib import Path -from typing import Any import requests from bs4 import BeautifulSoup @@ -48,6 +46,34 @@ MAX_IRM_ITEMS = 20 # 请求间隔(秒) REQUEST_DELAY = 0.5 +# 浏览器渲染并发上限。 +# 每次 _render_page 都会启动一个完整的 headless Chromium(实测约 850MB / 10 进程), +# 而 watchlist 是 15 只股票全量并发 → 峰值需求 ≈12.5GB,远超本机 7.9GB RAM, +# 靠 zram swap 侥幸时好时坏,一旦压穿即整机冻结(2026-09-06/08/10 早崩根因)。 +# 限流到 2 后峰值 ≈1.7GB。可用 env CNINFO_RENDER_CONCURRENCY 覆盖。 +DEFAULT_RENDER_CONCURRENCY = 2 + + +def _resolve_render_concurrency() -> int: + """解析渲染并发上限(env CNINFO_RENDER_CONCURRENCY 可覆盖,非法值回退默认)。""" + raw = os.environ.get("CNINFO_RENDER_CONCURRENCY", "") + if not raw.strip(): + return DEFAULT_RENDER_CONCURRENCY + try: + return max(1, int(raw)) + except ValueError: + logger.warning( + "CNINFO_RENDER_CONCURRENCY 非法({!r}),回退默认 {}", + raw, DEFAULT_RENDER_CONCURRENCY, + ) + return DEFAULT_RENDER_CONCURRENCY + + +MAX_RENDER_CONCURRENCY = _resolve_render_concurrency() + +# 渲染信号量:限制同时存活的 Chromium 实例数(防内存压穿) +_render_sem = asyncio.Semaphore(MAX_RENDER_CONCURRENCY) + # --------------------------------------------------------------------------- # # 工具函数 @@ -91,7 +117,11 @@ def _build_stock_url(code: str, org_id: str) -> str: async def _render_page(url: str, timeout_ms: int = 60000, delay_ms: int = 20) -> str: - """用 Crawl4AI 渲染 SPA 页面,返回 HTML 字符串。""" + """用 Crawl4AI 渲染 SPA 页面,返回 HTML 字符串。 + + 通过模块级信号量 _render_sem 限制并发:每次渲染都会启动一个完整的 + headless Chromium(约 850MB),必须限流以免同时驻留过多浏览器压穿内存。 + """ from crawl4ai import AsyncWebCrawler, BrowserConfig, CacheMode, CrawlerRunConfig bconf = BrowserConfig(headless=True, verbose=False) @@ -100,8 +130,11 @@ async def _render_page(url: str, timeout_ms: int = 60000, page_timeout=timeout_ms, delay_before_return_html=delay_ms, ) - async with AsyncWebCrawler(config=bconf) as c: - result = await c.arun(url=url, config=rconf) + # 限流:等待空闲渲染槽位(槽位内包含浏览器启动→渲染→关闭全过程) + async with _render_sem: + logger.debug("获得渲染槽位(并发上限 {}): {}", MAX_RENDER_CONCURRENCY, url[:70]) + async with AsyncWebCrawler(config=bconf) as c: + result = await c.arun(url=url, config=rconf) return getattr(result, "html", "") or "" diff --git a/scheduler/pipeline.py b/scheduler/pipeline.py index 083b6c2..e40a6f5 100644 --- a/scheduler/pipeline.py +++ b/scheduler/pipeline.py @@ -57,7 +57,7 @@ STEP_TIMEOUTS: dict[str, int] = { "embedding": 300, # M5 向量化 "qdrant": 300, # M6 入库(数据量大时需较长时间) "report": 30, # 日报生成+上传 - "cninfo_crawl": 900, # cninfo watchlist URL 驱动(SPA 渲染,每只约 25s) + "cninfo_crawl": 900, # cninfo watchlist URL 驱动(SPA 渲染;限流 2 并发后实测约 440s) } # 步骤对应的 uv run 命令(参数中 {date} 会被替换为实际日期) diff --git a/tests/test_cninfo.py b/tests/test_cninfo.py new file mode 100644 index 0000000..4db204d --- /dev/null +++ b/tests/test_cninfo.py @@ -0,0 +1,167 @@ +"""cninfo 抓取并发限流测试。 + +背景: + 每次 `_render_page` 都会启动一个完整的 headless Chromium(实测约 850MB / + 10 进程),而 `crawl_watchlist` 对 15 只股票全量并发 → 峰值需求 ≈12.5GB, + 远超本机 7.9GB RAM,曾导致 2026-09-06/08/10 三次整机冻结。 + 修复方式:模块级信号量限制同时存活的浏览器数(默认 2)。 + +本测试用假 AsyncWebCrawler 验证并发上限,不启动真实浏览器。 +""" + +from __future__ import annotations + +import asyncio +from types import SimpleNamespace + +import pytest + + +def _install_fake_crawler(monkeypatch: pytest.MonkeyPatch, tracker: dict) -> None: + """把 crawl4ai.AsyncWebCrawler 换成记录并发峰值的假实现。""" + import crawl4ai + + class FakeCrawler: + def __init__(self, config=None) -> None: # noqa: ARG002 + pass + + async def __aenter__(self) -> FakeCrawler: + tracker["current"] += 1 + tracker["max"] = max(tracker["max"], tracker["current"]) + return self + + async def __aexit__(self, *exc: object) -> bool: + tracker["current"] -= 1 + return False + + async def arun(self, url: str, config=None) -> SimpleNamespace: # noqa: ARG002 + # 模拟渲染耗时,制造并发窗口 + await asyncio.sleep(0.05) + return SimpleNamespace(html=f"{url}") + + monkeypatch.setattr(crawl4ai, "AsyncWebCrawler", FakeCrawler) + + +# --------------------------------------------------------------------------- # +# 并发上限解析 +# --------------------------------------------------------------------------- # + +def test_resolve_concurrency_default(monkeypatch: pytest.MonkeyPatch) -> None: + """未设 env 时使用默认值 2。""" + monkeypatch.delenv("CNINFO_RENDER_CONCURRENCY", raising=False) + from crawler.cninfo import DEFAULT_RENDER_CONCURRENCY, _resolve_render_concurrency + + assert DEFAULT_RENDER_CONCURRENCY == 2 + assert _resolve_render_concurrency() == 2 + + +def test_resolve_concurrency_env_override(monkeypatch: pytest.MonkeyPatch) -> None: + """env CNINFO_RENDER_CONCURRENCY 可覆盖。""" + monkeypatch.setenv("CNINFO_RENDER_CONCURRENCY", "4") + from crawler.cninfo import _resolve_render_concurrency + + assert _resolve_render_concurrency() == 4 + + +def test_resolve_concurrency_invalid_falls_back(monkeypatch: pytest.MonkeyPatch) -> None: + """非法值回退默认值,不抛异常。""" + monkeypatch.setenv("CNINFO_RENDER_CONCURRENCY", "abc") + from crawler.cninfo import _resolve_render_concurrency + + assert _resolve_render_concurrency() == 2 + + +def test_resolve_concurrency_clamped_to_one(monkeypatch: pytest.MonkeyPatch) -> None: + """0/负数夹到 1,避免信号量死锁。""" + monkeypatch.setenv("CNINFO_RENDER_CONCURRENCY", "0") + from crawler.cninfo import _resolve_render_concurrency + + assert _resolve_render_concurrency() == 1 + + +# --------------------------------------------------------------------------- # +# 渲染并发限流 +# --------------------------------------------------------------------------- # + +def test_render_concurrency_capped(monkeypatch: pytest.MonkeyPatch) -> None: + """6 个并发渲染请求,同时存活的浏览器数不超过信号量上限 2。""" + from crawler import cninfo + + tracker = {"current": 0, "max": 0} + _install_fake_crawler(monkeypatch, tracker) + # 独立信号量,避免跨测试污染模块级状态 + monkeypatch.setattr(cninfo, "_render_sem", asyncio.Semaphore(2)) + + async def run() -> list[str]: + return await asyncio.gather(*[ + cninfo._render_page(f"https://example.com/{i}") for i in range(6) + ]) + + results = asyncio.run(run()) + assert len(results) == 6, "全部请求都应完成(限流不应丢请求)" + assert tracker["max"] <= 2, f"并发峰值 {tracker['max']} 超过上限 2" + + +def test_render_concurrency_one_is_serial(monkeypatch: pytest.MonkeyPatch) -> None: + """上限为 1 时完全串行执行。""" + from crawler import cninfo + + tracker = {"current": 0, "max": 0} + _install_fake_crawler(monkeypatch, tracker) + monkeypatch.setattr(cninfo, "_render_sem", asyncio.Semaphore(1)) + + async def run() -> list[str]: + return await asyncio.gather(*[ + cninfo._render_page(f"https://example.com/{i}") for i in range(4) + ]) + + results = asyncio.run(run()) + assert len(results) == 4 + assert tracker["max"] == 1, "上限 1 时不应出现并发" + + +def test_render_page_returns_html(monkeypatch: pytest.MonkeyPatch) -> None: + """渲染正常返回 html 字符串。""" + from crawler import cninfo + + tracker = {"current": 0, "max": 0} + _install_fake_crawler(monkeypatch, tracker) + monkeypatch.setattr(cninfo, "_render_sem", asyncio.Semaphore(2)) + + html = asyncio.run(cninfo._render_page("https://example.com/x")) + assert html == "https://example.com/x" + + +def test_render_semaphore_released_on_error(monkeypatch: pytest.MonkeyPatch) -> None: + """渲染抛异常时也要释放槽位(否则后续渲染会死锁)。""" + import crawl4ai + + from crawler import cninfo + + class BoomCrawler: + def __init__(self, config=None) -> None: # noqa: ARG002 + pass + + async def __aenter__(self) -> BoomCrawler: + raise RuntimeError("浏览器启动失败") + + async def __aexit__(self, *exc: object) -> bool: + return False + + async def arun(self, url: str, config=None): # noqa: ARG002 + raise RuntimeError("unreachable") + + monkeypatch.setattr(crawl4ai, "AsyncWebCrawler", BoomCrawler) + monkeypatch.setattr(cninfo, "_render_sem", asyncio.Semaphore(1)) + + async def run() -> None: + # 第一次失败后,槽位必须已释放,第二次才能拿到 + with pytest.raises(RuntimeError): + await cninfo._render_page("https://example.com/fail") + with pytest.raises(RuntimeError): + await asyncio.wait_for( + cninfo._render_page("https://example.com/fail2"), timeout=2 + ) + + asyncio.run(run()) + assert cninfo._render_sem._value == 1, "异常后槽位未释放"