perf(mining): LLM簇命名串行→并发(Semaphore限5)——实测42簇串行命名540s占批次全程551s,并发后应~1/5;params种子段(τ/topk/min_cluster/big_split/keep_batches/embed_batch/cache_col,build.sh 2b幂等导入,已存在不覆盖);README同步挖掘架构文档

This commit is contained in:
yumoqing 2026-09-11 23:23:06 +08:00
parent f1eab971cb
commit bc6790a03a
4 changed files with 117 additions and 34 deletions

View File

@ -42,6 +42,10 @@ agent 禁止自批自审:确认与审批结论只能由人工任务回流驱
| opp_submit_report | 提交人工确认(门禁) |
| opp_list_approvals | 审批单列表 |
| opp_diagnose | 产线诊断(报告分布/待办/爬虫连通性) |
| opp_start_mining | 启动需求挖掘批次(后台:拉取→向量化→聚类→命名;scope=top 全量TopX / targeted 指定类型) |
| opp_mining_status | 批次状态轮询(pulling/embedding/clustering/naming/done/failed+原因) |
| opp_list_clusters | 类别排名表(热度/需求数/占比/命名证据)= TopX 热门需求类别 |
| opp_cluster_detail | 类内需求明细样例(每条带来源 URL 可查证) |
角色:`agent.opp_analyst`(商机分析师)、`agent.opp_writer`(研发报告撰写工程师)。
slash 命令:`/hot` `/ai` `/reports` `/oppdiag`。
@ -114,6 +118,31 @@ CRUD 页(`opp_reports/` `opp_approvals/`)保留供后台运维直接编辑
- `opp_reports`:研发报告(software/title/content/status/confirm_task_id/confirmed_by)
- `opp_approvals`:研发审批(report_id/status/note/resolved_by)
- `opp_mining_batches`:需求挖掘批次(org_id 隔离;scope=top/targeted;状态机
pulling→embedding→clustering→naming→done/failed;stats_json 进度统计;error_msg 可行动报错)
- `opp_demand_snap`:需求快照(batch_id+src_id 去重;cluster_id 聚类回填;embed_status)
- `opp_clusters`:需求类别(heat_rank 热度排名/share 占比/naming_evidence 命名证据/centroid_snap_id)
## 需求挖掘架构(P1,2026-09-11)
众包需求(爬虫平台 record_type=demand)→ 语义聚类 → TopX/指定类型分析:
- **embedding 走模型调用**:rag `/rag/api/embed.dspy`(凭据单点 rag_engine_configs)
← pipeline_service.rag_client.rag_embed_texts(project_id)(Bearer,10条/批)。
- **向量与近邻压给 VDB(Milvus)**(upapp.rag-vdb):kNN=/v1/query+vector+pagerows;
标量过滤=expr;协议细节见技能 rag-module-operations。
- **两层 collection**:共享缓存 `opp_demand_emb_cache`(基础数据全机构共用,
一条需求只嵌一次)+ 批次工作集 `opp_mine_<batch16>`(分析结果按批次隔离,
每机构保留最近 opp_mine_keep_batches 个,旧批自动 drop)。
- **聚类算法**:kNN(top30) + τ阈值(默认0.75) + 纯 Python 并查集(零原生依赖,
确定性可复现);>150 条大簇质心二次细分(τ+0.10,防众包模板标题过度聚合——
2075条实测 τ=0.75 拐点:42簇 top1=8.2%);小组(<5)质心并入最近大簇或归其他。
- **LLM 只做簇命名**(llm_call purpose=utility 60s 短超时),失败规则兜底(高频词)。
- **机构隔离**(用户定夺):基础数据共享、分析结果隔离——批次/快照/类别挂 org_id,
查询强制过滤(E2E 伪 org 拒绝实测通过)。
- 算法参数全走 appbase params 表:opp_mine_tau/topk/min_cluster/big_split/keep_batches、
opp_embed_batch、opp_emb_cache_col(代码默认兜底)。
- 验证脚本:scripts/p0_probe_vdb_embed*.py(VDB/embed 协议探测)、
p1_cluster_analyze.py(多τ对比+缓存,调参用)、p1_create_tables.py(幂等建表)。
## 技能

View File

@ -98,6 +98,40 @@ async def seed():
asyncio.run(seed())
PYEOF
# ── 2b. params 种子:挖掘算法参数(幂等:已存在不覆盖,运行期可在 params 表调)──
cd "$APP_ROOT"
"$PYTHON" - "$SCRIPT_DIR" "$APP_ROOT" <<'PYEOF'
import sys, os, json, asyncio
CDIR = sys.argv[1]
sys.path.insert(0, os.getcwd())
from appPublic.jsonConfig import getConfig
from sqlor.dbpools import DBPools
from appPublic.uniqueID import getID
data = json.load(open(os.path.join(CDIR, 'init', 'data.json')))
params_seed = data.get('params', [])
async def seed_params():
if not params_seed:
return
config = getConfig(sys.argv[2], {'workdir': sys.argv[2]})
db = DBPools(config.databases)
n = 0
async with db.sqlorContext('pipeline') as sor:
for pitem in params_seed:
chk = await sor.sqlExe(
"SELECT id FROM params WHERE params_name=${n}$", {"n": pitem['name']})
if not chk:
await sor.C('params', {'id': getID(),
'params_name': pitem['name'],
'params_value': str(pitem['value'])})
n += 1
await sor.sqlExe("COMMIT", {})
print('params seeded:', n, '(existing kept)')
asyncio.run(seed_params())
PYEOF
# ── 3. wwwroot 软链 + pip install ──
rm -f "$APP_ROOT/wwwroot/pipeline-opportunity"
ln -sf "$SCRIPT_DIR/wwwroot" "$APP_ROOT/wwwroot/pipeline-opportunity"

