Commit 7b4bd4a8 authored by Data Governance Dev's avatar Data Governance Dev

feat(web): Phase 3 - Web 端治理平台(FastAPI + Vue 3)

把工作流包装成浏览器化的产品:填表 → 测试连接 → 一键治理 → 实时进度。

核心架构
- web/app.py         FastAPI 入口 + CORS + 静态资源挂载
- web/start.py       启动前依赖检查(可选 anthropic / dmPython)
- web/start.bat      Windows 一键启动脚本
- web/api/routes.py  12 个 REST 接口(jobs/history/standards/reports...)

调度与运行
- web/core/orchestrator.py 8 步流程调度 + 步骤级日志 + 异常 traceback
- web/core/job_manager.py  后台任务管理 + LogBuffer + SSE 订阅
- web/core/db_adapter.py   MySQL / 达梦 统一连接与查询
- web/core/step_impl/      8 个 Step 的 Web 版实现

LLM 增强(可降级)
- web/core/llm.py          Anthropic / OpenAI / MiniMax 多 provider 适配
- configs/llm.example.yaml provider / api_key / base_url / model 模板

SQL 模板体系(Python 不再写 SQL)
- web/sql/info_schema/     list_columns / list_tables  按 dialect 区分 MySQL / 达梦
- web/sql/verify/          count_table_rows
- web/sql/empty_fields/    count_with_nulls
- web/sql/xzqh/            check_dict_table_exists / find_xzqh_orphans
- web/sql/standards/       sample_field_values
- web/sql/health/          check_connection
- web/sql/loader.py        占位符 + dialect 自动加载

前端(无构建步骤)
- web/static/index.html    Vue 3 单页应用
- web/static/lib/          Vue + Element Plus 静态资源
- web/static/style.css

实时日志 / 历史归档 / 报告下载
- SSE 流(GET /api/jobs/{id}/logs/stream)
- outputs/<db>/<ts>/findings/_all_findings.json
- outputs/<db>/<ts>/reports/{md,docx}
- GET /api/history /api/reports/{db}/{ts}/{file}

调试能力(本次一并落地)
- 三路日志:SSE 实时流 + web/logs/app.log + outputs/<db>/<ts>/run.log
- 每步开始 / 完成带耗时、进度 N/M;异常带 traceback 与异常类型
- 数据库 [DB] 连接 / SQL 执行时间;路由 [POST /api/xxx] 入参摘要
- 启动脚本自带 web/logs/start.log

Why:CLI 工作流需要登录机器;Web 端把治理入口放到任何能开浏览器的设备上,
     配合 SSE 实时反馈,业务方也能直接用。
How:复用 workflow/ 中的 8 步逻辑与 standards/ 国标插件;新增 Pydantic 契约
     + job_manager + SSE,前后端用 REST + 事件流通信。
