Commit aeb29b90 authored by wangteng's avatar wangteng

性能优化(二)流式扫描:OFFSET 深分页 → 后端单连接流式扫描 + /pull 轮询

旧方案每页独立连接跑 LIMIT/OFFSET,页越深扫得越多(50万行=1000次建连+
1000次全表扫描排序)。改为 /start 起 daemon 线程单连接 iter_rows 流式扫
全表一遍(保留 ORDER BY 首列,全表只排序一次),命中行进 session 缓冲区,
前端由 3-worker 分页队列改为轮询 /pull 增量拉取(满额立即续拉)。

后端:
- queries.py:_scan_worker 扫描线程(cancelled 无锁读+关连接双保险、每
  1000 行发布进度、命中超 10 万条截断只计数);/pull 游标切片;/cancel
  摘 session+request_cancel;删 /page;_load_task_compiled 批量化(N+1→
  固定 4 条查询)+ 接入 compile_rule 预编译(regex 每行 re.compile、code
  类每行 compile+exec 的量级浪费消除);/start 的 COUNT 与字段注释合并
  单连接;并发扫描软上限 MAX_ACTIVE_SCANS=4
- session_manager.py:QuerySession 扩展流式扫描状态(scanned/bad_rows/
  scan_done/scan_error/truncated,统一持锁读写),request_cancel() 幂等
  取消(置标志+关连接),TTL 清理锁外停掉过期扫描线程
- db_adapter.py:DBConnection.close() 抽为 public 幂等(跨线程取消用);
  iter_rows finally 的 cur.close 包保护(取消关连接后 generator 提前退出)

前端:
- queries.js:fetchPage → pullQuery(after 游标)
- DataQualityView.vue:pollLoop 轮询器(404 分流本端取消/会话过期、网络
  异常 2s 退避×3、done 定稿 toast);进度条改「已扫描 X/Y 条(百分比)」;
  快照/续跑按 __row_index 去重适配;无规则任务提示修正(原误报"无数据")
- client.js:错误对象挂 HTTP status