View File

@ -30,5 +30,14 @@
{"k": "rejected", "v": "审批驳回"}
]
}
],
"params": [
{"name": "opp_mine_tau", "value": "0.75", "desc": "聚类相似度阈值(2075条实测拐点:0.75出42簇top1=8.2%)"},
{"name": "opp_mine_topk", "value": "30", "desc": "kNN邻居数"},
{"name": "opp_mine_min_cluster", "value": "5", "desc": "最小簇规模,低于则并入最近大簇或其他"},
{"name": "opp_mine_big_split", "value": "150", "desc": "超大簇质心二次细分阈值(防模板标题过度聚合)"},
{"name": "opp_mine_keep_batches", "value": "3", "desc": "每机构保留最近N个批次工作集collection,旧批自动drop"},
{"name": "opp_embed_batch", "value": "10", "desc": "rag /embed 单批文本上限"},
{"name": "opp_emb_cache_col", "value": "opp_demand_emb_cache", "desc": "共享embedding缓存collection(全机构复用)"}
]
}

View File

@ -389,11 +389,15 @@ def cluster_vectors(vecs, nbr, tau, min_cluster, big_split):
# ══════════════════ 算法层:LLM 簇命名(规则兜底)══════════════════
async def name_clusters(sor, batch, clusters, vecs, text_map, cfg, ctx):
"""为每个簇命名:LLM(utility) 优先,失败/异常用高频词规则兜底。
返回 list[{name, members, centroid_snap, samples, evidence}]。"""
from .opp_common import get_db as _g
results = []
"""为每个簇命名:LLM(utility) 并发调用 + 失败/异常规则兜底。
返回 list[{name, members, centroid_snap, samples, evidence}]。
性能实测教训(2026-09-11):42 簇串行 LLM 命名 540s 占批次全程 551s——
改并发(Semaphore 限 5,防上游限流)后命名阶段应降至 ~1/5。
"""
total = len(vecs)
# 阶段1(纯CPU,串行):算质心/样例/排名
prepared = []
for ci, members in enumerate(clusters):
cent = _centroid(vecs, members)
sims = sorted(((_cos(vecs[m], cent), m) for m in members), reverse=True)
@ -406,37 +410,44 @@ async def name_clusters(sor, batch, clusters, vecs, text_map, cfg, ctx):
samples.append(t)
if len(samples) >= 8:
break
name, evidence = "", ""
# LLM 命名(utility purpose,60s 短超时快速失败走兜底——记忆铁律)
try:
from pipeline_service.llm_bridge import llm_call
prompt = (
"以下是同一类软件众包需求的标题样例(已去重):\n" +
"\n".join("- " + s for s in samples[:8]) +
"\n\n请用不超过12个汉字给这一类需求起一个简洁准确的类别名(如「微信小程序开发」"
"「企业网站定制」「AI短视频制作」),只输出类别名本身,不要解释、不要标点。")
name = (await llm_call(prompt, purpose="utility", timeout=60,
org_id=ctx.get("org_id") or "0",
user_id=ctx.get("user_id") or "",
session_id=ctx.get("session_id") or "")).strip()
name = name.splitlines()[0][:40] if name else ""
evidence = "llm"
# LLM 可能输出带引号/序号,清洗
name = name.strip("「」\"'。.  ")
except Exception as e:
logger.debug("LLM 命名失败,走规则兜底: %s", e)
name = ""
if not name:
name = _rule_name(samples)
evidence = "rule"
# 质心代表样本 → snap_id
centroid_vid = sims[0][1] if sims else ""
results.append({
"rank": ci + 1, "name": name, "members": members,
"size": len(members), "share": round(len(members) / total, 4) if total else 0,
"centroid_vid": centroid_vid, "samples": samples[:8], "evidence": evidence,
prepared.append({
"rank": ci + 1, "members": members, "size": len(members),
"share": round(len(members) / total, 4) if total else 0,
"centroid_vid": sims[0][1] if sims else "",
"samples": samples[:8], "name": "", "evidence": "rule",
})
return results
# 阶段2:LLM 命名(并发限5;单簇失败不影响整体,走规则兜底)
from pipeline_service.llm_bridge import llm_call
sem = asyncio.Semaphore(5)
async def _name_one(item):
async with sem:
try:
prompt = (
"以下是同一类软件众包需求的标题样例(已去重):\n" +
"\n".join("- " + s for s in item["samples"][:8]) +
"\n\n请用不超过12个汉字给这一类需求起一个简洁准确的类别名(如「微信小程序开发」"
"「企业网站定制」「AI短视频制作」),只输出类别名本身,不要解释、不要标点。")
raw = await llm_call(prompt, purpose="utility", timeout=60,
org_id=ctx.get("org_id") or "0",
user_id=ctx.get("user_id") or "",
session_id=ctx.get("session_id") or "")
name = (raw or "").strip().splitlines()[0][:40] if raw else ""
# LLM 可能输出带引号/序号,清洗
name = name.strip("「」\"'。.  ")
if name:
item["name"] = name
item["evidence"] = "llm"
except Exception as e:
logger.debug("LLM 命名失败,走规则兜底: %s", e)
await asyncio.gather(*(_name_one(it) for it in prepared))
# 阶段3:规则兜底
for it in prepared:
if not it["name"]:
it["name"] = _rule_name(it["samples"])
return prepared
def _rule_name(samples):