parent 3747e885
# 数据治理 Web 工具
> 基于 FastAPI + Vue 3 + Element Plus 的本地浏览器化治理平台
---
## 🚀 快速开始
### 1. 安装依赖(首次运行)
```bash
pip install -r web/requirements.txt
```
依赖包括:
- `fastapi` / `uvicorn` — Web 服务
- `pydantic` / `sse-starlette` — 数据模型 + SSE 流
- `pymysql` / `dmPython` — MySQL / 达梦驱动
- `anthropic` — Claude API(LLM 增强)
- `python-docx` — Word 报告生成
### 2. 配置 LLM Key(可选)
LLM 用于字段注释推测 / 合并建议 / 长度推断 / 空字段分类等智能判断。
未配置时**自动降级**为规则推理,主流程不受影响。
#### 方式 A:写到配置文件(推荐)
编辑 [web/configs/llm.yaml](web/configs/llm.yaml):
```yaml
provider: minimax # 改为你用的 provider
api_key: sk-请填入你的_MiniMax_API_Key # ← 把这里替换为真实 Key
base_url: https://api.minimaxi.com/anthropic
model: MiniMax-Text-01
```
更多 provider 示例见 [web/configs/llm.example.yaml](web/configs/llm.example.yaml)。
#### 方式 B:环境变量(适合容器/CI)
```bash
# Windows CMD
set LLM_API_KEY=sk-xxxxx
set LLM_PROVIDER=minimax
# Windows PowerShell
$env:LLM_API_KEY="sk-xxxxx"
$env:LLM_PROVIDER="minimax"
# Linux / Git Bash
export LLM_API_KEY=sk-xxxxx
export LLM_PROVIDER=minimax
```
#### 优先级
```
环境变量 > web/configs/llm.yaml > 内置默认值
```
#### 支持的 Provider
| provider | 接入方式 | base_url | 模型示例 |
|----------|---------|----------|---------|
| `minimax` | Anthropic SDK + 自定义 base_url | `https://api.minimaxi.com/anthropic` | `MiniMax-Text-01` / `MiniMax-01` |
| `anthropic` | Anthropic SDK 官方 | (默认) | `claude-sonnet-5` / `claude-haiku-4-5-20251001` |
| `openai` | OpenAI SDK | (默认) | `gpt-4o` / `gpt-4-turbo` |
启动时会打印当前生效的 LLM 配置:
```
[LLM] provider=minimax model=MiniMax-Text-01
[LLM] base_url=https://api.minimaxi.com/anthropic
[LLM] api_key=sk-xxxx...xx(已配置)
```
前端右上角可通过 `GET /api/llm/status` 查看当前配置。
### 3. 启动服务
**方式 A:Python 启动**
```bash
python web/start.py
# 或
python -m web.start
```
**方式 B:Windows 一键脚本**
```cmd
web\start.bat
```
**方式 C:直接 uvicorn**
```bash
python -m uvicorn web.app:app --host 0.0.0.0 --port 8765
```
### 4. 浏览器访问
打开 **http://localhost:8765**
- Web 界面:填表 → 测试连接 → 开始治理
- API 文档:http://localhost:8765/docs
- OpenAPI JSON:http://localhost:8765/openapi.json
---
## 🎯 功能特性
| 特性 | 说明 |
|------|------|
| 🖥️ **Web 界面** | 浏览器访问,Vue 3 + Element Plus 响应式布局 |
| 🔌 **双数据库** | MySQL(PyMySQL)+ 达梦(dmPython)统一接口 |
| 🤖 **LLM 增强** | Claude API 用于字段注释推测 / 合并建议 / 长度推断 / 空字段分类 |
| 📊 **多级表格** | 7 大章节各一张可展开表格,KPI 卡片概览 |
| 📜 **国标展示** | 4 项内置标准可视化展示在页面(无需查询) |
| 📡 **实时日志** | SSE 流式推送治理进度,含 INFO/WARN/ERROR 三色 |
| 📥 **报告下载** | Markdown / Word 双格式,含历史报告检索 |
| 🗂️ **历史归档** | 每次运行按时间戳归档到 `outputs/<db>/<ts>/` |
| 🛡️ **优雅降级** | LLM 失败 / DB 失败不卡死,规则推理兜底 |
---
## 🗂️ 目录结构
```
数据治理/
├── web/ ← Web 工具
│ ├── app.py ← FastAPI 入口
│ ├── start.py ← 启动脚本(带依赖检查)
│ ├── start.bat ← Windows 一键启动
│ ├── requirements.txt
│ ├── README.md ← 本文件
│ ├── sql/ ← 所有 SQL 模板(Python 不再写 SQL)
│ │ ├── loader.py ← 模板加载器
│ │ ├── info_schema/ ← INFORMATION_SCHEMA 查询
│ │ ├── verify/ xzqh/ standards/ empty_fields/ health/
│ ├── core/ ← 核心模块
│ │ ├── db_adapter.py ← MySQL + 达梦 统一接口
│ │ ├── llm.py ← Claude API 封装
│ │ ├── job_manager.py ← 后台任务 + SSE 推送
│ │ ├── orchestrator.py ← 8 步流程调度
│ │ ├── models.py ← Pydantic 契约
│ │ └── step_impl/ ← 8 个 Step 的 Web 版实现
│ ├── api/
│ │ └── routes.py ← 12 个 REST 接口
│ ├── static/
│ │ ├── index.html ← Vue 3 单页前端
│ │ └── style.css
│ ├── configs/
│ │ └── defaults.yaml
│ ├── outputs/ ← 历史治理产出(按 db/ts 分目录)
│ └── logs/
│ └── app.log
│
├── workflow/ ← 原 CLI 工作流(仍可用,独立运行)
│
└── standards/ ← 国标插件(Web 与 CLI 共用)
├── base.py
├── registry.py
├── std_001_id_card.py ← GB 11643-1999
├── std_002_uscc.py ← GB 32100-2015
├── std_003_mobile.py ← YD/T 1313
└── std_004_xzqh.py ← GB/T 2260
```
---
## 🔧 API 端点速查
| 方法 | 路径 | 说明 |
|------|------|------|
| GET | `/api/health` | 健康检查 |
| GET | `/api/steps` | 列出 8 步元信息 |
| GET | `/api/standards` | 列出已注册国标(页面用) |
| POST | `/api/connect/test` | 测试连接(返回表数量) |
| POST | `/api/jobs` | 提交治理任务 |
| GET | `/api/jobs/{id}` | 查询任务状态(含进度百分比) |
| DELETE | `/api/jobs/{id}` | 取消任务 |
| GET | `/api/jobs/{id}/result` | 获取完整结果 |
| GET | `/api/jobs/{id}/logs/stream` | SSE 日志流 |
| GET | `/api/history` | 历史治理记录 |
| GET | `/api/reports/{db}/{ts}/{file}` | 下载报告(md/docx) |
| GET | `/api/history/{db}/{ts}/result` | 历史结果数据 |
---
## 🤖 LLM 增强位置
| Step | LLM 用途 | 失败降级 |
|------|----------|----------|
| Step 2 | 表合并建议文案 + 风险评估 | 模板字符串 |
| Step 4 | 高空字段成因分类 | 跳过 LLM 分类 |
| Step 5 | 无规则命中的字段注释推测 | 仅靠规则映射表 |
| Step 6 | 自定义字段长度推断 | 标记"待业务确认" |
LLM 调用封装在 [web/core/llm.py](web/core/llm.py),所有方法都设计为**可降级**:
未配置 `ANTHROPIC_API_KEY` 时自动跳过,不阻断主流程。
---
## 🛡️ 双数据库适配说明
通过 [web/core/db_adapter.py](web/core/db_adapter.py) 统一抽象:
| 项 | MySQL | 达梦 |
|----|-------|------|
| 驱动 | `pymysql` | `dmPython` |
| 占位符 | `%s` | `:name` (NamedParam) |
| 标识符引用 | `` ` `` (反引号) | `"` (双引号 ANSI) |
| 字段大小写 | 视配置 | 默认大写 → 统一规范化为小写 |
| 长文本类型 | `text/longtext/...` | `text/clob/longvarchar` |
| `INFORMATION_SCHEMA` | 完全支持 | 完全支持(字段名差异已处理) |
> ⚠️ 达梦驱动 `dmPython` 需要本地先安装达梦客户端,否则 `pip install dmPython` 会失败。
---
## ⚠️ 已知事项
1. **MySQL 兼容**:所有 `INFORMATION_SCHEMA` 查询已用 ANSI 写法 + 适配层规范化,达梦可直接复用
2. **报告 bug 修复**:新版 Web 工具修正了原 `reporter.py` 的多处字段名不匹配(如 `over_provision_ratio`、`risk/suggestion` 等)
3. **单进程**:JobManager 使用进程内 asyncio 队列,并发任务排队执行;如需多任务并行可改为 Celery
4. **数据安全**:数据库密码仅在内存中传递,不写日志;提交前可勾选「启用 LLM」控制是否走外网
---
## 🗃️ SQL 模板管理
所有 SQL 语句都集中在 `web/sql/` 目录下,**Python 代码不再直接拼 SQL**。
### 模板语法(占位符)
| 语法 | 含义 | 示例 |
|------|------|------|
| `${var}` | 简单字符串替换 | `${schema}` |
| `${var \| quote}` | 按当前 dialect 自动加引号 | `${table \| quote}` → MySQL 反引号 / 达梦双引号 |
| `${list \| join:","}` | 列表按分隔符拼接 | `${cols \| join:","}` |
### Dialect 加载规则
加载器优先找 `<name>.<dialect>.sql`,找不到则用 `<name>.sql` 兜底。
### 使用示例
```python
from web.sql.loader import get_sql_loader
loader = get_sql_loader()
sql = loader.render("info_schema/list_columns", dialect="mysql", schema="mydb")
# SELECT ... FROM INFORMATION_SCHEMA.COLUMNS c WHERE c.TABLE_SCHEMA = `mydb`
sql = loader.render("info_schema/list_columns", dialect="dameng", schema="mydb")
# SELECT ... FROM ALL_TAB_COLUMNS c WHERE c.OWNER = "mydb"
```
详细文档:[sql/README.md](sql/README.md)
---
## 📝 后续可扩展
- [ ] 添加更多国标插件(邮箱、银行卡、邮编等)
- [ ] 接入 jaydebeapi 作为达梦 JDBC 替代驱动
- [ ] 报告编辑器(在线编辑 Markdown / 一键导出)
- [ ] 多任务并行(Celery + Redis)
- [ ] 用户系统(FastAPI Users + JWT)
- [ ] Docker 镜像打包
---
## 📚 相关文档
- [docs/DESIGN.md](docs/DESIGN.md) — **完整设计书**(架构、模块、LLM 判断矩阵)
- [../CLAUDE.md](../CLAUDE.md) — 顶层项目目标
\ No newline at end of file
"""web 包入口"""
__version__ = "1.0.0"
\ No newline at end of file
"""web.api 路由模块"""
from .routes import router
__all__ = ["router"]
\ No newline at end of file
This diff is collapsed.
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""Web 工具主入口
启动方式:
python -m web.app
# 或:
cd web && python app.py
环境变量:
ANTHROPIC_API_KEY Claude API Key(启用 LLM 增强)
WEB_PORT 端口(默认 8765)
WEB_HOST 主机(默认 0.0.0.0)
"""
from __future__ import annotations
import logging
import os
import sys
from pathlib import Path
# ── 把项目根目录加入 sys.path ──
WEB_DIR = Path(__file__).resolve().parent
PROJECT_ROOT = WEB_DIR.parent
sys.path.insert(0, str(PROJECT_ROOT))
from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse, HTMLResponse
from fastapi.staticfiles import StaticFiles
# ── 日志 ──
LOG_DIR = WEB_DIR / "logs"
LOG_DIR.mkdir(parents=True, exist_ok=True)
logging.basicConfig(
level=os.environ.get("LOG_LEVEL", "INFO"),
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
handlers=[
logging.FileHandler(LOG_DIR / "app.log", encoding="utf-8"),
logging.StreamHandler(),
],
)
logger = logging.getLogger(__name__)
logger.info("=" * 60)
logger.info("数据治理 Web 工具 · 进程启动")
logger.info(f"Python: {sys.version.split()[0]}, 平台: {sys.platform}")
logger.info(f"工作目录: {PROJECT_ROOT}")
logger.info(f"日志目录: {LOG_DIR}")
# ── 创建 FastAPI 应用 ──
app = FastAPI(
title="数据治理 Web 工具",
description="基于 FastAPI + Vue 3 的数据库治理平台(支持 MySQL / 达梦)",
version="1.0.0",
)
logger.info("FastAPI 应用已创建")
# CORS(开发模式)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
logger.debug("CORS 中间件已挂载(允许所有来源)")
# ── 注册路由 ──
from web.api import router as api_router
app.include_router(api_router)
logger.info(f"已注册 {len(app.routes)} 个路由")
# ── 启动时打印 LLM 配置(方便确认 Key 是否生效) ──
@app.on_event("startup")
async def show_llm_config():
logger.info("FastAPI 应用启动事件触发")
from web.core.llm import LLMConfig
cfg = LLMConfig.load()
if cfg.available:
masked_key = cfg.api_key[:6] + "..." + cfg.api_key[-4:] if len(cfg.api_key) > 12 else "***"
print(f"[LLM] provider={cfg.provider} model={cfg.effective_model}")
print(f"[LLM] base_url={cfg.effective_base_url or '(default)'}")
print(f"[LLM] api_key={masked_key}(已配置)")
logger.info(f"LLM 配置: provider={cfg.provider}, model={cfg.effective_model}")
else:
print("[LLM] 未配置 API Key → LLM 增强不可用,将使用规则推理")
logger.warning("LLM 未配置 API Key → LLM 增强不可用,将使用规则推理")
@app.on_event("shutdown")
async def shutdown_log():
logger.info("FastAPI 应用关闭")
# ── 静态资源 / 入口页 ──
STATIC_DIR = WEB_DIR / "static"
if STATIC_DIR.exists():
app.mount("/static", StaticFiles(directory=str(STATIC_DIR)), name="static")
@app.get("/", response_class=HTMLResponse)
async def root():
"""服务入口 HTML"""
index = STATIC_DIR / "index.html"
if index.exists():
return HTMLResponse(index.read_text(encoding="utf-8"))
return HTMLResponse("<h1>前端文件未找到</h1><p>请检查 web/static/index.html</p>")
# ── 启动 ──
def main():
import uvicorn
host = os.environ.get("WEB_HOST", "0.0.0.0")
port = int(os.environ.get("WEB_PORT", "8765"))
print("=" * 60)
print(f" 数据治理 Web 工具")
print(f" 访问: http://localhost:{port}")
print(f" API 文档: http://localhost:{port}/docs")
print("=" * 60)
uvicorn.run("web.app:app", host=host, port=port, reload=False, log_level="info")
if __name__ == "__main__":
main()
\ No newline at end of file
# ============================================================
# Web 工具默认配置
# 用户通过前端表单提交的内容会覆盖这里的连接信息
# ============================================================
# 数据库类型选项
db_types:
- mysql
- dameng
# 采样与阈值(与 CLI 工作流一致)
thresholds:
empty_field:
high: 0.80
mid: 0.50
min_rows_to_check: 5
similarity:
merge_threshold: 0.80
# 标准库
standards:
enabled: true
auto_apply: true
strict_mode: false
# LLM(Claude API)
llm:
model: claude-sonnet-5 # 默认模型
max_tokens: 1024
temperature: 0.2
# API Key 从环境变量 ANTHROPIC_API_KEY 读取,不在此保存
# Web 服务
web:
host: 0.0.0.0
port: 8765
title: "数据治理工具"
\ No newline at end of file
# ============================================================
# LLM 配置文件示例 — 多种 provider 用法
# ============================================================
# 复制本文件为 llm.yaml,按需修改。
# ============================================================
# ── 方式 1:MiniMax(推荐本项目用户使用) ──
provider: minimax
api_key: sk-请填入你的_MiniMax_API_Key
base_url: https://api.minimaxi.com/anthropic
model: MiniMax-Text-01
# ── 方式 2:Anthropic Claude ──
# provider: anthropic
# api_key: sk-ant-xxxxx
# model: claude-sonnet-5
# (base_url 留空,用 Anthropic 官方)
# ── 方式 3:OpenAI ──
# provider: openai
# api_key: sk-xxxxx
# model: gpt-4o
# base_url: https://api.openai.com/v1
# ── 方式 4:自部署 / 第三方代理(OpenAI 兼容) ──
# provider: openai
# api_key: 自定义
# base_url: https://your-proxy.com/v1
# model: your-model-name
# 生成参数(通用)
max_tokens: 1024
temperature: 0.2
timeout: 30
\ No newline at end of file
"""web.core 核心模块
- db_adapter: MySQL / 达梦 适配
- llm: Claude / OpenAI 客户端
- job_manager: 后台任务 + SSE
- orchestrator: 调度 8 步流程
- models: Pydantic 数据模型
"""
from .db_adapter import DBConfig, DBConnection, open_db, test_connection, quote_ident
from .llm import LLMClient, LLMConfig, get_llm_client, LLMUnavailable
__all__ = [
"DBConfig",
"DBConnection",
"open_db",
"test_connection",
"quote_ident",
"LLMClient",
"LLMConfig",
"get_llm_client",
"LLMUnavailable",
]
\ No newline at end of file
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
"""Pydantic 数据模型 — Web API 的请求/响应契约"""
from __future__ import annotations
from datetime import datetime
from typing import Any, Optional
from pydantic import BaseModel, Field
# ── 请求 ────────────────────────────────────────────────
class ConnectRequest(BaseModel):
"""前端「新建治理任务」表单提交"""
db_type: str = Field(..., description="mysql / dameng")
host: str
port: int = 3306
user: str
password: str
database: str
charset: str = "utf8mb4"
connect_timeout: int = 10
# 可选:选择要跑的步骤(默认全 8 步)
steps: Optional[list[int]] = None
# 可选:是否启用 LLM 增强(默认 True,缺 Key 时自动降级)
enable_llm: bool = True
# 可选:报告标题
report_title: Optional[str] = None
class TestConnectionRequest(ConnectRequest):
"""只测连通性,不启动任务"""
pass
# ── 响应 ────────────────────────────────────────────────
class TestConnectionResponse(BaseModel):
ok: bool
message: str
db_type: Optional[str] = None
database_name: Optional[str] = None # 改名为避免与 BaseModel.schema() 冲突
tables_count: Optional[int] = None
class JobCreatedResponse(BaseModel):
job_id: str
status: str
created_at: datetime
steps_planned: list[int]
class JobStatusResponse(BaseModel):
job_id: str
status: str # pending / running / completed / failed / cancelled
created_at: datetime
started_at: Optional[datetime] = None
finished_at: Optional[datetime] = None
current_step: Optional[int] = None
current_step_title: Optional[str] = None
progress_pct: float = 0.0 # 0-100
steps_planned: list[int] = []
steps_completed: list[int] = []
steps_failed: list[int] = []
error: Optional[str] = None
db_info: Optional[dict] = None
class StepInfo(BaseModel):
num: int
title: str
description: str
requires_db: bool
llm_enhanced: bool = False
class StepsListResponse(BaseModel):
steps: list[StepInfo]
class StandardInfo(BaseModel):
id: str
name: str
applies_to_fields: list[str]
description: str
class StandardsListResponse(BaseModel):
standards: list[StandardInfo]
# ── 治理结果(前端多级表格消费) ──
class OverviewStats(BaseModel):
database: str
host: str
db_type: str
total_tables: int = 0
total_fields: int = 0
high_empty_fields: int = 0
missing_comments: int = 0
length_issues: int = 0
standard_violations: int = 0
merge_candidates: int = 0
redundancy_fields: int = 0
executed_at: Optional[datetime] = None
class JobResultResponse(BaseModel):
job_id: str
status: str
overview: OverviewStats
# 各章节原始数据(前端按章节渲染多级表格)
sections: dict[str, Any] = {}
# 已生成报告路径
reports: dict[str, Optional[str]] = {}
class LogEntry(BaseModel):
ts: datetime
level: str
step: Optional[str] = None
message: str
class JobLogsResponse(BaseModel):
job_id: str
entries: list[LogEntry]
class RunHistoryItem(BaseModel):
run_id: str # 时间戳目录名
database: str
host: str
db_type: Optional[str] = "mysql"
executed_at: datetime
duration_sec: Optional[float] = None
stats: dict[str, Any] = {}
has_markdown: bool
has_docx: bool
class RunHistoryResponse(BaseModel):
items: list[RunHistoryItem]
\ No newline at end of file
This diff is collapsed.
"""Step 实现包 — Web 版(兼容 MySQL + 达梦,接 LLM)
每个 step_* 函数接收通用 cfg + 上游 step_outputs,返回结构化 dict。
所有 Step 通过 web.core.orchestrator.run_governance_workflow() 调度。
"""
\ No newline at end of file
"""Step 1: 获取数据字典(MySQL + 达梦 适配版)
通过 db_adapter 统一接口获取:
- data_dictionary: 全部字段元数据
- table_summary: 表级汇总
- overview_meta: 概览统计
"""
from __future__ import annotations
import logging
from typing import Callable
from ..db_adapter import DBConfig, open_db
logger = logging.getLogger(__name__)
def run_step1(cfg: DBConfig, log: Callable | None = None) -> dict:
"""执行 Step 1,返回 data_dictionary + table_summary + overview_meta"""
if log:
log("INFO", f"连接数据库: {cfg.db_type}://{cfg.user}@{cfg.host}:{cfg.port}/{cfg.database}", step="1")
with open_db(cfg) as db:
if log:
log("INFO", "已建立连接,开始读取 information_schema.COLUMNS ...", step="1")
columns = db.list_columns(cfg.database)
if log:
log("INFO", f" · COLUMNS: {len(columns)} 行", step="1")
if log:
log("INFO", "读取 information_schema.TABLES ...", step="1")
tables = db.list_tables(cfg.database)
if log:
log("INFO", f" · TABLES: {len(tables)} 行", step="1")
# 类型规整
fixed_int = 0
for r in columns:
for k in ("char_max_length", "numeric_precision", "numeric_scale", "ordinal_position"):
if r.get(k) is not None:
try:
r[k] = int(r[k])
fixed_int += 1
except (TypeError, ValueError):
pass
if log and fixed_int:
log("DEBUG", f" · 字段类型规整: 转换 {fixed_int} 个数值字段", step="1")
for r in tables:
for tf in ("create_time", "update_time"):
if r.get(tf) is not None and hasattr(r[tf], "strftime"):
r[tf] = r[tf].strftime("%Y-%m-%d %H:%M:%S")
for nf in ("table_rows", "data_length", "index_length"):
if r.get(nf) is not None:
try:
r[nf] = int(r[nf])
except (TypeError, ValueError):
pass
# 概览
table_names = {r["table_name"] for r in columns}
if log:
log("INFO", f"获取完成: {len(table_names)} 张表, {len(columns)} 个字段(覆盖 {len(tables)} 条表信息)", step="1")
overview_meta = {
"database": cfg.database,
"host": cfg.host,
"db_type": cfg.db_type,
"total_tables": len(table_names),
"total_fields": len(columns),
}
return {
"data_dictionary": columns,
"table_summary": tables,
"overview_meta": overview_meta,
}
\ No newline at end of file
"""Step 2: 表合并候选与冗余字段分析(离线 + LLM 增强)
基于 Step 1 产出,计算:
- merge_candidates: 结构高度相似的表组
- decommission_candidates: 0 行表 + 可能的新表替代
- redundancy_fields: 高频出现的疑似冗余字段
- 每个 merge_candidate 的 suggestion 由 LLM 优化(可选)
"""
from __future__ import annotations
import logging
from collections import Counter, defaultdict
from typing import Callable
from ..llm import LLMClient
logger = logging.getLogger(__name__)
COMMON_FIELDS = {
"id", "create_time", "update_time", "create_by", "update_by",
"remark", "del_flag", "tenant_id", "dept_id", "status",
"sort_order", "version",
}
NEW_PREFIXES = ["t_talent_", "t_ai_", "t_design_"]
OLD_PREFIXES = ["t_", "mall_", "project_"]
def run_step2(dict_data: dict, llm: LLMClient | None = None,
log: Callable | None = None) -> dict:
columns = dict_data.get("data_dictionary", [])
table_summary = dict_data.get("table_summary", [])
if not columns:
if log:
log("WARN", "未获取到任何字段元数据,跳过 Step 2", step="2")
return {"merge_candidates": [], "decommission_candidates": {}, "redundancy_fields": []}
by_table: dict[str, list[dict]] = defaultdict(list)
for row in columns:
by_table[row["table_name"]].append(row)
if log:
log("INFO", f"待分析表数: {len(by_table)}, 字段数: {len(columns)}", step="2")
# 1. 合并候选
if log:
log("INFO", "[1/3] 计算合并候选(相似度阈值 ≥95%)...", step="2")
merge = _find_merge(by_table, llm, log)
if log:
log("INFO", f" · 发现 {len(merge)} 组可合并表", step="2")
# 2. 废弃候选
if log:
log("INFO", "[2/3] 计算废弃候选(0 行表)...", step="2")
rows_by_table = {t["table_name"]: t.get("table_rows", 0) for t in table_summary}
decom = _find_decommission(rows_by_table, by_table.keys())
if log:
log("INFO", f" · 强废弃(有新版替代){decom['strong']['count']} 张, 弱废弃(无明确替代){decom['weak']['count']} 张", step="2")
# 3. 高频字段
if log:
log("INFO", "[3/3] 计算高频字段(出现 ≥10 张表)...", step="2")
redundancy = _find_redundancy(by_table)
if log:
log("INFO", f" · 高频字段 {len(redundancy)} 个", step="2")
if log:
log("INFO",
f"合并候选 {len(merge)} 组, 废弃候选 {decom['strong']['count'] + decom['weak']['count']} 张, "
f"高频字段 {len(redundancy)} 个",
step="2")
return {
"merge_candidates": merge,
"decommission_candidates": decom,
"redundancy_fields": redundancy,
}
def _table_columns(cols: list[dict]) -> set[str]:
return {c["column_name"] for c in cols}
def _similarity(a: set[str], b: set[str]) -> float:
if not a or not b:
return 0.0
return len(a & b) / len(a | b)
def _find_merge(by_table: dict, llm: LLMClient | None, log: Callable | None) -> list[dict]:
"""找出结构高度相似的表组(相似度 ≥95%)"""
names = sorted(by_table.keys())
visited = set()
groups = []
for i, t in enumerate(names):
if t in visited:
continue
ci = _table_columns(by_table[t])
cluster = [t]
for j in range(i + 1, len(names)):
t2 = names[j]
if t2 in visited:
continue
cj = _table_columns(by_table[t2])
if _similarity(ci, cj) >= 0.95:
cluster.append(t2)
if len(cluster) >= 2:
visited.update(cluster)
groups.append(cluster)
results = []
for idx, group in enumerate(groups, 1):
ref = _table_columns(by_table[group[0]])
diff = _table_columns(by_table[group[1]]) - ref if len(group) > 1 else set()
# 基础建议(模板化)
suggestion = _template_suggestion(group)
risk = ""
llm_verdict = None
# LLM 增强(如启用)
if llm and llm.available:
try:
if log:
log("INFO", f"调用 LLM 评估合并候选 MERGE-{idx:03d}", step="2")
llm_result = llm.suggest_table_merge(
tables=[
{"name": t, "comment": by_table[t][0].get("table_comment", ""),
"row_count": "?"}
for t in group
],
similarity=0.95,
common_columns=sorted(ref),
diff_columns=sorted(diff),
)
if llm_result:
llm_verdict = llm_result.get("verdict")
if llm_result.get("suggestion"):
suggestion = llm_result["suggestion"]
if llm_result.get("risks"):
risk = "; ".join(llm_result["risks"])
except Exception as e:
if log:
log("WARN", f"LLM 评估失败,使用规则建议: {e}", step="2")
results.append({
"id": f"MERGE-{idx:03d}",
"tables": group,
"common_columns": len(ref),
"similarity": "≥95%",
"similarity_value": 0.95,
"columns": sorted(ref),
"diff_columns": sorted(diff),
"suggestion": suggestion,
"risk": risk,
"verdict": llm_verdict or "merge",
"priority": "P1" if len(group) >= 3 else "P2",
})
return results
def _template_suggestion(group: list[str]) -> str:
if not group:
return "合并候选"
if len(group) == 2:
return f"合并为 {group[0]},新增 type 字段区分"
if len(group) <= 5:
return f"合并为 {group[0]},新增 record_type 字段区分"
return f"合并 {len(group)} 张同类表"
def _find_decommission(rows_by_table: dict, all_tables) -> dict:
zero_rows = [t for t, n in rows_by_table.items() if n == 0 and t in all_tables]
strong, weak = [], []
for t in sorted(zero_rows):
replacement = _find_replacement(t, all_tables)
entry = {"table": t, "rows": 0, "replacement": replacement}
(strong if replacement else weak).append(entry)
return {
"strong": {"description": "有新版表替代", "tables": strong, "count": len(strong)},
"weak": {"description": "无明确替代", "tables": weak, "count": len(weak)},
}
def _find_replacement(table: str, all_tables) -> str | None:
for prefix in OLD_PREFIXES:
if table.startswith(prefix):
stem = table[len(prefix):]
for new_prefix in NEW_PREFIXES:
cand = new_prefix + stem
if cand in all_tables and cand != table:
return cand
return None
def _find_redundancy(by_table: dict) -> list[dict]:
counter: Counter = Counter()
for cols in by_table.values():
for c in cols:
counter[c["column_name"]] += 1
logger.debug(f"字段频次统计: 共 {len(counter)} 个不同字段名")
return [
{
"field": f,
"table_count": n,
"risk": "高频出现,需评估是否为业务必要字段" if n >= 30 else "中频出现",
"suggestion": "建议评审是否需要统一到公共字典表" if n >= 30 else "建议评审",
}
for f, n in counter.most_common()
if n >= 10 and f not in COMMON_FIELDS
]
\ No newline at end of file
"""Step 3: 数据验证(连库)
连库验证:
- merge_candidates: 对每组前 5 张表跑 SELECT COUNT(*)
- redundancy_fields: 取高频字段,看实际出现在哪些表
- xzqh_orphan_check: 行政区划编码孤儿(如果存在字典表 c_bri_xzqh)
- data_quality_issues: 数据质量问题清单
所有 SQL 走 web/sql/ 模板,通过 db_adapter 调用。
"""
from __future__ import annotations
import logging
from typing import Callable
from ..db_adapter import DBConfig, open_db, quote_ident
from web.sql.loader import get_sql_loader
logger = logging.getLogger(__name__)
def run_step3(cfg: DBConfig, dict_data: dict,
merge_candidates: list, redundancy_fields: list,
log: Callable | None = None) -> dict:
confirmed_merges = []
redundancy_verified = []
xzqh_check = {"orphan_total": 0, "per_table": []}
quality_issues = []
if log:
log("INFO", f"[1/4] 验证合并候选(top {min(5, len(merge_candidates))})...", step="3")
try:
with open_db(cfg) as db:
loader = get_sql_loader()
# 1. 验证 merge 候选(实际行数)
for idx, mc in enumerate(merge_candidates[:5], 1):
row_counts = []
total = 0
for t in mc.get("tables", []):
try:
sql = loader.render(
"verify/count_table_rows",
dialect=cfg.db_type,
table=t, # 模板内 ${table | quote} 自动加引号
)
n = db.fetch_scalar(sql) or 0
row_counts.append({"table": t, "rows": int(n)})
total += int(n)
except Exception as e:
row_counts.append({"table": t, "error": str(e)})
if log:
log("WARN", f" · 表 {t} 行数统计失败: {e}", step="3")
confirmed_merges.append({
**mc,
"row_counts": row_counts,
"total_rows": total,
})
if log:
log("INFO",
f" · {mc.get('id')} ({idx}/{min(5, len(merge_candidates))}): "
f"{len(mc.get('tables', []))} 张表, 共 {total} 行",
step="3")
# 2. 验证冗余字段(出现表数 + 实际行数)
if log:
log("INFO", f"[2/4] 验证冗余字段(top {min(5, len(redundancy_fields))})...", step="3")
for rf in redundancy_fields[:5]:
field_name = rf.get("field")
tables = _find_field_in_tables(dict_data.get("data_dictionary", []), field_name)
redundancy_verified.append({
"field": field_name,
"occurrences": len(tables),
"tables_with_field": tables,
})
if log:
log("INFO", f" · 字段 {field_name} 出现在 {len(tables)} 张表", step="3")
# 3. 行政区划孤儿编码检查
if log:
log("INFO", "[3/4] 行政区划孤儿编码检查...", step="3")
xzqh_check = _check_xzqh_orphans(cfg, db, log)
# 4. 数据质量问题
if log:
log("INFO", "[4/4] 数据质量问题检查(占位)...", step="3")
quality_issues = []
except Exception as e:
if log:
log("ERROR", f"Step 3 连库验证失败: {type(e).__name__}: {e}", step="3")
logger.exception("Step 3 详细异常")
if log:
log("INFO",
f"Step 3 完成: 验证合并 {len(confirmed_merges)} 组, "
f"冗余字段 {len(redundancy_verified)} 个, 行政区划孤儿 {xzqh_check.get('orphan_total', 0)} 条",
step="3")
return {
"confirmed_merge_candidates": confirmed_merges,
"redundancy_fields": redundancy_verified,
"xzqh_orphan_check": xzqh_check,
"data_quality_issues": quality_issues,
}
def _find_field_in_tables(columns: list, field_name: str) -> list[str]:
return sorted({r["table_name"] for r in columns if r["column_name"] == field_name})
def _check_xzqh_orphans(cfg: DBConfig, db, log: Callable | None) -> dict:
"""检查行政区划字段是否存在孤儿编码(不在字典表 c_bri_xzqh 中)"""
DICT_TABLE = "c_bri_xzqh"
result = {"dict_table_rows": 0, "orphan_total": 0, "per_table": []}
loader = get_sql_loader()
# 检查字典表是否存在 + 行数
try:
sql = loader.render(
"xzqh/check_dict_table_exists",
dialect=cfg.db_type,
dict_table=DICT_TABLE, # 模板内 ${dict_table | quote} 自动加引号
)
result["dict_table_rows"] = int(db.fetch_scalar(sql) or 0)
except Exception:
if log:
log("INFO", f"字典表 {DICT_TABLE} 不存在,跳过行政区划校验", step="3")
return result
if result["dict_table_rows"] == 0:
return result
if log:
log("INFO", f"字典表 {DICT_TABLE} 存在,共 {result['dict_table_rows']} 行", step="3")
# 这里可扩展:对含 xzqhbm 字段的前 5 张大表跑 LEFT JOIN 找孤儿
# 模板:xzqh/find_xzqh_orphans.sql
return result
\ No newline at end of file
This diff is collapsed.
"""Step 5: 缺失注释字段检查 + LLM 推测
优先用 LLM 推测无注释字段的语义(替换原硬编码 COMMENT_HINTS);
LLM 不可用时降级为本地规则。
产出:
- summary: 总数 + 推测命中率
- by_table: 每表缺失注释数量
- predicted_comments: 推测出的注释(含置信度)
- unpredictable_sample: 未推测到的样本
"""
from __future__ import annotations
import logging
from collections import defaultdict
from typing import Callable
from ..llm import LLMClient
logger = logging.getLogger(__name__)
# 兜底规则(LLM 不可用时使用)
FALLBACK_HINTS = {
"id": "主键ID",
"create_time": "创建时间",
"update_time": "更新时间",
"create_by": "创建人",
"update_by": "更新人",
"remark": "备注",
"del_flag": "删除标记",
"tenant_id": "租户ID",
"dept_id": "部门ID",
"project_id": "项目ID",
"project_name": "项目名称",
"project_code": "项目编号",
"site_id": "工地ID",
"site_name": "工地名称",
"status": "状态",
"sort_order": "排序",
"start_time": "开始时间",
"end_time": "结束时间",
"type": "类型",
"code": "编码",
}
def run_step5(dict_data: dict, llm: LLMClient | None = None,
log: Callable | None = None) -> dict:
columns = dict_data.get("data_dictionary", [])
if not columns:
if log:
log("WARN", "未获取到任何字段元数据,跳过 Step 5", step="5")
return {"summary": {}, "by_table": [], "predicted_comments": [], "unpredictable_sample": []}
if log:
log("INFO", f"待检查字段总数: {len(columns)}", step="5")
missing = []
by_table: dict[str, int] = defaultdict(int)
predicted = []
unpredictable = []
# 先用规则快速批匹配
for r in columns:
comment = (r.get("column_comment") or "").strip()
if comment:
continue
by_table[r["table_name"]] += 1
fallback = FALLBACK_HINTS.get(r["column_name"])
entry = {
"table_name": r["table_name"],
"table_comment": r.get("table_comment", ""),
"column_name": r["column_name"],
"data_type": r.get("data_type", ""),
"column_type": r.get("column_type", ""),
"is_nullable": r.get("is_nullable", ""),
"predicted": None,
"confidence": "low",
"reason": "",
}
if fallback:
entry["predicted"] = fallback
entry["confidence"] = "high"
entry["reason"] = "字段名匹配内置规则"
predicted.append(entry)
else:
unpredictable.append(entry)
missing.append(entry)
if log:
log("INFO",
f" · 缺失注释字段: {len(missing)} 个 (覆盖 {len(by_table)} 张表)",
step="5")
log("INFO",
f" · 内置规则命中: {len(predicted)}, 需 LLM 推测: {len(unpredictable)}",
step="5")
# 用 LLM 处理 unpredictable(如果可用)
if llm and llm.available and unpredictable:
if log:
log("INFO", f"[LLM] 调用 LLM 推测 {len(unpredictable)} 个无规则命中的字段注释", step="5")
llm_predicted = []
still_unknown = []
for idx, entry in enumerate(unpredictable, 1):
try:
r = llm.predict_field_comment(
table_name=entry["table_name"],
table_comment=entry["table_comment"],
column_name=entry["column_name"],
data_type=entry["data_type"],
)
if r and r.get("comment"):
entry["predicted"] = r["comment"]
entry["confidence"] = r.get("confidence", "low")
entry["reason"] = "LLM 推测"
llm_predicted.append(entry)
# 从 predicted 列表的视角也算推测成功
predicted.append(entry)
if log and idx % 5 == 0:
log("DEBUG",
f" · LLM 推测进度 {idx}/{len(unpredictable)} "
f"(已成功 {len(llm_predicted)})",
step="5")
else:
still_unknown.append(entry)
except Exception as e:
if log:
log("WARN",
f"LLM 推测失败 ({entry['table_name']}.{entry['column_name']}): {e}",
step="5")
still_unknown.append(entry)
if log:
log("INFO", f"[LLM] 推测成功 {len(llm_predicted)} / {len(unpredictable)}", step="5")
unpredictable = still_unknown
# 按表聚合
by_table_list = sorted(
[{"table_name": t, "missing_count": c} for t, c in by_table.items()],
key=lambda x: x["missing_count"], reverse=True
)[:20]
if log:
log("INFO",
f"缺失注释 {len(missing)} 个字段,覆盖 {len(by_table)} 张表;"
f"推测成功 {len(predicted)}(规则 {sum(1 for p in predicted if p.get('reason') == '字段名匹配内置规则')}, "
f"LLM {sum(1 for p in predicted if p.get('reason') == 'LLM 推测')})",
step="5")
return {
"summary": {
"total_missing_comments": len(missing),
"tables_affected": len(by_table),
"predicted_count": len(predicted),
"predicted_by_rules": sum(1 for p in predicted if p.get("reason") == "字段名匹配内置规则"),
"predicted_by_llm": sum(1 for p in predicted if p.get("reason") == "LLM 推测"),
"unpredicted_count": len(unpredictable),
},
"by_table": by_table_list,
"predicted_comments": predicted[:50],
"unpredictable_sample": unpredictable[:50],
}
\ No newline at end of file
"""Step 6: 字段长度检查(离线 + LLM 增强)
识别字段定义长度超出标准所需(如身份证号 18 位用 varchar(50) 存储)。
基础规则覆盖常见字段;自定义字段由 LLM 推断建议长度。
产出:
- summary: 总异常数
- issues: 每条具体异常(含浪费字节估算)
"""
from __future__ import annotations
import logging
import re
from typing import Callable
from ..llm import LLMClient
logger = logging.getLogger(__name__)
# 字段名模式 → 期望长度的规则
# 注意:期望长度需严格匹配国家标准
LENGTH_RULES: list[tuple[re.Pattern, int, str, str]] = [
# (pattern, expected_length, standard, description)
(re.compile(r"^(id_?card|id_?number|identity_?card)$", re.I), 18, "GB 11643-1999", "身份证号 18 位"),
(re.compile(r"^(id_?card_?no|_?sfz_?hm)$", re.I), 18, "GB 11643-1999", "身份证号 18 位"),
(re.compile(r"^(uscc|credit_?code|social_?credit_?code)$", re.I), 18, "GB 32100-2015", "统一社会信用代码 18 位"),
(re.compile(r"^mobile(_?phone)?$", re.I), 11, "YD/T 1313", "手机号 11 位"),
(re.compile(r"^phone(_?no)?$", re.I), 11, "YD/T 1313", "手机号 11 位"),
(re.compile(r"^tel(ephone)?$", re.I), 11, "YD/T 1313", "手机号 11 位"),
(re.compile(r"^(xzqhbm|adcode|district_?code)$", re.I), 6, "GB/T 2260", "行政区划代码 6 位"),
(re.compile(r"^province_?code$", re.I), 2, "GB/T 2260", "省级代码 2 位"),
(re.compile(r"^city_?code$", re.I), 4, "GB/T 2260", "市级代码 4 位"),
(re.compile(r"^region_?code$", re.I), 6, "GB/T 2260", "区县级代码 6 位"),
(re.compile(r"^zip_?code$", re.I), 6, "GB/T 23703", "邮政编码 6 位"),
(re.compile(r"^email$", re.I), 50, "RFC 5321", "电子邮件 ≤254 位,常用 ≤50"),
(re.compile(r"^bank_?card(_?no)?$", re.I), 19, "JR/T 0002", "银行卡号 ≤19 位"),
(re.compile(r"^post_?code$", re.I), 6, "GB/T 23703", "邮政编码 6 位"),
]
def run_step6(dict_data: dict, llm: LLMClient | None = None,
log: Callable | None = None) -> dict:
columns = dict_data.get("data_dictionary", [])
if not columns:
if log:
log("WARN", "未获取到任何字段元数据,跳过 Step 6", step="6")
return {"summary": {"total_issues": 0}, "issues": []}
if log:
log("INFO",
f"内置长度规则: {len(LENGTH_RULES)} 条, 待扫描字段: {len(columns)}",
step="6")
issues = []
rule_match_count = 0
for r in columns:
col_name = r["column_name"]
dt = (r.get("data_type") or "").lower()
max_len = r.get("char_max_length")
if dt not in ("varchar", "char") or not max_len:
continue
for pat, expected, std, desc in LENGTH_RULES:
if pat.match(col_name):
rule_match_count += 1
if max_len > expected:
issues.append({
"table_name": r["table_name"],
"table_comment": r.get("table_comment", ""),
"column_name": col_name,
"data_type": dt,
"column_type": r.get("column_type", ""),
"actual_length": max_len,
"required_length": expected,
"wasted_bytes_per_row": max_len - expected,
"standard": std,
"rule": desc,
"issue": f"定义 {max_len} 位超出标准 {expected} 位",
"suggestion": f"改为 VARCHAR({expected})",
"column_comment": r.get("column_comment", ""),
})
break
if log:
log("INFO",
f"[1/2] 规则匹配 {rule_match_count} 次,发现异常 {len(issues)} 条",
step="6")
# 自定义字段(如命名不符合常见模式但 varchar 过大) — 简单启发式
# 仅对没有匹配规则且长度 ≥100 的字段提醒
custom_count = 0
for r in columns:
col_name = r["column_name"]
dt = (r.get("data_type") or "").lower()
max_len = r.get("char_max_length")
if dt not in ("varchar",) or not max_len or max_len < 100:
continue
if any(pat.match(col_name) for pat, *_ in LENGTH_RULES):
continue
if any(i["table_name"] == r["table_name"] and i["column_name"] == col_name for i in issues):
continue
# 大字段且无规则命中 → 候选给 LLM
issues.append({
"table_name": r["table_name"],
"table_comment": r.get("table_comment", ""),
"column_name": col_name,
"data_type": dt,
"column_type": r.get("column_type", ""),
"actual_length": max_len,
"required_length": None,
"wasted_bytes_per_row": None,
"standard": "(自定义)",
"rule": "未匹配任何标准规则",
"issue": f"字段长度 {max_len} 较大,未匹配任何标准",
"suggestion": "需业务确认或 LLM 推断",
"column_comment": r.get("column_comment", ""),
"needs_llm": True,
})
custom_count += 1
if log:
log("INFO",
f"[2/2] 自定义大字段 {custom_count} 条(候选给 LLM)",
step="6")
# LLM 处理 needs_llm 标记的字段(前 20 个)
if llm and llm.available:
to_llm = [i for i in issues if i.get("needs_llm")][:20]
if to_llm and log:
log("INFO", f"[LLM] 调用 LLM 推断 {len(to_llm)} 个自定义字段的合理长度", step="6")
llm_ok = 0
for idx, issue in enumerate(to_llm, 1):
try:
r = llm.suggest_length_rule(
column_name=issue["column_name"],
column_comment=issue["column_comment"],
current_length=issue["actual_length"],
table_name=issue["table_name"],
)
if r:
rec = r.get("recommended")
if rec and rec < issue["actual_length"]:
issue["required_length"] = rec
issue["suggestion"] = f"建议改为 VARCHAR({rec})({r.get('reasoning', '')})"
issue["llm_recommended"] = rec
issue["llm_reasoning"] = r.get("reasoning", "")
issue["wasted_bytes_per_row"] = issue["actual_length"] - rec
llm_ok += 1
if log:
log("DEBUG",
f" · [{idx}/{len(to_llm)}] "
f"{issue['table_name']}.{issue['column_name']}: "
f"{issue['actual_length']} → {rec}",
step="6")
except Exception as e:
if log:
log("WARN",
f"LLM 长度推断失败 ({issue['table_name']}.{issue['column_name']}): {e}",
step="6")
if to_llm and log:
log("INFO", f"[LLM] 推断完成: 采纳 {llm_ok}/{len(to_llm)}", step="6")
if log:
log("INFO",
f"字段长度异常 {len(issues)} 条(其中 {custom_count} 条自定义大字段)",
step="6")
return {
"summary": {
"total_issues": len(issues),
"rules_applied": len(LENGTH_RULES),
"wasted_bytes_total": sum(i.get("wasted_bytes_per_row") or 0 for i in issues),
},
"issues": issues[:200],
}
\ No newline at end of file
"""Step 7: 国家标准字段校验(连库)
按字段名匹配 standards 插件,抽样校验实际数据:
- 身份证号(GB 11643)
- 统一社会信用代码(GB 32100)
- 手机号(YD/T 1313)
- 行政区划代码(GB/T 2260)
抽样 SQL 走 web/sql/standards/sample_field_values.sql 模板。
"""
from __future__ import annotations
import logging
from typing import Callable
from ..db_adapter import DBConfig, open_db, quote_ident
from standards.registry import discover_standards, list_all_standards
from standards.base import BaseStandard, ValidationResult
from web.sql.loader import get_sql_loader
logger = logging.getLogger(__name__)
SAMPLE_LIMIT = 500 # 每字段最多抽样
MAX_TABLES_PER_FIELD = 5 # 每字段最多检查前 N 张表
def run_step7(cfg: DBConfig, dict_data: dict, log: Callable | None = None) -> dict:
"""执行 Step 7:国标校验"""
columns = dict_data.get("data_dictionary", [])
if not columns:
return {"summary": {}, "violations_by_standard": [], "violations": []}
# 字段名 → 标准实例的映射
standards_by_field = discover_standards()
all_standards_info = list_all_standards()
# 按字段名索引列
by_name: dict[str, list[dict]] = {}
for r in columns:
by_name.setdefault(r["column_name"], []).append(r)
by_standard: dict[str, dict] = {}
total_violations = 0
total_fields_checked = 0
violations_flat: list[dict] = []
if not standards_by_field:
if log:
log("WARN", "未发现任何标准插件", step="7")
return {
"summary": {"standards_applied": 0, "fields_checked": 0, "violations": 0},
"violations_by_standard": [],
"violations": [],
}
if log:
log("INFO",
f"已加载 {len(all_standards_info)} 个标准插件, "
f"匹配字段 {len(standards_by_field)} 个, "
f"最大每字段抽样 {SAMPLE_LIMIT} 行 / 每字段最多 {MAX_TABLES_PER_FIELD} 张表",
step="7")
loader = get_sql_loader()
try:
with open_db(cfg) as db:
for std_idx, (std_field, std_instance) in enumerate(standards_by_field.items(), 1):
if std_field not in by_name:
if log:
log("DEBUG",
f" · [{std_idx}/{len(standards_by_field)}] 标准字段 {std_field} 未出现在任何表中, 跳过",
step="7")
continue
tables_with_field = by_name[std_field][:MAX_TABLES_PER_FIELD]
std_id = std_instance.standard_id
std_record = by_standard.setdefault(std_id, {
"standard": std_id,
"standard_name": std_instance.standard_name,
"fields": [],
"total_invalid": 0,
"total_sampled": 0,
})
if log:
log("INFO",
f" · [{std_idx}/{len(standards_by_field)}] "
f"{std_id} ({std_field}) → {len(tables_with_field)} 张表",
step="7")
for col_record in tables_with_field:
table = col_record["table_name"]
field_name = col_record["column_name"]
total_fields_checked += 1
try:
sql = loader.render(
"standards/sample_field_values",
dialect=cfg.db_type,
table=table, # 模板内 ${table | quote} 自动加引号
field=field_name,
limit=SAMPLE_LIMIT,
)
rows = db.fetchall(sql)
except Exception as e:
std_record["fields"].append({
"table": table,
"field": field_name,
"error": str(e),
"violations": 0,
"sampled": 0,
})
if log:
log("WARN",
f" - {table}.{field_name} 抽样失败: {e}",
step="7")
continue
sampled = len(rows)
invalid_samples = []
invalid_count = 0
for row in rows:
val = row.get("val")
if val is None:
continue
result: ValidationResult = std_instance.validate(str(val))
if not result.valid:
invalid_count += 1
if len(invalid_samples) < 10:
invalid_samples.append({
"value": str(val)[:50],
"reason": result.reason,
})
violations_flat.append({
"table_name": table,
"column_name": field_name,
"rule_type": std_id,
"standard_name": std_instance.standard_name,
"column_type": col_record.get("column_type", ""),
"value": str(val)[:50],
"error": result.reason,
})
total_violations += invalid_count
std_record["total_invalid"] += invalid_count
std_record["total_sampled"] += sampled
std_record["fields"].append({
"table": table,
"field": field_name,
"sampled": sampled,
"violations": invalid_count,
"violation_rate": round(invalid_count / sampled, 4) if sampled else 0,
"samples": invalid_samples,
})
logger.debug(
f"{table}.{field_name}: 抽样 {sampled}, 违规 {invalid_count}"
)
except Exception as e:
if log:
log("ERROR", f"Step 7 连库校验失败: {type(e).__name__}: {e}", step="7")
logger.exception("Step 7 详细异常")
by_standard_list = []
for std_id, rec in by_standard.items():
by_standard_list.append({
"standard": rec["standard"],
"standard_name": rec["standard_name"],
"fields_count": len(rec["fields"]),
"total_invalid": rec["total_invalid"],
"total_sampled": rec["total_sampled"],
"violation_rate": round(rec["total_invalid"] / rec["total_sampled"], 4) if rec["total_sampled"] else 0,
"details": rec["fields"],
})
if log:
log("INFO", f"校验完成: {len(by_standard_list)} 个标准, {total_fields_checked} 个字段, {total_violations} 条违规", step="7")
return {
"summary": {
"standards_applied": len(all_standards_info),
"standards_used": len(by_standard_list),
"fields_checked": total_fields_checked,
"violations": total_violations,
},
"violations_by_standard": by_standard_list,
"violations": violations_flat[:500], # 给前端表格用
"standards_info": all_standards_info,
}
\ No newline at end of file
This diff is collapsed.
This diff is collapsed.
# ============================================================
# 数据治理 Web 工具 — Python 依赖
# ============================================================
# Web 框架
fastapi>=0.110.0
uvicorn[standard]>=0.27.0
pydantic>=2.5.0
# 数据库
pymysql>=1.1.0 # MySQL 驱动
dmPython>=2.5.0 # 达梦数据库驱动(需本地有达梦客户端)
# LLM
anthropic>=0.39.0 # Claude API 官方 SDK
# 报告生成
python-docx>=1.1.0
# 配置
pyyaml>=6.0.1
# 异步流式输出(可选,FastAPI 自带 sse-starlette 替代)
sse-starlette>=2.1.0
\ No newline at end of file
# SQL 模板目录
所有 SQL 语句集中在此目录,Python 代码不再直接拼 SQL。
## 目录结构
```
sql/
├── loader.py ← 加载器(get_sql_loader() 单例)
├── info_schema/
│ ├── list_columns.mysql.sql ← INFORMATION_SCHEMA.COLUMNS (MySQL)
│ ├── list_columns.dameng.sql ← ALL_TAB_COLUMNS (达梦)
│ ├── list_tables.mysql.sql
│ └── list_tables.dameng.sql
├── verify/
│ └── count_table_rows.sql ← SELECT COUNT(*) FROM tbl
├── xzqh/
│ ├── check_dict_table_exists.sql
│ └── find_xzqh_orphans.sql
├── empty_fields/
│ └── count_with_nulls.sql ← 动态列空值统计
└── standards/
└── sample_field_values.sql ← 字段抽样(国标校验用)
```
## 占位符语法
| 语法 | 含义 | 示例 |
|------|------|------|
| `${var}` | 简单字符串替换 | `${schema}` |
| `${var \| quote}` | 用 `quote_ident()` 加引号(按当前 dialect 自动选反引号或双引号) | `${table \| quote}` |
| `${list \| join:","}` | 列表按分隔符拼接 | `${cols \| join:","}` |
| `${var \| upper}` | 转大写 | `${schema \| upper}` |
| `${var \| lower}` | 转小写 | `${schema \| lower}` |
## Dialect 加载规则
加载器按以下优先级找模板:
1. `<name>.<dialect>.sql`(如 `list_columns.mysql.sql`)
2. `<name>.sql`(跨方言通用版)
## 在 Python 中使用
```python
from web.sql.loader import get_sql_loader
loader = get_sql_loader()
sql = loader.render("info_schema/list_columns", dialect="mysql", schema="mydb")
# SELECT ... FROM INFORMATION_SCHEMA.COLUMNS c WHERE c.TABLE_SCHEMA = `mydb`
```
## 添加新 SQL 模板
1. 在合适子目录下创建 `<name>.sql`(或 `<name>.<dialect>.sql`)
2. 用 `${...}` 占位符标记变量
3. Python 代码调 `loader.render("子目录/name", dialect=..., **params)`
## 注意事项
- 所有标识符(表名/字段名)必须用 `${var | quote}`,由加载器自动按方言加引号
- 字面量(数字、字符串)可直接替换 `${var}`
- 列表类占位符用 `${list | join:"sep"}`
- 模板中可写 SQL 注释(`--`),方便维护
\ No newline at end of file
-- ============================================================================
-- 统计单表多个列的 NULL/空值数量 + 总行数
-- 调用方:web/core/step_impl/step4_empty_fields.py
-- 参数:
-- ${table} 表名(自动加引号)
-- ${empty_clauses} 已构造好的 SUM(CASE WHEN...) 子句列表(逗号分隔)
-- 说明:
-- ${empty_clauses} 由 Python 端根据 data_type 过滤后动态构造,例如:
-- SUM(CASE WHEN "col1" IS NULL OR TRIM("col1") = '' THEN 1 ELSE 0 END) AS empty_c0,
-- SUM(CASE WHEN "col2" IS NULL OR TRIM("col2") = '' THEN 1 ELSE 0 END) AS empty_c1
-- ============================================================================
SELECT COUNT(*) AS total, ${empty_clauses}
FROM ${table | quote}
\ No newline at end of file
-- ============================================================================
-- 健康检查:确认连接可用,返回固定值 1
-- 调用方:web/core/db_adapter.py → test_connection()
-- 跨方言通用
-- ============================================================================
SELECT 1
\ No newline at end of file
-- ============================================================================
-- 列出指定 schema 下所有列的元数据 (达梦方言)
-- 调用方:web/core/db_adapter.py → list_columns()
-- 参数:schema 数据库/模式名(绑定参数,由调用方通过 params 传入,勿拼接)
-- 说明:达梦的 ALL_TAB_COLUMNS.COMMENTS 列名与 MySQL 的 COLUMN_COMMENT 不同
-- 达梦的 OWNER 通常为大写
-- ============================================================================
SELECT
c.TABLE_NAME AS table_name,
c.COLUMN_NAME AS column_name,
c.COLUMN_ID AS ordinal_position,
c.DATA_TYPE || CASE WHEN c.DATA_LENGTH IS NOT NULL THEN '(' || c.DATA_LENGTH || ')' END AS column_type,
c.DATA_TYPE AS data_type,
c.DATA_LENGTH AS char_max_length,
c.DATA_PRECISION AS numeric_precision,
c.DATA_SCALE AS numeric_scale,
c.NULLABLE AS is_nullable,
c.DATA_DEFAULT AS column_default,
c.COMMENTS AS column_comment,
NULL AS extra
FROM ALL_TAB_COLUMNS c
WHERE c.OWNER = ?
ORDER BY c.TABLE_NAME, c.COLUMN_ID
\ No newline at end of file
This diff is collapsed.
-- ============================================================================
-- 列出指定 schema 下所有表 (达梦方言)
-- 调用方:web/core/db_adapter.py → list_tables()
-- 参数:schema 数据库/模式名(绑定参数,由调用方通过 params 传入,勿拼接)
-- ============================================================================
SELECT TABLE_NAME,
'BASE TABLE' AS TABLE_TYPE,
NULL AS ENGINE,
NUM_ROWS AS TABLE_ROWS,
BYTES AS DATA_LENGTH,
0 AS INDEX_LENGTH,
COMMENTS AS TABLE_COMMENT,
CREATED AS CREATE_TIME,
LAST_DDL_TIME AS UPDATE_TIME
FROM ALL_TABLES
WHERE OWNER = UPPER(?)
ORDER BY TABLE_NAME
\ No newline at end of file
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
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