Commit 8013322d authored by Data Governance Dev's avatar Data Governance Dev

feat(web3): 查询只返回不合规行 + 改为全表流式扫描

判定口径修正:查询结果应是「不符合」校验规则的数据,之前合规行也一起返回了。

- 字段内多条规则 = AND:任一条不满足该字段即不合规,且所有未通过的规则
  都进 issues(去掉原来命中第一条就 break 的逻辑),备注列能看全原因
- 行内多个字段 = AND:任一字段不合规整行就被查出来
- 全部规则通过的行直接丢弃,不再进 rows
- total 从 len(raw_rows) 改为 bad_rows,与 rows 对齐

全表扫描(原来固定 LIMIT 1000,只校验前 1000 行):

- db_adapter 新增 DBConnection.iter_rows(),fetchmany 流式逐行读。
  MySQL 必须用 SSDictCursor —— 普通 Cursor 是客户端缓冲的,execute 时
  整个结果集就进内存了,fetchmany 省不了内存;达梦/Oracle 驱动默认
  就是服务端游标,普通 cursor 即可
- SQL 去掉 LIMIT,读取与校验合成一个流式 pass,只有命中的不合规行进内存
- 请求参数 limit(扫描上限)→ max_rows(返回上限,默认 5000);
  超上限后 continue 而非 break,继续扫完全表,保证 scanned / bad_rows
  始终是全表真值,避免计数被截断成假数字
- 响应新增 truncated 标志 + 截断提示文案

顺带修复:

- SELECT 的列从「有规则的字段」改为「有规则 ∪ show_default」,
  否则勾了默认展示但没配规则的列在结果表里全空
- field_list_out 去掉对 Rule 的重复查询,抽 rules_by_field 复用

前端:runQuery 参数改 max_rows,文案改「全表扫描」,truncated 时顶部
常驻横幅 + 信息栏提示仅展示前 N 条。

