From bc6790a03a19daf1a6199ae470fbaaa4a6165408 Mon Sep 17 00:00:00 2001 From: yumoqing Date: Fri, 11 Sep 2026 23:23:06 +0800 Subject: [PATCH] =?UTF-8?q?perf(mining):=20LLM=E7=B0=87=E5=91=BD=E5=90=8D?= =?UTF-8?q?=E4=B8=B2=E8=A1=8C=E2=86=92=E5=B9=B6=E5=8F=91(Semaphore?= =?UTF-8?q?=E9=99=905)=E2=80=94=E2=80=94=E5=AE=9E=E6=B5=8B42=E7=B0=87?= =?UTF-8?q?=E4=B8=B2=E8=A1=8C=E5=91=BD=E5=90=8D540s=E5=8D=A0=E6=89=B9?= =?UTF-8?q?=E6=AC=A1=E5=85=A8=E7=A8=8B551s,=E5=B9=B6=E5=8F=91=E5=90=8E?= =?UTF-8?q?=E5=BA=94~1/5;params=E7=A7=8D=E5=AD=90=E6=AE=B5(=CF=84/topk/min?= =?UTF-8?q?=5Fcluster/big=5Fsplit/keep=5Fbatches/embed=5Fbatch/cache=5Fcol?= =?UTF-8?q?,build.sh=202b=E5=B9=82=E7=AD=89=E5=AF=BC=E5=85=A5,=E5=B7=B2?= =?UTF-8?q?=E5=AD=98=E5=9C=A8=E4=B8=8D=E8=A6=86=E7=9B=96);README=E5=90=8C?= =?UTF-8?q?=E6=AD=A5=E6=8C=96=E6=8E=98=E6=9E=B6=E6=9E=84=E6=96=87=E6=A1=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- README.md | 29 +++++++++++ build.sh | 34 +++++++++++++ init/data.json | 9 ++++ pipeline_opportunity/opp_mining.py | 79 +++++++++++++++++------------- 4 files changed, 117 insertions(+), 34 deletions(-) diff --git a/README.md b/README.md index 4448037..19b284c 100644 --- a/README.md +++ b/README.md @@ -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_`(分析结果按批次隔离, + 每机构保留最近 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(幂等建表)。 ## 技能 diff --git a/build.sh b/build.sh index 31bf63a..465d81c 100644 --- a/build.sh +++ b/build.sh @@ -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" diff --git a/init/data.json b/init/data.json index 1f61780..c3b173e 100644 --- a/init/data.json +++ b/init/data.json @@ -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(全机构复用)"} ] } diff --git a/pipeline_opportunity/opp_mining.py b/pipeline_opportunity/opp_mining.py index 716b8d5..59ecb26 100644 --- a/pipeline_opportunity/opp_mining.py +++ b/pipeline_opportunity/opp_mining.py @@ -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):