Commit d442e8fb authored by Data Governance Dev's avatar Data Governance Dev

chore: 更新文档、配置默认值、任务管理及数据模型补充

parent 8ac75cef
......@@ -28,6 +28,20 @@
- 期间记录好操作记录,后须最好能抽出成一个固定工作流程
# LLM 依赖(必须 vs 可选)
LLM 不是「增强」,对部分步骤是核心依赖:
| 步骤 | LLM 模式 | 含义 |
|------|---------|------|
| Step 2 表合并与冗余字段分析 | **required** | 合并 verdict 与高频字段分类必须由 LLM 判断;缺 LLM = 任务失败 |
| Step 5 缺失注释字段推测 | **required** | 字段语义(拼音首字母缩写、英文组合)必须由 LLM 翻译;缺 LLM = 任务失败 |
| Step 4 空字段成因分类 | optional | 有 LLM 更准;缺了也能跑(走规则推理) |
| Step 6 自定义字段长度建议 | optional | 同上 |
任务提交前需要确保 `web/configs/llm.yaml` 中 `api_key` 已配置,
否则勾选 Step 2 / 5 的任务会被预检拦截、整体失败。
# 下一阶段
- 目前已经把过去工作整理成工作流并且实现Web端调用,参考: `./web/README.md`
......
......@@ -186,17 +186,25 @@ python -m uvicorn web.app:app --host 0.0.0.0 --port 8765
---
## 🤖 LLM 增强位置
## 🤖 LLM 依赖矩阵
| Step | LLM 用途 | 失败降级 |
|------|----------|----------|
| Step 2 | 表合并建议文案 + 风险评估 | 模板字符串 |
| Step 4 | 高空字段成因分类 | 跳过 LLM 分类 |
| Step 5 | 无规则命中的字段注释推测 | 仅靠规则映射表 |
| Step 6 | 自定义字段长度推断 | 标记"待业务确认" |
LLM 不是「增强」,对部分步骤是核心依赖:
LLM 调用封装在 [web/core/llm.py](web/core/llm.py),所有方法都设计为**可降级**:
未配置 `ANTHROPIC_API_KEY` 时自动跳过,不阻断主流程。
| Step | LLM 模式 | 用途 | 缺 LLM 行为 |
|------|---------|------|------------|
| Step 2 表合并与冗余字段分析 | **required** | 合并 verdict + 高频字段分类 | 任务整体失败 |
| Step 5 缺失注释字段推测 | **required** | 字段语义推测(拼音/英文→中文) | 任务整体失败 |
| Step 4 空字段成因分类 | optional | top 20 高空字段成因分类 | 走规则推理(标 "LLM 未配置") |
| Step 6 自定义字段长度建议 | optional | 自定义字段合理长度推断 | 标记"待业务确认" |
### 设计要点
- **required** 步骤:任务启动时预检 LLM 可用性;不可用 → `RuntimeError` → 任务标 `JOB_FAILED`
- **optional** 步骤:LLM 可用就用,失败自动降级到规则推理
- **批处理**:Step 2 / 5 的 LLM 调用按 15 字段/批,避免 200 次串行 API
LLM 调用封装在 [web/core/llm.py](web/core/llm.py)。配置 Key 见上方「配置 LLM Key」
小节;未配置时勾选 Step 2/5 的任务会被拦截、整体失败。
---
......
......@@ -34,4 +34,4 @@ llm:
web:
host: 0.0.0.0
port: 8765
title: "数据治理工具"
\ No newline at end of file
title: "数据分析工具"
\ No newline at end of file
......@@ -43,16 +43,24 @@ JOB_CANCELLED = "cancelled"
@dataclass
class LogBuffer:
"""环形日志缓冲 + 异步发布(订阅者可 SSE 流式接收)"""
"""环形日志缓冲 + 异步发布(订阅者可 SSE 流式接收)
每条 entry 拥有单调递增的 _seq 字段,用作 SSE 事件 id:
- 前端 EventSource 自动在重连时通过 Last-Event-ID header 发送上次收到的 id
- 后端据此跳过已发条目,避免日志被重复 append
"""
max_size: int = 5000
def __post_init__(self):
self._entries: list[dict] = []
self._subscribers: list[asyncio.Queue] = []
self._lock = asyncio.Lock()
self._next_seq: int = 0 # 单调递增,每个 entry 唯一
async def append(self, entry: dict) -> None:
async with self._lock:
self._next_seq += 1
entry["_seq"] = self._next_seq
self._entries.append(entry)
if len(self._entries) > self.max_size:
self._entries = self._entries[-self.max_size:]
......@@ -72,6 +80,10 @@ class LogBuffer:
def snapshot(self) -> list[dict]:
return list(self._entries)
def snapshot_after(self, last_seq: int) -> list[dict]:
"""返回 _seq > last_seq 的快照。"""
return [e for e in self._entries if e.get("_seq", 0) > last_seq]
def subscribe(self) -> asyncio.Queue:
q: asyncio.Queue = asyncio.Queue(maxsize=1000)
self._subscribers.append(q)
......@@ -302,12 +314,20 @@ class JobManager:
return True
# ── 日志订阅(SSE 用) ──
async def stream_logs(self, job_id: str) -> AsyncIterator[dict]:
async def stream_logs(self, job_id: str, last_seq: int = 0) -> AsyncIterator[dict]:
"""SSE 日志流。
Args:
job_id: 任务 ID
last_seq: 客户端已收到的最后一条 entry 的 _seq(来自 SSE Last-Event-ID
header,浏览器 EventSource 在自动重连时会自动发送)。
0 表示首次连接,需要全量 snapshot;>0 表示续传,只发更新的部分。
"""
job = self.get(job_id)
if not job:
return
# 先发历史
for entry in job.logs.snapshot():
# 历史回放:跳过 last_seq 之前的(避免重连时重复推送)
for entry in job.logs.snapshot_after(last_seq):
yield entry
if job.status in (JOB_COMPLETED, JOB_FAILED, JOB_CANCELLED):
return
......
......@@ -70,7 +70,10 @@ class StepInfo(BaseModel):
title: str
description: str
requires_db: bool
# 旧字段保留为向后兼容;新代码请用 llm_mode
llm_enhanced: bool = False
# "none" | "optional" | "required"
llm_mode: str = "none"
class StepsListResponse(BaseModel):
......
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