测试:新增 web3/tests/test_run_query_only_bad_rows.py,临时 sqlite 当配置库
+ monkeypatch DBConnection 喂假数据,不依赖真实 MySQL/达梦,9 条用例全过。
parent f05e53ba
...@@ -500,6 +500,51 @@ class DBConnection: ...@@ -500,6 +500,51 @@ class DBConnection:
finally: finally:
cur.close() cur.close()
def iter_rows(
self,
sql: str,
params: tuple | None = None,
batch_size: int = 1000,
) -> Iterator[dict]:
"""流式逐行返回结果(列名已小写归一),用于全表扫描。
与 fetchall 的区别:**不把整张表读进内存**。调用方边拿边处理,
内存只由 batch_size 和调用方自己攒的东西决定。
方言差异:
- MySQL:普通 Cursor 是「客户端缓冲」的 —— cur.execute() 就把整个结果集
拉到本地了,再 fetchmany 只是从内存里切片,省不了内存。
必须换成 SSDictCursor(unbuffered / server-side)才是真流式。
代价:结果集没消费完之前这条连接不能干别的(这里是独占连接,无所谓)。
- 达梦 / Oracle:驱动默认就是服务端游标 + 数组取数,普通 cursor
配 fetchmany 即为流式。
"""
if self._conn is None:
raise RuntimeError("DBConnection 未打开,请用 with 语句")
if self._driver == "pymysql":
from pymysql.cursors import SSDictCursor
cur = self._conn.cursor(SSDictCursor)
else:
cur = self._conn.cursor()
started = time.monotonic()
n = 0
try:
cur.execute(sql, params or ())
col_names = [d[0] for d in (cur.description or [])]
while True:
batch = cur.fetchmany(batch_size)
if not batch:
break
for row in _rows_as_dicts(batch, col_names):
n += 1
yield row
finally:
elapsed_ms = (time.monotonic() - started) * 1000
logger.debug(f"[SQL] iter_rows {n} rows ({elapsed_ms:.1f}ms): {sql[:120]}...")
cur.close()
def fetchone(self, sql: str, params: tuple | None = None) -> Optional[dict]: def fetchone(self, sql: str, params: tuple | None = None) -> Optional[dict]:
cur = self._cursor() cur = self._cursor()
started = time.monotonic() started = time.monotonic()
......
This diff is collapsed.
...@@ -8,10 +8,14 @@ import { http } from './client' ...@@ -8,10 +8,14 @@ import { http } from './client'
/** /**
* 跑一次任务查询 + 校验 * 跑一次任务查询 + 校验
*
* 扫描范围是**全表**(后端不带 LIMIT,流式逐行读)。
* maxRows 只截断返回的行数,不影响 scanned / bad_rows —— 它们始终是全表真值。
*
* @param {number} taskId 任务 id * @param {number} taskId 任务 id
* @param {number} limit 最多扫多少行(后端默认 1000) * @param {number} maxRows 最多返回多少条不合规行(后端默认 5000)
* @returns {Promise<{ok, message?, scanned, bad_rows, total, rows: [], field_list: []}>} * @returns {Promise<{ok, message?, scanned, bad_rows, total, truncated, rows: [], field_list: []}>}
*/ */
export function runQuery(taskId, limit = 1000) { export function runQuery(taskId, maxRows = 5000) {
return http.post('/queries/run', { task_id: taskId, limit }) return http.post('/queries/run', { task_id: taskId, max_rows: maxRows })
} }
...@@ -26,8 +26,11 @@ ...@@ -26,8 +26,11 @@
<el-icon><InfoFilled /></el-icon> <el-icon><InfoFilled /></el-icon>
<span>当前任务:<b>{{ selectedTask }}</b></span> <span>当前任务:<b>{{ selectedTask }}</b></span>
<span class="meta" v-if="lastResult"> <span class="meta" v-if="lastResult">
共扫描 <b>{{ lastResult.scanned }}</b> 条 全表扫描 <b>{{ lastResult.scanned }}</b> 条
· 命中不合规 <b style="color:#f56c6c;">{{ lastResult.bad_rows }}</b> 条 · 命中不合规 <b style="color:#f56c6c;">{{ lastResult.bad_rows }}</b> 条
<template v-if="lastResult.truncated">
· 表内仅展示前 <b>{{ lastResult.rows.length }}</b> 条
</template>
</span> </span>
<span class="meta" v-else style="color:#909399;">尚未查询</span> <span class="meta" v-else style="color:#909399;">尚未查询</span>
</div> </div>
...@@ -152,7 +155,7 @@ async function onQuery() { ...@@ -152,7 +155,7 @@ async function onQuery() {
runningQuery.value = true runningQuery.value = true
queryError.value = '' queryError.value = ''
try { try {
const r = await runQuery(t.id, 1000) const r = await runQuery(t.id)
if (!r.ok) { if (!r.ok) {
// 后端明确返回失败:清空表 + 顶部横幅 + toast // 后端明确返回失败:清空表 + 顶部横幅 + toast
queryError.value = r.message || '查询失败' queryError.value = r.message || '查询失败'
...@@ -171,12 +174,16 @@ async function onQuery() { ...@@ -171,12 +174,16 @@ async function onQuery() {
showDefault: f.show_default !== false, // 任务配置「默认展示」未勾选 → 结果表不展示该列 showDefault: f.show_default !== false, // 任务配置「默认展示」未勾选 → 结果表不展示该列
})) }))
lastResult.value = r lastResult.value = r
if (r.scanned === 0) { if (r.truncated) {
// 全表扫描的计数是准的,但表里只装得下前 max_rows 条 → 顶部常驻提示
queryError.value = r.message
ElMessage.warning(r.message)
} else if (r.scanned === 0) {
ElMessage.warning('该表无数据(0 行)') ElMessage.warning('该表无数据(0 行)')
} else if (r.bad_rows === 0) { } else if (r.bad_rows === 0) {
ElMessage.success(`扫描 ${r.scanned} 条,全部合规`) ElMessage.success(`全表扫描 ${r.scanned} 条,全部合规`)
} else { } else {
ElMessage.success(`扫描 ${r.scanned} 条,命中 ${r.bad_rows} 条不合规`) ElMessage.success(`全表扫描 ${r.scanned} 条,命中 ${r.bad_rows} 条不合规`)
} }
} catch (e) { } catch (e) {
queryError.value = `请求失败:${e.message}` queryError.value = `请求失败:${e.message}`
......
"""
验证:POST /api/queries/run 只返回「不合规」的行。
判定口径(本次需求):
- 字段内多条规则是 AND —— 任一条不满足,该字段就不合规,且每条未通过的规则都进 issues
- 行内多个字段也是 AND —— 任一字段不合规,整行就要被查出来
- 全部规则都通过的行 **不返回**
做法:用临时 sqlite 当配置库(task/field/rule),monkeypatch 掉 DBConnection,
喂一批构造好的目标表数据,不依赖真实 MySQL / 达梦。
跑:python -m pytest web3/tests/test_run_query_only_bad_rows.py -v
"""
import json
import sys
from pathlib import Path
import pytest
from fastapi import FastAPI
from fastapi.testclient import TestClient
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
sys.path.insert(0, str(Path(__file__).resolve().parents[2]))
from web3.backend.db.database import Base, get_session # noqa: E402
from web3.backend.models.field import Field # noqa: E402
from web3.backend.models.rule import Rule # noqa: E402
from web3.backend.models.task import Task # noqa: E402
from web3.backend.models.task_group import TaskGroup # noqa: E402
from web3.backend.routers import queries as queries_mod # noqa: E402
# ── 目标表的假数据 ───────────────────────────────────────────
# 字段:id_card(有 2 条规则)/ phone(有 1 条规则)/ person_name(无规则但 show_default)
FAKE_TABLE_ROWS = [
# 0 全合规 → 不该被返回
{"id_card": "110101199003074512", "phone": "13800138000", "person_name": "张三"},
# 1 id_card 长度不对(规则 A 挂,规则 B 也挂)→ 返回,issues 应有 2 条
{"id_card": "1101", "phone": "13800138001", "person_name": "李四"},
# 2 id_card 合规但 phone 不合规 → 返回(另一个字段挂也要查出来)
{"id_card": "110101199003074512", "phone": "021-8888", "person_name": "王五"},
# 3 两个字段都挂 → 返回,errorCells 应含两个字段
{"id_card": "abc", "phone": "nope", "person_name": "赵六"},
# 4 全合规 → 不该被返回
{"id_card": "310101198801011234", "phone": "13900139000", "person_name": "钱七"},
]
class _FakeDBConnection:
"""替身:吞掉 DBConfig,iter_rows 流式吐 rows(默认 FAKE_TABLE_ROWS)。"""
last_sql = None
rows = FAKE_TABLE_ROWS
def __init__(self, cfg):
self.cfg = cfg
def __enter__(self):
return self
def __exit__(self, *exc):
return False
def iter_rows(self, sql, params=None, batch_size=1000):
_FakeDBConnection.last_sql = sql
for r in _FakeDBConnection.rows:
yield dict(r)
@pytest.fixture()
def client(tmp_path, monkeypatch):
_FakeDBConnection.rows = FAKE_TABLE_ROWS # 每个用例复位,避免相互污染
_FakeDBConnection.last_sql = None
engine = create_engine(
f"sqlite:///{tmp_path / 'test.sqlite'}",
connect_args={"check_same_thread": False},
)
Base.metadata.create_all(engine)
TestSession = sessionmaker(bind=engine, autocommit=False, autoflush=False)
with TestSession() as db:
db.add(TaskGroup(id=1, name="身份证"))
db.add(Task(
id=1, name="test-bad-rows", group_id=1,
source_table="t_user_info", db_type="MySQL",
conn_json=json.dumps({"host": "127.0.0.1", "port": 3306,
"user": "u", "password": "p", "db": "d"}),
))
db.flush()
# id_card:两条规则(AND)
db.add(Field(id=10, task_id=1, field_key="id_card", show_default=1, ord=0))
db.add(Rule(field_id=10, desc="必须 18 位", regex=r"^.{18}$", ord=0))
db.add(Rule(field_id=10, desc="必须全是数字或以 X 结尾", regex=r"^\d{17}[\dXx]$", ord=1))
# phone:一条规则
db.add(Field(id=11, task_id=1, field_key="phone", show_default=1, ord=1))
db.add(Rule(field_id=11, desc="手机号 11 位", regex=r"^1\d{10}$", ord=0))
# person_name:无规则,但默认展示 → 应被 SELECT 出来当上下文
db.add(Field(id=12, task_id=1, field_key="person_name", show_default=1, ord=2))
db.commit()
monkeypatch.setattr(queries_mod, "DBConnection", _FakeDBConnection)
app = FastAPI()
app.include_router(queries_mod.router, prefix="/api")
app.dependency_overrides[get_session] = lambda: TestSession()
with TestClient(app) as c:
yield c
def _run(client, **payload):
body_in = {"task_id": 1, **payload}
r = client.post("/api/queries/run", json=body_in)
assert r.status_code == 200, r.text
body = r.json()
assert body["ok"] is True, body.get("message")
return body
def test_only_bad_rows_returned(client):
"""5 行里 3 行不合规 → rows 只含那 3 行,合规行被丢弃。"""
body = _run(client)
assert body["scanned"] == 5, "scanned 应是实际扫到的行数"
assert body["bad_rows"] == 3
assert body["total"] == 3, "total 应跟 bad_rows 对齐(rows 只含不合规行)"
assert len(body["rows"]) == 3
names = [r["person_name"] for r in body["rows"]]
assert names == ["李四", "王五", "赵六"]
assert "张三" not in names and "钱七" not in names, "合规行不该返回"
# 每一行都必须有 issues,否则不该出现在结果里
for row in body["rows"]:
assert row["issues"], f"返回的行必须带不合规原因:{row}"
def test_multiple_rules_on_one_field_all_reported(client):
"""一个字段挂多条规则:任一条不满足就查出来,且所有未通过的规则都进 issues。"""
body = _run(client)
li_si = next(r for r in body["rows"] if r["person_name"] == "李四")
descs = [i["desc"] for i in li_si["issues"] if i["field"] == "id_card"]
assert len(descs) == 2, f"id_card 两条规则都不满足,应报 2 条,实际 {descs}"
assert li_si["errorCells"] == ["id_card"], "phone 合规,不该标红"
def test_any_field_failing_flags_the_row(client):
"""任务里多个字段有规则:只要一个字段的一条规则不满足,这行就要被查出来。"""
body = _run(client)
# 王五:id_card 合规、phone 不合规 → 依然要被查出来
wang = next(r for r in body["rows"] if r["person_name"] == "王五")
assert wang["errorCells"] == ["phone"]
assert [i["field"] for i in wang["issues"]] == ["phone"]
# 赵六:两个字段都挂 → errorCells 两个都在
zhao = next(r for r in body["rows"] if r["person_name"] == "赵六")
assert zhao["errorCells"] == ["id_card", "phone"]
def test_show_default_field_without_rules_is_selected(client):
"""没规则但勾了「默认展示」的字段也要 SELECT 出来,否则结果表这列全空。"""
_run(client)
sql = _FakeDBConnection.last_sql
assert "`person_name`" in sql, f"show_default 字段应进 SELECT:{sql}"
assert "`id_card`" in sql and "`phone`" in sql
def test_full_table_scan_no_limit_clause(client):
"""全表扫描:SQL 不该带 LIMIT。"""
_run(client)
sql = _FakeDBConnection.last_sql
assert "LIMIT" not in sql.upper(), f"应全表扫描,SQL 不该有 LIMIT:{sql}"
def test_max_rows_truncates_returned_rows_but_not_counts(client):
"""max_rows 只截断返回的行,不截断扫描 —— scanned / bad_rows 仍是全表真值。"""
body = _run(client, max_rows=2)
assert len(body["rows"]) == 2, "返回行数应被 max_rows 截断"
assert body["truncated"] is True
assert body["scanned"] == 5, "扫描不受 max_rows 影响,仍是全表 5 行"
assert body["bad_rows"] == 3, "计数不受 max_rows 影响,仍是全表 3 行不合规"
assert body["total"] == 3
assert body["message"], "被截断时应给出提示文案"
def test_not_truncated_when_under_max_rows(client):
"""没超上限时 truncated=False,且不带截断提示。"""
body = _run(client, max_rows=5000)
assert body["truncated"] is False
assert body["message"] is None
assert len(body["rows"]) == body["bad_rows"] == 3
def test_scans_beyond_old_1000_row_limit(client):
"""回归:旧实现固定 LIMIT 1000,只校验前 1000 行。现在要扫完整张表。"""
good = {"id_card": "110101199003074512", "phone": "13800138000", "person_name": "合规"}
bad = {"id_card": "x", "phone": "y", "person_name": "第2500行的脏数据"}
_FakeDBConnection.rows = [dict(good) for _ in range(2499)] + [bad]
body = _run(client)
assert body["scanned"] == 2500, "应扫完 2500 行,不是停在 1000"
assert body["bad_rows"] == 1
assert body["rows"][0]["person_name"] == "第2500行的脏数据", "第 1000 行之后的脏数据也要被查出来"
def test_row_index_points_at_position_in_full_scan(client):
"""__row_index 记的是「全表第几行」,截断后也不能错位。"""
body = _run(client, max_rows=1)
# 李四是全表第 2 行(0-based = 1)
assert body["rows"][0]["__row_index"] == 1
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