实测(37151 行任务):扫描 12.5s→0.9s;扫描期间 /api/health 均值
114ms→2.2ms;code 类规则求值 175x;命中结果与旧版逐字节一致。
Co-Authored-By: default avatarClaude Fable 5 <noreply@anthropic.com>
parent f7446793
......@@ -612,9 +612,17 @@ class DBConnection:
return self
def __exit__(self, exc_type, exc, tb):
if self._conn is not None:
self.close()
def close(self) -> None:
"""关闭连接(幂等)。2026-09-24:从 __exit__ 抽成 public —— 流式扫描的
取消路径要跨线程关连接(/cancel 或 TTL 清理在别的线程调用),与 worker
自己的 with 退出双重关闭是常态,必须幂等。
"""
conn, self._conn = self._conn, None
if conn is not None:
try:
self._conn.close()
conn.close()
logger.debug("[DB] 连接已关闭")
except Exception as e:
logger.warning(f"[DB] 关闭连接异常: {e}")
......@@ -682,7 +690,13 @@ class DBConnection:
finally:
elapsed_ms = (time.monotonic() - started) * 1000
logger.debug(f"[SQL] iter_rows {n} rows ({elapsed_ms:.1f}ms): {sql[:120]}...")
# 2026-09-24:连接可能已被取消线程关闭(流式扫描的 /cancel 路径),
# generator 提前退出走到这里时 cur.close() 会对已关连接抛异常、
# 掩盖原退出路径 —— 包一层按原计划收尾。
try:
cur.close()
except Exception:
pass
def fetchone(self, sql: str, params: tuple | None = None) -> Optional[dict]:
cur = self._cursor()
......
"""查询 session 管理器(web 自包含版)
为每次 ``POST /api/queries/start`` 在内存里建一个 ``QuerySession``,记录任务配置 +
DB 连接 + 分页参数的快照。前端按页号去 ``POST /api/queries/page`` 拉一页结果,可并发
拉多页,并主动 ``POST /api/queries/cancel`` 终止。
编译好的规则 + 流式扫描状态。后端起 daemon 线程单连接流式扫全表(2026-09-24 起,
替代逐页 LIMIT/OFFSET 深分页),前端经 ``POST /api/queries/pull`` 增量拉取命中行 +
进度,并主动 ``POST /api/queries/cancel`` 终止。
生命周期:
创建 ``/start`` 跑完 COUNT 后建 session 入 ``_REGISTRY``,UUID 作 token 返回
取消 ``/cancel`` 从 ``_REGISTRY`` pop 掉,后续 ``/page`` 立刻 409 短路
TTL 每次 ``/start`` 入口懒清理 ``>30min`` 的 session(无后台线程,简单即正义)
创建 ``/start`` 跑完 COUNT 后建 session 入 ``_REGISTRY``(先 put 后起扫描线程),
UUID 作 token 返回
拉取 ``/pull`` 按 after 游标从 bad_rows 缓冲区增量切片返回(不摘 session)
取消 ``/cancel`` 从 ``_REGISTRY`` pop 掉 + request_cancel()(置标志 + 关扫描
连接),后续 ``/pull`` 立刻 404 短路
TTL 每次 ``/start`` 入口懒清理 ``>30min`` 的 session(无后台线程,简单即正义);
过期 session 同样 request_cancel(),停掉可能还在扫的线程
隔离粒度:**per-client** 而非 **per-user**。
当前无 auth 模块 → 任何人拿到 ``session_id`` 都能拉/取消对应 session;
......@@ -22,10 +27,13 @@ from __future__ import annotations
import threading
import time
from dataclasses import dataclass, field
from typing import Optional
from typing import Any, Callable, Optional
from web.backend.core.db_adapter import DBConfig
# 编译期产出的规则校验闭包:check(value) -> bool(rule_runner.compile_rule 的返回类型)
RuleCheckFn = Callable[[Any], bool]
# ── 模块级存储 ────────────────────────────────────────────
# 改 per-user 时记得把 get/pop 都接受 user_id 校验。
......@@ -43,53 +51,87 @@ class CompiledFieldRules:
字段大小写已归一为小写(与列名对齐)。
"""
field_key: str
# 每条规则:(rule_type, regex_or_None, code_or_None, desc, skip_null, name, rule_id)
# 2026-08-21:追加 skip_null(值为空时是否跳过)
# 2026-08-24:追加 name(规则名称)→ 给 _evaluate_row 写到 issues 里,
# 前端「数据明细」按字段+规则名展示「违反哪条规则」
# 2026-08-24:再追加 rule_id(Rule ORM 主键)→ 前端「点 ! 问 LLM」时用来
# 在 col.ruleList[] 里精确定位原规则(含 code/regex/rule_type)
rules: list[tuple[str, Optional[str], Optional[str], str, bool, str, int]] = field(default_factory=list)
# 每条规则:(rule_type, check_fn_or_None, desc, name, rule_id)
# 2026-09-24:预编译改造 —— check_fn 是 rule_runner.compile_rule() 产出的
# check(value) -> bool 闭包(regex 存编译好的 Pattern,number/date/string 存
# exec 一次后取出的用户 check 函数;skip_null 短路已烘焙进去)。
# 之前存 regex/code 字符串、每行重复 compile+exec(含重建沙箱命名空间),
# 全表扫描量级差距巨大。
# None = 编译失败的坏规则:_evaluate_row 对每行判不合规并注明「规则编译失败」。
# desc:规则说明(issues 展示用);name:规则名称(2026-08-24 加,前端展示
# 「违反哪条规则」);rule_id:Rule ORM 主键(2026-08-24 加,前端「点 ! 问
# LLM」时在 col.ruleList[] 里精确定位原规则)。
rules: list[tuple[str, Optional["RuleCheckFn"], str, str, int]] = field(default_factory=list)
@dataclass
class QuerySession:
"""一次查询的所有上下文。
注意:
``cancelled`` 由 ``/cancel`` 置 True 并从 ``_REGISTRY`` 摘掉;
但 ``/page`` 在并发飞过来时仍可能拿到 in-flight 副本,所以 ``/page`` 入口要
二次检查 ``_REGISTRY`` 里是否还在(不在就 404),以及 session 上的 cancelled。
2026-09-24 流式扫描改造:session 不再被 /page 逐页引用,而是绑定一个后台扫描
线程(daemon)单连接流式扫全表,结果缓存在 bad_rows 缓冲区,前端经 /pull
增量拉取。生命周期:
/start 建 session 入 _REGISTRY + 起扫描线程(先 put 后 start)
/pull 按 after 游标切缓冲区增量返回 + 进度计数快照
/cancel request_cancel():置 cancelled + 关扫描连接(worker 批次间退出)
TTL 30min 懒清理;过期 session 同样 request_cancel()(停掉泄漏线程)
加锁纪律(重要):
- bad_rows / bad_rows_total / scanned / truncated / scan_done / scan_error /
scan_conn 一律持 self.lock 写;/pull 持锁切片(ms 级,与 worker 抢锁无竞争)
- cancelled 锁内置位;扫描 worker 循环里**无锁读**(GIL 下 bool 读原子,
只为尽快退出,读到旧值最多多扫一批)
"""
session_id: str
task_id: int
db_config: DBConfig
db_type: str
page_size: int
total_rows: int
total_pages: int
select_sql_base: str # SELECT cols FROM table ORDER BY first_field
compiled: list[CompiledFieldRules] = field(default_factory=list)
col_names: list[str] = field(default_factory=list)
field_list_out: list[dict] = field(default_factory=list) # 给前端表头用
cancelled: bool = False
pages_scanned: int = 0
bad_rows_total: int = 0
created_at: float = field(default_factory=time.time)
lock: threading.Lock = field(default_factory=threading.Lock)
def mark_page_done(self, bad_delta: int) -> None:
# ── 流式扫描状态(2026-09-24,除 scan_thread 外全部经 self.lock 读写)──
scanned: int = 0 # 已扫行数(worker 每 SCAN_PUBLISH_EVERY 行发布一次)
bad_rows: list[dict] = field(default_factory=list) # 命中行缓冲区(只追加;超上限停止收集)
bad_rows_total: int = 0 # 命中总数(含被截断丢弃的部分)
truncated: bool = False # 缓冲区超上限后置 True(前端提示「仅展示前 N 条」)
scan_done: bool = False # worker 终止标志(成功/异常/取消都置)
scan_error: Optional[str] = None # 非取消原因的异常描述
scan_conn: Optional[object] = None # worker 打开的 DBConnection 引用(取消用;避免循环 import 用 object 注解)
scan_thread: Optional[threading.Thread] = None
def request_cancel(self) -> None:
"""置取消标志 + best-effort 关闭扫描连接(幂等;/cancel 与 TTL 清理共用)。
关闭连接是为了打断阻塞中的 DB 读取(MySQL SSDictCursor 会被 close 唤醒);
连不上也无妨 —— worker 在下一个批次边界看到 cancelled 退出。
"""
with self.lock:
self.pages_scanned += 1
self.bad_rows_total += bad_delta
self.cancelled = True
conn = self.scan_conn
if conn is not None:
try:
conn.close()
except Exception:
pass
# ── CRUD ──────────────────────────────────────────────────
def put(session: QuerySession) -> None:
"""插入 session,入口处顺便清一遍过期 session"""
"""插入 session,入口处顺便清一遍过期 session(含停掉其 in-flight 扫描线程)"""
with _REGISTRY_LOCK:
_purge_expired_locked()
expired = _purge_expired_locked()
_REGISTRY[session.session_id] = session
# 锁外取消:request_cancel 内部会抢 session.lock + 关 DB 连接,
# 放在 registry 锁内一旦与 /pull(持 session.lock)交叠容易拉长临界区
for s in expired:
s.request_cancel()
def get(session_id: str) -> Optional[QuerySession]:
......@@ -109,14 +151,25 @@ def size() -> int:
return len(_REGISTRY)
def active_scan_count() -> int:
"""还在扫描中的 session 数(/start 的并发软上限用,防连点/泄漏)。
判定 = 未 done 且未取消;对 scan_done 的无锁读在「worker 正要置 done」的
窗口里可能多算一个,作为软上限的误差可接受。
"""
with _REGISTRY_LOCK:
return sum(1 for s in _REGISTRY.values() if not s.scan_done and not s.cancelled)
# ── 内部 ──────────────────────────────────────────────────
def _purge_expired_locked() -> None:
"""把超过 TTL 的 session 删掉。调用方必须已持有 _REGISTRY_LOCK。
def _purge_expired_locked() -> list[QuerySession]:
"""把超过 TTL 的 session 摘掉,返回被摘的 session 列表(调用方锁外 request_cancel)。
Lazy cleanup —— 没有后台线程;只在 /start 入口触发一次。
单进程 uvicorn 这个开销可忽略(session 数个位数)。
"""
cutoff = time.time() - _TTL_SECONDS
expired = [sid for sid, s in _REGISTRY.items() if s.created_at < cutoff]
for sid in expired:
_REGISTRY.pop(sid, None)
\ No newline at end of file
expired = [s for sid, s in _REGISTRY.items() if s.created_at < cutoff]
for s in expired:
_REGISTRY.pop(s.session_id, None)
return expired
\ No newline at end of file
"""web 任务查询 + 校验 API(分页版)
"""web 任务查询 + 校验 API(流式扫描版,2026-09-24 重构)
端点:
POST /api/queries/start 创建查询 session:COUNT → 编译规则 → 入 session_manager
POST /api/queries/page 按页号拉一页,跑规则,只返回不命中行的(每页独立连接)
POST /api/queries/cancel 摘掉 session(best-effort 中断)
POST /api/queries/start 创建查询 session:COUNT → 编译规则 → 起后台扫描线程
POST /api/queries/pull 按 after 游标增量拉取命中行 + 进度(前端轮询)
POST /api/queries/cancel 摘掉 session + 关扫描连接(best-effort 中断)
设计:
- 复用 web.backend.core.db_adapter 建连(每页新连接,简单)
- 复用 web.backend.core.db_adapter 建连
- 复用 routers/db._normalize_db_type 做中文 label → 内部小写
- 列名按 task.field_list 显式 SELECT(不 SELECT * 防止暴露未配置列)
- 校验引擎:rule_runner.run_rule() 按 rule_type dispatch
- regex : re.compile(regex).search(value),空 regex 回退到 `^.+$`(非空校验)
- number/date : 用户提供的 def check(value) -> bool,安全沙箱 exec
- 校验引擎:rule_runner.compile_rule() 在 /start 编译一次 → _evaluate_row 逐行调闭包
- regex : 编译好的 Pattern.search(value),空 regex 回退到 `^.+$`(非空校验)
- number/date/string : exec 一次取出用户 check 函数,逐行只调用
判定口径(只返回「不合规」的行):
- 字段级:该字段的多条规则是 AND —— 任一条不满足,该字段即不合规,
且所有未通过的规则都会进 issues(备注列列全原因),不止第一条
- 行级:任一字段不合规 → 整行不合规 → 返回;全部字段都合规的行直接丢弃
分页策略:
- SQL 不带 LIMIT,由 paginate_sql(base, db_type, offset, limit) 拼方言
- ORDER BY 必加,否则 OFFSET 语义不一致(MySQL 有默认顺序,达梦/Oracle 无)
- 每页新开 DBConnection → fetchall(不用 iter_rows,页内是有限结果集)
- 每页 ~50ms 重连开销 vs 长连接 cursor 生命周期管理复杂度,选前者
流式扫描策略(替代 2026-08-21 的逐页 LIMIT/OFFSET):
- 旧方案每页独立连接 + OFFSET 深分页:第 N 页要跳过 (N-1)*page_size 行,
页越深扫得越多,且 ORDER BY 首列通常无索引每页都全表排序 ——
50 万行 = 1000 次建连 + 1000 次扫描排序(页均扫 25 万行)
- 新方案 /start 起一个 daemon 线程,单连接 iter_rows 流式扫全表一遍
(MySQL 走 SSDictCursor 真流式;达梦/Oracle 服务端游标 fetchmany),
ORDER BY 首列保留(全表只排序一次,结果顺序与旧版一致),命中行进
session.bad_rows 缓冲区,前端经 /pull 增量拉取
- 中断:cancelled 标志 + 关闭连接双保险(见 session_manager.request_cancel)
Session 隔离:
- per-client 而非 per-user(当前无 auth,详见 session_manager.py 注释)
- UUID 作 token,TTL 30min(懒清理)
- UUID 作 token,TTL 30min(懒清理,含停掉过期扫描线程)
"""
from __future__ import annotations
import re
import threading
import time
import uuid
from datetime import datetime, timezone
......@@ -41,11 +45,12 @@ from pydantic import BaseModel, Field as PydField
from sqlalchemy.orm import Session
from web.backend._logging import get_logger
from web.backend.core.db_adapter import DBConfig, DBConnection, paginate_sql, quote_ident
from web.backend.core.rule_runner import RuleRunError, run_rule
from web.backend.core.db_adapter import DBConfig, DBConnection, quote_ident
from web.backend.core.rule_runner import RuleRunError, compile_rule
from web.backend.core.session_manager import (
CompiledFieldRules,
QuerySession,
active_scan_count,
get as sm_get,
pop as sm_pop,
put as sm_put,
......@@ -61,13 +66,19 @@ from web.backend.routers.tasks import _parse_conn_json
router = APIRouter(prefix="/queries", tags=["queries"])
logger = get_logger("backend.routers.queries")
# ── 流式扫描参数(2026-09-24)──────────────────────────────
MAX_BUFFERED_BAD_ROWS = 100_000 # 命中行缓冲区上限(防「全部不合规」时内存失控),超限只计数
PULL_MAX_ROWS = 5000 # 单次 /pull 返回行数上限(前端满额应立即续拉)
SCAN_PUBLISH_EVERY = 1000 # 扫描线程每 N 行持锁发布一次 scanned(进度条粒度)
MAX_ACTIVE_SCANS = 4 # 同时扫描的 session 软上限(防重复 /start 泄漏线程)
# ── Pydantic ──────────────────────────────────────────────
class StartQueryRequest(BaseModel):
task_id: int = PydField(..., description="任务 id(后端从 DB 读全部配置)")
page_size: int = PydField(
500, ge=1, le=10000,
description="每页行数(默认 500);前端并发拉页,受单页大小影响",
description="兼容保留:流式扫描版不再分页,此字段被忽略",
)
......@@ -82,23 +93,23 @@ class StartQueryResponse(BaseModel):
field_list: list[dict] = PydField(default_factory=list, description="顺带回当前任务的字段+规则数")
class PageQueryRequest(BaseModel):
class PullQueryRequest(BaseModel):
session_id: str = PydField(..., description="/start 返回的 session token")
page_no: int = PydField(..., ge=1, description="1-based 页号")
after: int = PydField(0, ge=0, description="已消费的 bad_rows 数(游标),首次传 0")
class PageQueryResponse(BaseModel):
class PullQueryResponse(BaseModel):
ok: bool
message: Optional[str] = None
session_id: str
page_no: int
page_size: int
row_start: int # (page_no-1)*page_size + 1
bad_rows: list[dict] = PydField(default_factory=list)
scanned_delta: int = 0 # 本页扫过的行数(应 = page_size,最后一页可能更少)
bad_delta: int = 0 # 本页命中的不合规行数
cancelled: bool = False # 当前 session 是否已取消(前端据此停拉剩余页)
done: bool = False # 本页是最后一页
bad_rows: list[dict] = PydField(default_factory=list) # bad_rows[after : after+PULL_MAX_ROWS]
next_after: int # after + len(bad_rows),下次请求带上
scanned: int # 已扫描行数(进度分子)
total_rows: int # COUNT 结果(进度分母)
bad_total: int # 命中总数(含被截断丢弃部分)
done: bool # 扫描线程已终止(成功/异常/取消)
error: Optional[str] = None # 非取消原因的扫描异常
cancelled: bool # 已被 /cancel(前端据此静默退出)
truncated: bool # 命中超缓冲区上限(前端提示仅展示前 N 条)
class CancelQueryRequest(BaseModel):
......@@ -108,8 +119,8 @@ class CancelQueryRequest(BaseModel):
class CancelQueryResponse(BaseModel):
ok: bool
cancelled: bool
pages_scanned: int
bad_rows_so_far: int
scanned: int = 0
bad_rows_so_far: int = 0
# ── 公共工具 ──────────────────────────────────────────────
......@@ -129,13 +140,42 @@ def _load_task_compiled(db: Session, task_id: int):
fields = db.query(Field).filter(Field.task_id == task.id).order_by(Field.ord).all()
# 新任务优先使用独立规则库关联;没有关联的历史任务继续读取旧 rule 表。
rules_by_field = {}
for f in fields:
links = db.query(FieldRule).filter(FieldRule.field_id == f.id).order_by(FieldRule.ord).all()
library_rules = [db.get(RuleLibrary, link.rule_library_id) for link in links]
rules_by_field[f.id] = [r for r in library_rules if r] or (
db.query(Rule).filter(Rule.field_id == f.id).order_by(Rule.ord).all()
# 2026-09-24:批量化 —— 之前每字段 2~3 次查询(N+1:FieldRule 列表 + 逐条
# db.get(RuleLibrary) + Rule 兜底),50 字段任务就是 150+ 次 SQLite 往返;
# 现在固定 4 条查询(fields / FieldRule 全量 / RuleLibrary 批量 / Rule 兜底批量),
# Python 侧分组,语义不变:有 FieldRule 链接 → 按 ord 取 RuleLibrary;否则回退老 Rule。
field_ids = [f.id for f in fields]
links: list[FieldRule] = []
if field_ids:
links = (
db.query(FieldRule)
.filter(FieldRule.field_id.in_(field_ids))
.order_by(FieldRule.field_id, FieldRule.ord)
.all()
)
lib_ids = sorted({l.rule_library_id for l in links})
libs_by_id: dict[int, RuleLibrary] = {}
if lib_ids:
libs_by_id = {x.id: x for x in db.query(RuleLibrary).filter(RuleLibrary.id.in_(lib_ids)).all()}
links_by_field: dict[int, list[FieldRule]] = {}
for l in links:
links_by_field.setdefault(l.field_id, []).append(l)
# 只对「没有任何 FieldRule 链接」的字段查老 Rule 表(保持旧回退语义)
legacy_field_ids = [fid for fid in field_ids if fid not in links_by_field]
legacy_rules_by_field: dict[int, list[Rule]] = {}
if legacy_field_ids:
for r in (
db.query(Rule)
.filter(Rule.field_id.in_(legacy_field_ids))
.order_by(Rule.field_id, Rule.ord)
.all()
):
legacy_rules_by_field.setdefault(r.field_id, []).append(r)
rules_by_field: dict[int, list] = {}
for f in fields:
library_rules = [libs_by_id[l.rule_library_id] for l in links_by_field.get(f.id, [])]
library_rules = [r for r in library_rules if r]
rules_by_field[f.id] = library_rules or legacy_rules_by_field.get(f.id, [])
fields_with_rules: list[tuple[Field, list[Rule]]] = [
(f, rules_by_field[f.id]) for f in fields if rules_by_field[f.id]
]
......@@ -165,30 +205,33 @@ def _load_task_compiled(db: Session, task_id: int):
for f in fields
]
# 编译:regex 提前 compile;number/date/string 不编译(每行 exec)
# 编译:2026-09-24 起全部规则在 /start 时经 rule_runner.compile_rule() 编译一次
# (regex → 编译好的 Pattern;number/date/string → exec 一次取出 check 闭包),
# _evaluate_row 逐行只调闭包 —— 之前 regex 每行 re.compile、code 类每行
# compile+exec 重建沙箱命名空间,全表扫描量级浪费。
# 注意:r.rule_type / r.desc / r.regex / r.code 都是 SQLAlchemy Mapped 字段,
# 运行时是 str,但静态类型是 InstrumentedAttribute —— 用 str() 包一下让 Pylance 也满意
# 2026-08-21:snapshot 追加 skip_null(值为空时是否跳过该规则)
# 2026-08-24:snapshot 再追加 name(规则名称)→ _evaluate_row 写到 issues → 数据明细展示
# snapshot 元组:(rule_type, check_fn_or_None, desc, name, rule_id)
# check_fn=None = 编译失败的坏规则 → 每行判不合规并注明「规则编译失败」
compiled_fields: list[CompiledFieldRules] = []
for f, rules in fields_with_rules:
snapshots: list[tuple[str, Optional[str], Optional[str], str, bool, str, int]] = []
snapshots: list[tuple[str, Optional[object], str, str, int]] = []
for r in rules:
rt = str(r.rule_type or "regex")
desc = str(getattr(r, "description", "") or getattr(r, "desc", "") or "")
skip_null = bool(r.skip_null)
name = str(r.name or "")
rule_id = int(r.id or 0)
if rt == "regex":
regex_src = str(r.regex or "").strip() or r"^.+$"
try:
compiled_pat = re.compile(regex_src)
snapshots.append((rt, compiled_pat.pattern, None, desc, skip_null, name, rule_id))
except re.error as e:
logger.warning(f"[queries] 规则 id={r.id} regex 编译失败:{e}")
snapshots.append((rt, None, None, desc, skip_null, name, rule_id)) # None 标记"坏规则"
else:
snapshots.append((rt, None, str(r.code or ""), desc, skip_null, name, rule_id))
fn = compile_rule(
rt,
regex=(str(r.regex or "") if rt == "regex" else None),
code=(str(r.code or "") if rt != "regex" else None),
skip_null=bool(r.skip_null),
)
snapshots.append((rt, fn, desc, name, rule_id))
except RuleRunError as e:
logger.warning(f"[queries] 规则 id={r.id} 编译失败:{e}")
snapshots.append((rt, None, desc, name, rule_id)) # None 标记"坏规则"
compiled_fields.append(CompiledFieldRules(field_key=f.field_key.lower(), rules=snapshots))
# SELECT 列名:按 ord 升序的去重小写列表
......@@ -228,41 +271,29 @@ def _build_select_sql(cfg: DBConfig, col_names: list[str], table: str) -> str:
def _evaluate_row(compiled_fields: list[CompiledFieldRules], raw: dict[str, Any]) -> tuple[list[str], list[dict]]:
"""对一行业务数据跑所有规则,返 (error_cells, issues)。
2026-09-24:改为调 /start 时预编译好的 check 闭包(rule_runner.compile_rule),
本函数不再做任何 compile/exec。skip_null 短路与 date 归一化已烘焙进闭包。
error_cells: 不合规字段 key 列表(去重保序)
issues: 每条不通过的规则的 {field, name, desc}
- name: 2026-08-24 新增,规则名称(前端展示「违反哪条规则」用);
老规则空串时由 UI 用 desc 兜底
- desc: 规则说明(自然语言;执行失败时是「规则执行失败:xxx」)
- desc: 规则说明(自然语言;编译失败时是「规则编译失败:xxx」、
执行失败时是「规则执行失败:xxx」)
"""
issues: list[dict] = []
error_cells: list[str] = []
for cf in compiled_fields:
val = raw.get(cf.field_key)
# 2026-08-21:字段级「跳过空值」开关(per-field,非 per-rule)
# —— 检测引擎目前是 per-rule;如果未来要 per-field,把 skip_null 上提到 Field 模型
# 现在是 per-rule:跑每条规则时单独判断
for rule_type, compiled_regex, code, desc, skip_null, name, rule_id in cf.rules:
passed: Optional[bool] = None
fail_desc: Optional[str] = None
if rule_type == "regex":
if compiled_regex is None:
for rule_type, fn, desc, name, rule_id in cf.rules:
if fn is None:
# 坏规则(/start 编译失败):每行判不合规,让用户在结果里看到原因
passed = False
fail_desc = f"规则正则编译失败:{desc}"
elif skip_null and (val is None or val == ""):
# 2026-08-21:跳过空值开关 —— 空值直接合规,跳过正则
passed = True
else:
pat = re.compile(compiled_regex)
val_str = "" if val is None else str(val)
passed = pat.search(val_str) is not None
if not passed:
fail_desc = desc
fail_desc = f"规则编译失败:{desc}"
else:
try:
# 2026-08-21:number/date/string 类型把 skip_null 传给 run_rule,由沙箱入口短路
passed = run_rule(rule_type, None, code, val, skip_null=skip_null)
if not passed:
fail_desc = desc
passed = fn(val)
fail_desc = None if passed else desc
except RuleRunError as e:
passed = False
fail_desc = f"规则执行失败:{e}"
......@@ -281,9 +312,72 @@ def _evaluate_row(compiled_fields: list[CompiledFieldRules], raw: dict[str, Any]
return error_cells, issues
# ── 扫描线程(2026-09-24 流式扫描核心)────────────────────
def _scan_worker(sess: QuerySession) -> None:
"""daemon 线程:单连接 iter_rows 流式扫全表 + 逐行规则求值 + 命中行入缓冲区。
- ORDER BY 首列保留:全表只排序一次,结果顺序与旧分页版一致
- cancelled 无锁读(GIL 下 bool 原子),每行检查、批次间最快退出;
request_cancel() 还会直接关连接打断阻塞中的读取(双保险)
- 命中行 append 进 sess.bad_rows(持锁);超过 MAX_BUFFERED_BAD_ROWS 后
只计数不收集(sess.truncated=True),防「全部不合规」时内存失控
- 任何非取消异常 → sess.scan_error;取消引发的连接异常不算错误
"""
row_index = 0
scan_error: Optional[str] = None
started = time.monotonic()
try:
with DBConnection(sess.db_config) as dbc:
with sess.lock:
if sess.cancelled:
return # 启动前就被取消(/cancel 抢在线程 run 之前)
sess.scan_conn = dbc # 暴露给 request_cancel(跨线程关连接)
for raw in dbc.iter_rows(sess.select_sql_base):
if sess.cancelled:
break
row_index += 1
err_cells, issues = _evaluate_row(sess.compiled, raw)
if issues:
row_out = dict(raw)
row_out["errorCells"] = err_cells
row_out["issues"] = issues
# 全局行号(0-based 扫描序):前端做 row-key(append 模式必用)
row_out["__row_index"] = row_index - 1
with sess.lock:
sess.bad_rows_total += 1
if len(sess.bad_rows) < MAX_BUFFERED_BAD_ROWS:
sess.bad_rows.append(row_out)
else:
sess.truncated = True
if row_index % SCAN_PUBLISH_EVERY == 0:
with sess.lock:
sess.scanned = row_index
with sess.lock:
sess.scanned = row_index
except Exception as e:
if not sess.cancelled:
# 取消路径:request_cancel 关连接 → iter_rows 抛错 → 这里吞掉(不算错误)
scan_error = f"{type(e).__name__}: {e}"
logger.exception(f"[scan] session={sess.session_id[:8]} 扫描异常")
finally:
with sess.lock:
sess.scan_conn = None # 摘引用(cancel 者不再拿得到;close 幂等)
sess.scan_done = True
if scan_error:
sess.scan_error = scan_error
logger.info(
f"[scan] done session={sess.session_id[:8]} scanned={row_index} "
f"bad={sess.bad_rows_total} truncated={sess.truncated} "
f"cancelled={sess.cancelled} error={scan_error} "
f"({time.monotonic() - started:.1f}s)"
)
# ── 端点 ──────────────────────────────────────────────────
@router.post("/start", response_model=StartQueryResponse, summary="创建查询 session(跑 COUNT + 编译规则)")
async def start_query(req: StartQueryRequest, db: Session = Depends(get_session)):
# 2026-09-24:端点为 def —— COUNT / 扫描控制都是同步 DB 调用,def 走 FastAPI
# 线程池,不冻结事件循环(真正的扫描在 _scan_worker 线程里,与请求线程解耦)。
@router.post("/start", response_model=StartQueryResponse, summary="创建查询 session(跑 COUNT + 编译规则 + 起扫描线程)")
def start_query(req: StartQueryRequest, db: Session = Depends(get_session)):
task, fields_with_rules, field_list_out, compiled_fields, col_names = _load_task_compiled(db, req.task_id)
if task is None:
return StartQueryResponse(ok=False, message=f"任务不存在:{req.task_id}", task_id=req.task_id)
......@@ -299,30 +393,35 @@ async def start_query(req: StartQueryRequest, db: Session = Depends(get_session)
task_id=task.id, field_list=field_list_out,
)
# 并发软上限:防连点 / 泄漏(TTL 会兜底回收,这里挡住突刺)
if active_scan_count() >= MAX_ACTIVE_SCANS:
return StartQueryResponse(
ok=False, message=f"并发扫描任务过多(≥{MAX_ACTIVE_SCANS}),请稍后再试",
task_id=req.task_id,
)
cfg = _build_db_config(task)
started = time.monotonic()
logger.info("─" * 60)
logger.info(
f"POST /api/queries/start task_id={task.id} table={task.source_table!r} "
f"fields={len(col_names)} rules_total={sum(len(cf.rules) for cf in compiled_fields)} "
f"page_size={req.page_size}"
f"fields={len(col_names)} rules_total={sum(len(cf.rules) for cf in compiled_fields)}"
)
# 1) COUNT(短连接,干净)
# 1) COUNT + 2) 字段注释 —— 2026-09-24 合并到同一条连接(原来开两条,
# 每条 ~50ms 建连 + info_schema 往返)。COUNT 失败整体失败(保持旧行为),
# 注释拉取失败只降级 warn(保持旧行为)。
total_rows = 0
try:
with DBConnection(cfg) as dbc:
count_sql = f"SELECT COUNT(*) AS n FROM {quote_ident(task.source_table, cfg.db_type)}"
count_row = dbc.fetchone(count_sql)
total_rows = int((count_row or {}).get("n") or 0)
except Exception as e:
msg = f"COUNT 失败:{type(e).__name__}: {e}"
logger.exception(msg)
return StartQueryResponse(ok=False, message=msg, task_id=task.id, field_list=field_list_out)
# 2) 补字段注释(info_schema;不影响主流程)
if (conn := _parse_conn_json(task.conn_json)) and (schema := (conn.get("schema") or conn.get("db") or "").strip()):
conn = _parse_conn_json(task.conn_json)
schema = ((conn or {}).get("schema") or (conn or {}).get("db") or "").strip()
if schema:
try:
with DBConnection(cfg) as dbc:
col_meta = dbc.list_columns(schema, task.source_table) or []
comments_by_key: dict[str, str] = {
(m.get("column_name") or "").lower(): (m.get("column_comment") or "")
......@@ -335,26 +434,31 @@ async def start_query(req: StartQueryRequest, db: Session = Depends(get_session)
f_out["field_comment"] = c
except Exception as e:
logger.warning(f"[queries] 获取字段注释失败,跳过:{type(e).__name__}: {e}")
except Exception as e:
msg = f"COUNT 失败:{type(e).__name__}: {e}"
logger.exception(msg)
return StartQueryResponse(ok=False, message=msg, task_id=task.id, field_list=field_list_out)
# 3) 拼 base SQL(ORDER BY 必加,否则 OFFSET 语义不一致)
# 3) 拼 base SQL(保留 ORDER BY 首列:流式扫一遍 + 全表只排序一次,
# 结果顺序与旧分页版一致;不再拼 LIMIT/OFFSET)
base_sql = _build_select_sql(cfg, col_names, task.source_table)
total_pages = max(1, (total_rows + req.page_size - 1) // req.page_size)
# 4) 建 session
# 4) 建 session + 起扫描线程(先 put 后 start,保证 /pull 一定能找到)
sess = QuerySession(
session_id=uuid.uuid4().hex,
task_id=task.id,
db_config=cfg,
db_type=cfg.db_type,
page_size=req.page_size,
total_rows=total_rows,
total_pages=total_pages,
select_sql_base=base_sql,
compiled=compiled_fields,
col_names=col_names,
field_list_out=field_list_out,
)
sm_put(sess)
t = threading.Thread(target=_scan_worker, args=(sess,), name=f"scan-{sess.session_id[:8]}", daemon=True)
sess.scan_thread = t
t.start()
# 2026-08-21:回写「最近检测时间」—— 任务被启动过查询就算检测行为
# 用户切到任务配置页能看到时间;中途中断也照样记录(代表「最近关注过」)
......@@ -370,8 +474,7 @@ async def start_query(req: StartQueryRequest, db: Session = Depends(get_session)
elapsed_ms = (time.monotonic() - started) * 1000
logger.info(
f"✅ session={sess.session_id} total_rows={total_rows} total_pages={total_pages} "
f"(耗时 {elapsed_ms:.0f}ms)"
f"✅ session={sess.session_id} total_rows={total_rows}(扫描线程已启动,耗时 {elapsed_ms:.0f}ms)"
)
logger.info("─" * 60)
return StartQueryResponse(
......@@ -379,81 +482,65 @@ async def start_query(req: StartQueryRequest, db: Session = Depends(get_session)
session_id=sess.session_id,
task_id=task.id,
total_rows=total_rows,
total_pages=total_pages,
page_size=req.page_size,
total_pages=0, # 兼容保留:流式版无页概念(前端不再读)
page_size=0,
field_list=field_list_out,
)
@router.post("/page", response_model=PageQueryResponse, summary="拉一页")
async def fetch_page(req: PageQueryRequest, db: Session = Depends(get_session)):
@router.post("/pull", response_model=PullQueryResponse, summary="增量拉取扫描结果(前端轮询)")
def pull_query(req: PullQueryRequest):
sess = sm_get(req.session_id)
if sess is None:
raise HTTPException(status_code=404, detail=f"session 不存在或已过期:{req.session_id}")
if req.page_no > sess.total_pages:
raise HTTPException(status_code=400, detail=f"page_no 越界:{req.page_no} > total_pages={sess.total_pages}")
cfg = sess.db_config
offset = (req.page_no - 1) * sess.page_size
sql = paginate_sql(sess.select_sql_base, cfg.db_type, offset, sess.page_size)
logger.info(
f"[queries] session={sess.session_id[:8]} page={req.page_no}/{sess.total_pages} "
f"offset={offset} size={sess.page_size}"
# 锁内只做切片 + 计数快照(≤PULL_MAX_ROWS 行,ms 级),
# 与扫描线程每 SCAN_PUBLISH_EVERY 行一次的持锁写几乎无竞争
with sess.lock:
buffered = len(sess.bad_rows)
if req.after > buffered:
raise HTTPException(
status_code=400,
detail=f"after={req.after} 越界(缓冲区现有 {buffered} 条)",
)
# 每页新连接(避开 cursor 生命周期复杂度,~50ms 重连可接受)
bad_rows: list[dict[str, Any]] = []
scanned_delta = 0
try:
with DBConnection(cfg) as dbc:
rows = dbc.fetchall(sql)
scanned_delta = len(rows)
for idx, raw in enumerate(rows):
err_cells, issues = _evaluate_row(sess.compiled, raw)
if not issues:
continue
row_out = dict(raw)
row_out["errorCells"] = err_cells
row_out["issues"] = issues
# 全局行号:(page_no-1)*page_size + idx → 前端做 row-key(append 模式必用)
row_out["__row_index"] = offset + idx
bad_rows.append(row_out)
except Exception as e:
# 单页拉失败不摘 session —— 让前端重试或走 /cancel
msg = f"拉取 page={req.page_no} 失败:{type(e).__name__}: {e}"
logger.exception(msg)
raise HTTPException(status_code=500, detail=msg)
sess.mark_page_done(bad_delta=len(bad_rows))
return PageQueryResponse(
chunk = sess.bad_rows[req.after : req.after + PULL_MAX_ROWS]
resp = PullQueryResponse(
ok=True,
session_id=sess.session_id,
page_no=req.page_no,
page_size=sess.page_size,
row_start=offset + 1,
bad_rows=bad_rows,
scanned_delta=scanned_delta,
bad_delta=len(bad_rows),
cancelled=False, # session 已被 cancel 时上面已 404
done=req.page_no == sess.total_pages,
bad_rows=chunk,
next_after=req.after + len(chunk),
scanned=sess.scanned,
total_rows=sess.total_rows,
bad_total=sess.bad_rows_total,
done=sess.scan_done,
error=sess.scan_error,
cancelled=sess.cancelled,
truncated=sess.truncated,
)
if chunk:
logger.info(
f"[queries] pull session={sess.session_id[:8]} after={req.after} "
f"+{len(chunk)} rows (scanned={resp.scanned} bad={resp.bad_total}"
f"{' truncated' if resp.truncated else ''})"
)
return resp
@router.post("/cancel", response_model=CancelQueryResponse, summary="取消查询 session")
async def cancel_query(req: CancelQueryRequest):
def cancel_query(req: CancelQueryRequest):
sess = sm_pop(req.session_id)
if sess is None:
# 已过期 / 已取消 / 不存在 —— 算 idempotent success
return CancelQueryResponse(ok=True, cancelled=False, pages_scanned=0, bad_rows_so_far=0)
sess.cancelled = True
return CancelQueryResponse(ok=True, cancelled=False, scanned=0, bad_rows_so_far=0)
with sess.lock:
scanned, bad_total = sess.scanned, sess.bad_rows_total
sess.request_cancel() # 置标志 + 关扫描连接(worker 批次间退出)
logger.info(
f"[queries] cancel session={sess.session_id[:8]} "
f"pages_scanned={sess.pages_scanned} bad_rows={sess.bad_rows_total}"
f"[queries] cancel session={sess.session_id[:8]} scanned={scanned} bad_rows={bad_total}"
)
return CancelQueryResponse(
ok=True,
cancelled=True,
pages_scanned=sess.pages_scanned,
bad_rows_so_far=sess.bad_rows_total,
scanned=scanned,
bad_rows_so_far=bad_total,
)
......@@ -51,7 +51,10 @@ async function request(path, { method = 'GET', body, params } = {}, { signal } =
const body = await res.json()
msg = errorDetailText(body.detail || body.message) || msg
} catch (_) { /* 非 JSON */ }
throw new Error(msg)
// 2026-09-24:挂上 HTTP status —— 轮询器等调用方要按 404(会话过期)分流
const err = new Error(msg)
err.status = res.status
throw err
}
if (res.status === 204) return null
return res.json()
......
/**
* 数据查询 / 校验 API(分页版,2026-08-21 重构)
* 数据查询 / 校验 API(流式扫描版,2026-09-24 重构)
*
* 老端点 POST /api/queries/run(一次性同步返回)已删。
* 旧分页三段式(start + 逐页 fetchPage + cancel)已删 —— 每页独立
* LIMIT/OFFSET 深分页在大表上退化严重(页越深扫得越多)。
* 新三段式:
* - startQuery(taskId, { pageSize, signal }) 创建 session + 返 total_rows/total_pages/field_list
* - fetchPage(sessionId, pageNo, { signal }) 拉一页(可并发,多页同时拉)
* - cancelQuery(sessionId, { signal }) 主动取消,best-effort
* - startQuery(taskId, { signal }) 创建 session:COUNT + 编译规则 + 起后台扫描线程
* - pullQuery(sessionId, after, { signal }) 增量拉取命中行 + 进度(前端轮询,满额立即续拉)
* - cancelQuery(sessionId, { signal }) 主动取消:摘 session + 关扫描连接
*
* 三段都支持 AbortSignal:传 signal 后前端 abort 会让 fetch 直接 reject
* (网络中断 + 节省带宽),同时 server 端 /cancel 摘掉 session 让后续 /page 立刻 404。
* (网络中断 + 节省带宽),同时 server 端 /cancel 摘掉 session 让后续 /pull 立刻 404。
*/
import { http } from './client'
/**
* 创建查询 session(跑 COUNT + 编译规则)
* 创建查询 session(COUNT + 编译规则 + 起扫描线程)
*
* @param {number} taskId
* @param {{ pageSize?: number, signal?: AbortSignal }} [opts]
* @param {{ signal?: AbortSignal }} [opts]
* @returns {Promise<{
* ok: boolean,
* message?: string,
* session_id: string,
* task_id: number,
* total_rows: number,
* total_pages: number,
* page_size: number,
* total_rows: number, // 进度分母
* total_pages: number, // 兼容保留,流式版恒 0
* page_size: number, // 兼容保留,流式版恒 0
* field_list: Array<{
* id: number, field_key: string, show_default: boolean, ord: number,
* rules: number, rule_list: Array<{ id, desc, rule_type, regex?, code? }>,
......@@ -33,37 +34,37 @@ import { http } from './client'
* }>,
* }>}
*/
export function startQuery(taskId, { pageSize = 500, signal } = {}) {
return http.post(
'/queries/start',
{ task_id: taskId, page_size: pageSize },
{ signal },
)
export function startQuery(taskId, { signal } = {}) {
return http.post('/queries/start', { task_id: taskId }, { signal })
}
/**
* 拉一页(单页独立连接,跑规则,只返不合规行)
* 增量拉取扫描结果(前端轮询;返回的 next_after 作为下次的 after 游标)
*
* 满额返回(bad_rows.length === 5000)时应立即再拉(清缓冲);不足额等下一轮询周期。
* done=true 表示扫描线程已终止,这是最后一次有效数据。
*
* @param {string} sessionId /start 返回的 session token
* @param {number} pageNo 1-based
* @param {number} after 已消费的 bad_rows 数(游标,首次 0)
* @param {{ signal?: AbortSignal }} [opts]
* @returns {Promise<{
* ok: boolean,
* session_id: string,
* page_no: number,
* page_size: number,
* row_start: number,
* bad_rows: Array<{ __row_index: number, errorCells: string[], issues: Array, ...业务列 }>,
* scanned_delta: number,
* bad_delta: number,
* cancelled: boolean, // session 已被 /cancel 时 true
* done: boolean, // 本页是最后一页
* next_after: number, // after + len(bad_rows)
* scanned: number, // 已扫描行数(进度分子)
* total_rows: number, // COUNT 结果(进度分母)
* bad_total: number, // 命中总数(含截断丢弃部分)
* done: boolean, // 扫描线程已终止
* error?: string, // 非取消原因的扫描异常
* cancelled: boolean, // 已被 /cancel(前端据此静默退出)
* truncated: boolean, // 命中超缓冲区上限(提示「仅展示前 N 条」)
* }>}
*/
export function fetchPage(sessionId, pageNo, { signal } = {}) {
export function pullQuery(sessionId, after, { signal } = {}) {
return http.post(
'/queries/page',
{ session_id: sessionId, page_no: pageNo },
'/queries/pull',
{ session_id: sessionId, after },
{ signal },
)
}
......@@ -76,7 +77,7 @@ export function fetchPage(sessionId, pageNo, { signal } = {}) {
* @returns {Promise<{
* ok: boolean,
* cancelled: boolean,
* pages_scanned: number,
* scanned: number,
* bad_rows_so_far: number,
* }>}
*/
......
......@@ -37,8 +37,7 @@
<el-icon><InfoFilled /></el-icon>
<span>当前任务:<b>{{ selectedTask }}</b></span>
<span class="meta" v-if="runningQuery">
扫描中
<b>{{ scannedPages }}</b> / <b>{{ totalPages }}</b> 页
已扫描 <b>{{ scannedCount.toLocaleString() }}</b> / <b>{{ totalRows.toLocaleString() }}</b> 条({{ progressPercent }}%)
· 命中不合规 <b style="color:#f56c6c;">{{ rows.length }}</b> 条
</span>
<span class="meta" v-else-if="lastResult">
......@@ -82,7 +81,7 @@ import { computed, ref, watch, onBeforeUnmount } from 'vue'
import { useRoute, useRouter } from 'vue-router'
import { ElMessage, ElMessageBox } from 'element-plus'
import { listTasks } from '@/api/tasks'
import { startQuery, fetchPage, cancelQuery } from '@/api/queries'
import { startQuery, pullQuery, cancelQuery } from '@/api/queries'
import { exportToExcel } from '@/utils/excel'
import ResultTable from '@/components/ResultTable.vue'
import ExplainIssueDialog from '@/components/ExplainIssueDialog.vue'
......@@ -136,11 +135,11 @@ async function loadTasksForCurrentType() {
const cacheByType = new Map()
function takeSnapshot(typeKey) {
// lastResult 兜底:appendPage 已经每页都更新,但 takeSnapshot 可能在第一次
// appendPage 跑前调用(极小窗口:用户刚点查询就立刻切走)→ lastResult 还是 null。
// lastResult 兜底:pollLoop 每次 pull 都更新,但 takeSnapshot 可能在第一次
// pull 跑前调用(极小窗口:用户刚点查询就立刻切走)→ lastResult 还是 null。
// 给一个 partial 兜底,让切回时 top 信息条不显示「尚未查询」。
const partial = runningQuery.value
? { scanned: scannedRows.value, bad_rows: rows.value.length, rows: rows.value }
? { scanned: scannedCount.value, bad_rows: rows.value.length, rows: rows.value }
: null
return {
selectedTask: selectedTask.value,
......@@ -148,16 +147,11 @@ function takeSnapshot(typeKey) {
rows: rows.value,
resultFields: resultFields.value,
lastResult: lastResult.value || partial,
// 续跑所需:in-flight 状态下用,新 session 重新拉剩下的页
// 续跑所需(2026-09-24 流式版):新 session 会全表重扫,已收到的行靠
// __row_index > lastRowIndex 去重(旧分页版按 receivedPages 页号去重已删)
inFlight: runningQuery.value ? {
taskId: taskOptions.value.find((t) => t.name === selectedTask.value)?.id,
pageSize: PAGE_SIZE,
// Map → entries 数组(Map 不能 JSON 序列化;后续 new Map(entries) 恢复)
receivedPages: [...receivedPages.value.entries()],
scannedPages: scannedPages.value,
scannedRows: scannedRows.value,
totalPages: totalPages.value,
totalRows: totalRows.value,
lastRowIndex: rows.value.reduce((m, r) => Math.max(m, r.__row_index ?? -1), -1),
} : null,
}
}
......@@ -173,6 +167,8 @@ function applySnapshot(snap) {
rows.value = []
resultFields.value = []
lastResult.value = null
scannedCount.value = 0
totalRows.value = 0
return
}
// 切回来:恢复上次离开时的状态
......@@ -181,54 +177,36 @@ function applySnapshot(snap) {
rows.value = snap.rows
resultFields.value = snap.resultFields
lastResult.value = snap.lastResult
scannedCount.value = snap.lastResult?.scanned || 0
// 恢复 in-flight 进度(fire-and-forget 启动续跑)
if (snap.inFlight) {
receivedPages.value = new Map(snap.inFlight.receivedPages)
scannedPages.value = snap.inFlight.scannedPages
scannedRows.value = snap.inFlight.scannedRows
totalPages.value = snap.inFlight.totalPages
totalRows.value = snap.inFlight.totalRows
// 不 await —— applySnapshot 是同步语义,让 watch(type) 不等续跑完成
continueQuery(snap.inFlight)
}
}
// 续跑(fire-and-forget):切回时如果上次 in-flight,自动起新 session 拉剩下的页
// 续跑(fire-and-forget):切回时如果上次 in-flight,自动起新 session 重扫全表
// 为什么不复用老 sessionId:watch(type) 头部的 abortInFlight 已经 cancelQuery 老 session,
// 后端会拒后续 /page 请求。只能重新 startQuery 拿新 sessionId。
// 为什么不直接调 onQuery:onQuery 会 resetQueryStatus 清空 receivedPages/rows,
// 重复拉已经收到的页,浪费带宽(部分 task 总扫描数上万)。
// 这里保留 receivedPages(已收到的页)+ 启动新 worker 只拉剩下的页。
// 后端会拒后续 /pull 请求。只能重新 startQuery 拿新 sessionId。
// 为什么不直接调 onQuery:onQuery 会 resetQueryStatus 清空 rows,
// 已收到的命中行会闪没。这里保留 rows + pollLoop 按 __row_index 去重追加。
// 注:watch(type) 已在 runningQuery 时禁止切分组(见下方),此路径实际不可达,
// 作为防御保留;重扫的行序与旧 session 一致(ORDER BY 首列相同)。
async function continueQuery(snap) {
// snap: { taskId, pageSize, receivedPages, scannedPages, scannedRows, totalPages, totalRows }
// snap: { taskId, lastRowIndex }
const ac = new AbortController()
queryAbort.value = ac
runningQuery.value = true
cancelRequested = false
try {
const start = await startQuery(snap.taskId, { pageSize: snap.pageSize, signal: ac.signal })
const start = await startQuery(snap.taskId, { signal: ac.signal })
if (!start.ok) {
queryError.value = start.message || '续跑失败'
return
}
sessionId.value = start.session_id
totalPages.value = start.total_pages
totalRows.value = start.total_rows
// 已收到的页就不重拉了
const received = new Set(receivedPages.value.keys())
const queue = Array.from({ length: start.total_pages }, (_, i) => i + 1)
.filter((p) => !received.has(p))
if (queue.length === 0) {
finalize()
return
}
const workers = Array.from(
{ length: Math.min(PAGE_CONCURRENCY, queue.length) },
() => worker(queue, ac.signal),
)
await Promise.all(workers)
if (ac.signal.aborted) return
finalize()
await pollLoop(ac, { lastRowIndex: snap.lastRowIndex })
} catch (e) {
if (e.name === 'AbortError' || ac.signal.aborted) return
queryError.value = `续跑失败:${e.message || e}`
......@@ -236,102 +214,150 @@ async function continueQuery(snap) {
runningQuery.value = false
queryAbort.value = null
}
function finalize() {
const totalBad = rows.value.length
lastResult.value = {
scanned: totalRows.value,
bad_rows: totalBad,
rows: rows.value,
}
}
}
// ── 分页查询 orchestrator(2026-08-21 重构) ──
const PAGE_CONCURRENCY = 3 // 并发 /page 数;> 3 可能压垮 DB
const PAGE_SIZE = 500 // 每页行数,与后端默认对齐
// ── 流式扫描轮询器(2026-09-24,替代 3-worker 分页队列) ──
// 后端 /start 起扫描线程单连接流式扫全表,前端 pollLoop 每 500ms 拉一次增量:
// 满额返回(=PULL_CHUNK_MAX,缓冲有积压)→ 立即续拉清缓冲
// 未满额 → 等 PULL_INTERVAL_MS 再拉
// done=true → 定稿(toast + lastResult)后退出
const PULL_INTERVAL_MS = 500 // 轮询周期
const PULL_CHUNK_MAX = 5000 // 与后端 PULL_MAX_ROWS 对齐(满额判断)
const runningQuery = ref(false)
const cancelling = ref(false)
const queryAbort = ref(null) // AbortController 实例
const sessionId = ref(null)
const totalPages = ref(0)
const totalRows = ref(0)
const scannedPages = ref(0)
const scannedRows = ref(0)
const receivedPages = ref(new Map()) // page_no → bad_rows[](按页号缓存)
const scannedCount = ref(0)
const lastResult = ref(null) // 最近一次查询结果(给顶部信息条用)
const queryError = ref('')
const rows = ref([]) // 喂给 ResultTable 的行
const resultFields = ref([])
// 本端刚发过 cancel(onCancel / abortInFlight / unmount)→ 之后 /pull 的 404 静默退出
// (非本端取消的 404 = 会话被别人取消 / TTL 过期,要给用户提示)
let cancelRequested = false
// 进度百分比:封顶 99%(COUNT 后表若有插入,scanned 可能超过 total;
// 定稿后信息条切到 lastResult,不再显示百分比)
const progressPercent = computed(() =>
totalRows.value > 0
? Math.min(99, Math.floor((scannedCount.value / totalRows.value) * 100))
: 0)
// 行 key:用后端注入的全局行号(append 模式下保证 Vue DOM 复用正确)
const rowKeyFn = (row) => row.__row_index
// 把一页结果合并进 rows:按页号升序拼,保证索引稳定
// 可中断的 sleep:abort 时立即 resolve,让循环顶部的 aborted 检查尽快退出
function sleep(ms, signal) {
return new Promise((resolve) => {
const t = setTimeout(resolve, ms)
signal?.addEventListener('abort', () => { clearTimeout(t); resolve() }, { once: true })
})
}
// 把一次 pull 的增量合并进 rows(顺序 append,__row_index 全局行号天然有序)
// 顺便实时更新 lastResult(in-flight 期间顶部信息条也能看到「已扫 X / 已命中 Y」进度)
// 收尾时 onQuery / continueQuery 会用 totalRows 覆盖,scannedRows 字段不参与最终展示
function appendPage(page) {
receivedPages.value.set(page.page_no, page.bad_rows)
const ordered = [...receivedPages.value.keys()].sort((a, b) => a - b)
rows.value = ordered.flatMap((n) => receivedPages.value.get(n) || [])
scannedPages.value = receivedPages.value.size
scannedRows.value += page.scanned_delta
function appendRows(chunk) {
rows.value.push(...chunk)
lastResult.value = {
scanned: scannedRows.value,
scanned: scannedCount.value,
bad_rows: rows.value.length,
rows: rows.value,
}
}
// 轮询器状态机(onQuery / continueQuery 共用):
// 正常 → 追加行 / 更新进度 / done 时定稿退出
// cancelled → 后端响应里 cancelled=true(极小窗口)→ 静默退出
// 404 → session 已摘:本端取消过(cancelRequested)静默;否则横幅提示
// 网络异常 → 2s 退避重试,连续 3 次 → 横幅退出(成功响应后清零)
// AbortError → 用户中断 / 新查询 / unmount → 静默退出
async function pollLoop(ac, { lastRowIndex = -1 } = {}) {
let after = 0
let failures = 0
while (!ac.signal.aborted) {
let res
try {
res = await pullQuery(sessionId.value, after, { signal: ac.signal })
failures = 0
} catch (e) {
if (e.name === 'AbortError' || ac.signal.aborted) return
if (e.status === 404) {
if (!cancelRequested) queryError.value = '会话已过期或已被取消,请重新扫描'
return
}
failures += 1
if (failures >= 3) {
queryError.value = `拉取结果失败:${e.message || e}`
return
}
await sleep(2000, ac.signal)
continue
}
if (res.cancelled) return // 被取消(极小窗口,静默)
scannedCount.value = res.scanned
// 增量行:continueQuery 场景按 __row_index 去重(新 session 全表重扫)
if (res.bad_rows.length) {
const fresh = lastRowIndex < 0
? res.bad_rows
: res.bad_rows.filter((r) => (r.__row_index ?? 0) > lastRowIndex)
if (fresh.length) appendRows(fresh)
}
after = res.next_after
if (res.done) {
scannedCount.value = res.scanned
if (res.error) {
queryError.value = `扫描失败:${res.error}`
try { await cancelQuery(sessionId.value) } catch (_) { /* best-effort 释放 */ }
} else {
// 定稿:bad_rows 用 bad_total(含截断丢弃部分),rows 是实际展示的
lastResult.value = { scanned: res.scanned, bad_rows: res.bad_total, rows: rows.value }
if (res.truncated) {
ElMessage.warning(
`命中过多(共 ${res.bad_total.toLocaleString()} 条),仅展示前 ${rows.value.length.toLocaleString()} 条`,
)
} else if (res.bad_total === 0) {
ElMessage.success(`全表扫描 ${res.scanned.toLocaleString()} 条,全部合规`)
} else {
ElMessage.success(`全表扫描 ${res.scanned.toLocaleString()} 条,命中 ${res.bad_total.toLocaleString()} 条不合规`)
}
}
return
}
// 未完成:满额立即续拉(后端缓冲有积压),否则等一个轮询周期
if (res.bad_rows.length < PULL_CHUNK_MAX) {
await sleep(PULL_INTERVAL_MS, ac.signal)
}
}
}
function resetQueryStatus() {
sessionId.value = null
totalPages.value = 0
totalRows.value = 0
scannedPages.value = 0
scannedRows.value = 0
receivedPages.value = new Map()
scannedCount.value = 0
rows.value = []
lastResult.value = null
queryError.value = ''
cancelRequested = false
}
// 「中断」专用:只清进度(session + 页号计数),保留 rows + resultFields + lastResult
// 「中断」专用:只清进度(session + 计数),保留 rows + resultFields + lastResult
// —— 用户已经看到的不合规行不能丢;下一次点「查询」会走完整 onQuery,自然覆盖
function clearProgressOnly() {
sessionId.value = null
totalPages.value = 0
totalRows.value = 0
scannedPages.value = 0
scannedRows.value = 0
receivedPages.value = new Map()
}
// worker:模块级函数(onQuery / continueQuery 都用得到)
// 从 queue 拿一个页号去拉,拉到就追加
async function worker(q, signal) {
while (q.length && !signal.aborted) {
const pageNo = q.shift()
try {
const page = await fetchPage(sessionId.value, pageNo, { signal })
appendPage(page)
} catch (e) {
// abort 三种识别:
// 1) e.name === 'AbortError' —— fetch 自己识别的 abort(最常见)
// 2) signal.aborted === true —— abortInFlight 已发,但 fetch 已经发出请求
// 后端拿到取消通知返回 404/410 等非 AbortError 错误,catch 里再兜一次
// 3) fetch throw TypeError + signal.aborted —— 部分浏览器实现
if (e.name === 'AbortError' || signal.aborted) return
// 单页拉失败:记错误但不让整个查询崩;让其它页继续拉
queryError.value = `page=${pageNo} 拉取失败:${e.message || e}`
console.warn(queryError.value)
}
}
scannedCount.value = 0
}
async function abortInFlight() {
// 取消上一轮还在飞的查询(用户连续点 / 切换任务时)
if (sessionId.value) {
cancelRequested = true // 本端取消:后续 /pull 404 静默
try { await cancelQuery(sessionId.value) } catch (_) { /* best-effort */ }
}
queryAbort.value?.abort()
......@@ -353,15 +379,20 @@ async function onQuery() {
runningQuery.value = true
try {
// 1) /start:COUNT + 编译规则 + 建 session
const start = await startQuery(t.id, { pageSize: PAGE_SIZE, signal: ac.signal })
// 1) /start:COUNT + 编译规则 + 建 session + 起后端扫描线程
const start = await startQuery(t.id, { signal: ac.signal })
if (!start.ok) {
queryError.value = start.message || '查询失败'
ElMessage.error(queryError.value)
return
}
if (!start.session_id && start.message) {
// ok=true 但没建 session:目前只有「任务未配置任何校验规则」这一种
// (旧版会误走 total_rows=0 分支提示「该表无数据」,2026-09-24 修正)
ElMessage.warning(start.message)
return
}
sessionId.value = start.session_id
totalPages.value = start.total_pages
totalRows.value = start.total_rows
resultFields.value = (start.field_list || []).map((f) => ({
// 关键:key 跟 db_adapter.normalize_column_name 对齐成小写
......@@ -376,35 +407,15 @@ async function onQuery() {
ruleList: f.rule_list || [],
showDefault: f.show_default !== false,
}))
if (start.total_pages === 0) {
if (start.total_rows === 0) {
// COUNT=0:表里没数据
lastResult.value = { scanned: 0, bad_rows: 0, rows: [] }
ElMessage.warning('该表无数据(0 行)')
return
}
// 2) 并发 worker 拉所有页(PAGE_CONCURRENCY 并发)
const queue = Array.from({ length: start.total_pages }, (_, i) => i + 1)
const workers = Array.from(
{ length: Math.min(PAGE_CONCURRENCY, queue.length) },
() => worker(queue, ac.signal),
)
await Promise.all(workers)
if (ac.signal.aborted) return // 中断流程不弹成功 toast
// 3) 收尾:聚合总数 + toast
const totalBad = rows.value.length
lastResult.value = {
scanned: totalRows.value,
bad_rows: totalBad,
rows: rows.value,
}
if (totalBad === 0) {
ElMessage.success(`全表扫描 ${totalRows.value} 条,全部合规`)
} else {
ElMessage.success(`全表扫描 ${totalRows.value} 条,命中 ${totalBad} 条不合规`)
}
// 2) 轮询增量直到 done(toast / lastResult 定稿都在 pollLoop 内)
await pollLoop(ac)
} catch (e) {
if (e.name === 'AbortError' || ac.signal.aborted) return
queryError.value = `请求失败:${e.message || e}`
......@@ -418,15 +429,16 @@ async function onQuery() {
async function onCancel() {
if (!sessionId.value) return
cancelling.value = true
cancelRequested = true // 本端取消:/pull 404 静默
try {
try { await cancelQuery(sessionId.value) } catch (_) { /* best-effort */ }
queryAbort.value?.abort() // 让在飞的 fetch 立即 reject
queryAbort.value?.abort() // 让在飞的 pull 立即 reject
// 注意:不要调 resetQueryStatus —— 那会把 rows 也清空。
// 中断后用户已经看到的不合规行要留住(用户重跑时会自然覆盖)
const hit = rows.value.length
const pages = scannedPages.value
const scanned = scannedCount.value
clearProgressOnly()
ElMessage.info(`已中断(已扫 ${pages} 页,命中 ${hit} 条)`)
ElMessage.info(`已中断(已扫 ${scanned.toLocaleString()} 条,命中 ${hit} 条)`)
} finally {
runningQuery.value = false
cancelling.value = false
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment