diff --git a/mysql.ddl.sql b/mysql.ddl.sql index d92c1f6..49be3d3 100644 --- a/mysql.ddl.sql +++ b/mysql.ddl.sql @@ -1,5 +1,5 @@ --- /home/ymq/work/pipeline/pipeline-opportunity/models/opp_reports.json +-- /home/ymq/work/repos/pipeline-opportunity/models/opp_reports.json @@ -18,6 +18,7 @@ CREATE TABLE opp_reports `content` longtext CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci comment '报告正文', `status` VARCHAR(24) NOT NULL DEFAULT 'draft' comment '报告状态', `confirm_task_id` VARCHAR(32) comment '确认任务ID', + `ppt_path` VARCHAR(500) comment 'PPT文件路径', `confirmed_by` VARCHAR(64) comment '确认人', `created_by` VARCHAR(64) comment '创建人', `created_at` TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL comment '创建时间', @@ -35,7 +36,7 @@ comment '研发报告' ; --- /home/ymq/work/pipeline/pipeline-opportunity/models/opp_approvals.json +-- /home/ymq/work/repos/pipeline-opportunity/models/opp_approvals.json @@ -68,3 +69,112 @@ engine=innodb comment '研发审批' ; + +-- /home/ymq/work/repos/pipeline-opportunity/models/opp_clusters.json + + + + + +-- 建库时请用以下语句,支持emoji字符 +-- CREATE DATABASE mydb CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; +drop table if exists opp_clusters; +CREATE TABLE opp_clusters +( + + `id` VARCHAR(32) NOT NULL comment '主键ID', + `batch_id` VARCHAR(32) NOT NULL comment '批次ID', + `org_id` VARCHAR(32) NOT NULL comment '机构ID', + `name` VARCHAR(128) comment '类别名称(LLM命名)', + `doc_count` int NOT NULL DEFAULT '0' comment '类内需求数', + `share` double(6,4) DEFAULT '0' comment '占比', + `heat_rank` int DEFAULT '0' comment '热度排名', + `naming_evidence` longtext CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci comment '命名证据(样例标题JSON)', + `centroid_snap_id` VARCHAR(32) comment '质心需求快照ID(代表样本)', + `created_at` TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL comment '创建时间' + + +,primary key(id) + + +) +CHARACTER SET utf8mb4 +COLLATE utf8mb4_unicode_ci +engine=innodb +comment '需求类别(聚类结果)' +; + + +-- /home/ymq/work/repos/pipeline-opportunity/models/opp_mining_batches.json + + + + + +-- 建库时请用以下语句,支持emoji字符 +-- CREATE DATABASE mydb CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; +drop table if exists opp_mining_batches; +CREATE TABLE opp_mining_batches +( + + `id` VARCHAR(32) NOT NULL comment '主键ID', + `org_id` VARCHAR(32) NOT NULL comment '机构ID', + `scope` VARCHAR(16) NOT NULL DEFAULT 'top' comment '挖掘范围(top=TopX全量/targeted=指定类型)', + `params_json` longtext CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci comment '批次参数(days/sources/record_type/keyword/category)', + `status` VARCHAR(24) NOT NULL DEFAULT 'pulling' comment '批次状态(pulling/embedding/clustering/naming/done/failed)', + `vdb_col` VARCHAR(64) comment '批次向量工作集collection名', + `stats_json` longtext CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci comment '批次统计(需求数/嵌入数/簇数/失败数)', + `error_msg` longtext CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci comment '失败原因(可行动报错)', + `created_by` VARCHAR(64) comment '创建人', + `created_at` TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL comment '创建时间', + `updated_at` TIMESTAMP DEFAULT CURRENT_TIMESTAMP comment '更新时间' + + +,primary key(id) + + +) +CHARACTER SET utf8mb4 +COLLATE utf8mb4_unicode_ci +engine=innodb +comment '需求挖掘批次' +; + + +-- /home/ymq/work/repos/pipeline-opportunity/models/opp_demand_snap.json + + + + + +-- 建库时请用以下语句,支持emoji字符 +-- CREATE DATABASE mydb CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; +drop table if exists opp_demand_snap; +CREATE TABLE opp_demand_snap +( + + `id` VARCHAR(32) NOT NULL comment '主键ID', + `batch_id` VARCHAR(32) NOT NULL comment '批次ID', + `org_id` VARCHAR(32) NOT NULL comment '机构ID', + `src_id` VARCHAR(64) NOT NULL comment '爬虫平台需求ID(去重键)', + `source` VARCHAR(32) comment '数据来源(zbj/epwk/freelancer/pph)', + `title` VARCHAR(500) NOT NULL comment '需求标题', + `item_category` VARCHAR(128) comment '站内类目', + `budget_wan` double(15,2) comment '预算(万元)', + `url` VARCHAR(500) comment '需求原文链接', + `publish_time` VARCHAR(32) comment '发布时间', + `cluster_id` VARCHAR(32) comment '所属类别ID(聚类后回填)', + `embed_status` VARCHAR(16) DEFAULT 'new' comment '向量化状态(cached/new/failed)', + `created_at` TIMESTAMP DEFAULT CURRENT_TIMESTAMP NOT NULL comment '创建时间' + + +,primary key(id) + + +) +CHARACTER SET utf8mb4 +COLLATE utf8mb4_unicode_ci +engine=innodb +comment '需求快照' +; + diff --git a/pipeline_opportunity/opp_ability.py b/pipeline_opportunity/opp_ability.py index 233ae35..6b0be69 100644 --- a/pipeline_opportunity/opp_ability.py +++ b/pipeline_opportunity/opp_ability.py @@ -80,6 +80,15 @@ OPP_TOOLS = [ parameters={"status": "状态过滤(可选)"}, category="report"), ToolDefinition(name="opp_diagnose", description="诊断商机产线:报告/审批分布、待办人工任务、爬取平台连通性", parameters={}, category="agent"), + # ── 需求挖掘(众包需求→语义聚类→TopX/指定类型,P1)── + ToolDefinition(name="opp_start_mining", description="启动需求挖掘批次(后台执行,立即返回batch_id):从众包需求池拉取→向量化→语义聚类→LLM命名类别。scope=top 全量聚类出TopX热门类别排名;scope=targeted 指定类型分析(必须带keyword如'合同管理',在范围内聚类找子型)。结果用 opp_mining_status 轮询,完成后用 opp_list_clusters 看类别排名", + parameters={"scope": "top=全量TopX排名 / targeted=指定类型分析", "keyword": "targeted模式的类型关键词(如:合同管理)", "category": "targeted模式的站内类目过滤(可选)", "days": "需求时间窗天数(默认365)", "sources": "限定来源zbj/epwk/freelancer/pph(可选,默认全部)"}, category="mining"), + ToolDefinition(name="opp_mining_status", description="查需求挖掘批次状态(pulling拉取/embedding向量化/clustering聚类/naming命名/done完成/failed失败+原因)。batch_id空=列本机构最近批次", + parameters={"batch_id": "批次ID(空=最近批次列表)"}, category="mining"), + ToolDefinition(name="opp_list_clusters", description="查看挖掘批次的类别排名表(热度排名/需求数/占比/命名证据),即TopX热门需求类别。批次必须done", + parameters={"batch_id": "批次ID"}, category="mining"), + ToolDefinition(name="opp_cluster_detail", description="查看某类别的类内需求明细样例(每条带来源URL可查证),用于写可行性报告时引用证据", + parameters={"cluster_id": "类别ID", "limit": "样例条数(默认20)"}, category="mining"), ] # ══════════════════ 产线 prompt 片段 ══════════════════ @@ -128,6 +137,18 @@ opp_create_report / opp_update_report 的正文必须用以下四个二级标题 - `## 成果`:可交付成果(建议的产品形态、里程碑、预期成果文件如 PPT/方案) 旧格式(无四节标题)的存量报告仍可查看,全文归入「输出」节展示。 +## 需求挖掘(众包需求 → TopX 类别 / 指定类型分析) +用户问「市场上什么软件需求最多」「帮我分析XX类需求」「做需求挖掘」时走这条链: +1. opp_start_mining 启动批次——scope=top(全量聚类出热门类别排名)或 + scope=targeted + keyword(用户指定类型,如「合同管理」「小程序」,在范围内聚类找子型)。 + 批次后台执行(拉取→向量化→聚类→命名),启动后**立即返回 batch_id**,不要傻等。 +2. opp_mining_status 轮询进度(每轮回复用户当前状态;failed 时如实转述原因,禁隐瞒)。 +3. done 后 opp_list_clusters 看类别排名(热度/需求数/占比),向用户呈现 TopX 表格。 +4. 用户选定类别后 opp_cluster_detail 取类内样例(每条带来源URL),结合数据写分析报告; + 需要立项走可行性研究时,用 opp_create_report 建报告(数据基础引用挖掘批次与类别证据)。 +硬规则:聚类结果只能来自 opp_* mining 工具返回,禁止自己编类别/编排名/编数量; +类别命名证据(naming_evidence 样例)要能支撑命名,命名可疑时如实说明。 + ## 硬规则 - **统计优先用聚合工具,明细只取样例(2026-09-11 实测教训)**: 数量/排名/预算合计/平均价这类**统计**一律用聚合工具——opp_hot_software(招标侧)、 @@ -416,6 +437,63 @@ async def _h_diagnose(sor, p, ctx): }, ensure_ascii=False, default=str) +# ══════════════════ 需求挖掘 handler(P1)══════════════════ + +async def _h_start_mining(sor, p, ctx): + from .opp_mining_flow import start_mining + scope = (p.get("scope") or "top").strip() + bid, err = await start_mining( + sor, ctx, scope=scope, + days=int(p.get("days") or 365), + sources=(p.get("sources") or "").strip(), + keyword=(p.get("keyword") or "").strip(), + category=(p.get("category") or "").strip()) + if err: + return _fmt(False, err) + return _fmt(True, "挖掘批次已启动 batch_id=%s(后台执行:%s)。用 opp_mining_status 轮询进度,done 后用 opp_list_clusters 看类别排名" % ( + bid, "全量TopX" if scope == "top" else "指定类型[%s/%s]" % (p.get("keyword") or "", p.get("category") or ""))) + + +async def _h_mining_status(sor, p, ctx): + from .opp_mining_flow import mining_status + ok, res = await mining_status(sor, ctx, batch_id=(p.get("batch_id") or "").strip()) + if not ok: + return _fmt(False, res) + return json.dumps(res, ensure_ascii=False, default=str) + + +async def _h_list_clusters(sor, p, ctx): + from .opp_mining_flow import list_clusters + bid = (p.get("batch_id") or "").strip() + if not bid: + return _fmt(False, "batch_id 必填(先用 opp_mining_status 查批次)") + ok, res = await list_clusters(sor, ctx, bid) + if not ok: + return _fmt(False, res) + rows = [{"排名": c["heat_rank"], "类别": c["name"], "需求数": c["doc_count"], + "占比": "%.1f%%" % (float(c["share"] or 0) * 100), "cluster_id": c["id"], + "命名证据": (json.loads(c["naming_evidence"] or "{}").get("samples") or [])[:3]} + for c in res] + return json.dumps({"batch_id": bid, "类别数": len(rows), "排名": rows}, + ensure_ascii=False, default=str) + + +async def _h_cluster_detail(sor, p, ctx): + from .opp_mining_flow import cluster_detail + cid = (p.get("cluster_id") or "").strip() + if not cid: + return _fmt(False, "cluster_id 必填") + ok, res = await cluster_detail(sor, ctx, cid, limit=int(p.get("limit") or 20)) + if not ok: + return _fmt(False, res) + cl = res["cluster"] + samples = [{"标题": s["title"], "来源": s["source"], "预算(万)": s["budget_wan"], + "url": s["url"], "发布": s["publish_time"]} for s in res["samples"]] + return json.dumps({"类别": cl["name"], "需求数": cl["doc_count"], + "占比": "%.1f%%" % (float(cl["share"] or 0) * 100), + "样例(带来源可查证)": samples}, ensure_ascii=False, default=str) + + OPP_HANDLERS = { "opp_hot_software": _h_hot_software, "opp_daily_ai_tenders": _h_daily_ai, @@ -435,6 +513,11 @@ OPP_HANDLERS = { "opp_submit_report": _h_submit_report, "opp_list_approvals": _h_list_approvals, "opp_diagnose": _h_diagnose, + # 需求挖掘(P1) + "opp_start_mining": _h_start_mining, + "opp_mining_status": _h_mining_status, + "opp_list_clusters": _h_list_clusters, + "opp_cluster_detail": _h_cluster_detail, } diff --git a/pipeline_opportunity/opp_mining.py b/pipeline_opportunity/opp_mining.py new file mode 100644 index 0000000..716b8d5 --- /dev/null +++ b/pipeline_opportunity/opp_mining.py @@ -0,0 +1,465 @@ +# -*- coding:utf-8 -*- +"""商机产线需求挖掘能力(P1):众包需求 → 语义聚类 → TopX/指定类型。 + +分层边界(pipeline-extension-patterns): + · embedding 走 pipeline_service.rag_client.rag_embed_texts(凭据单点在 rag)。 + · 向量近邻全部压给 VDB(Milvus)(upapp.rag-vdb),产线侧零向量计算。 + · 聚类 = kNN + τ阈值 + 纯 Python 并查集(无原生依赖,Nuitka 友好,确定性可复现)。 + · LLM 只做簇命名(purpose=utility),失败规则兜底不崩。 + +机构隔离(用户定夺 2026-09-11): + · 基础数据共享——爬虫需求 + embedding 缓存 collection(opp_demand_emb_cache)全机构共用, + 一条需求只嵌一次,跨机构/跨批次复用,成本摊薄。 + · 分析结果隔离——批次/快照/类别全部挂 org_id,所有查询按 ctx.org_id 强制过滤。 + +VDB 协议(实测,详见 rag-module-operations 技能): + createcollection/upsert/query(kNN=vector+pagerows)/expr标量过滤/dropcollection; + kNN+expr 组合可用;output_fields 可取回 vector(缓存→工作集拷贝)。 +""" +import asyncio +import hashlib +import json +import logging + +from .opp_common import get_db, new_id, rows_to_dicts, get_param, get_crawler_config + +logger = logging.getLogger("pipeline.opp_mining") + +DIM = 1024 +# 共享 embedding 缓存 collection(基础数据,全机构共用) +CACHE_COL_DEFAULT = "opp_demand_emb_cache" +# 参数默认兜底(params 表优先,禁硬编码为唯一来源) +P = { + "opp_mine_tau": "0.75", + "opp_mine_topk": "30", + "opp_mine_min_cluster": "5", + "opp_mine_big_split": "150", # 超过此规模的簇质心二次细分 + "opp_mine_keep_batches": "3", # 每机构保留最近 N 个批次工作集 collection + "opp_embed_batch": "10", # rag /embed 单批上限(实测 dashscope 兼容 10) +} + +# 状态机(非法迁移拒绝,对齐 REPORT_TRANSITIONS 范式) +BATCH_TRANSITIONS = { + "pulling": {"embedding", "failed", "done"}, + "embedding": {"clustering", "failed"}, + "clustering": {"naming", "failed"}, + "naming": {"done", "failed"}, + "done": set(), + "failed": {"pulling"}, # 失败可重跑 +} + + +def can_batch_transition(cur, nxt): + return nxt in BATCH_TRANSITIONS.get(cur, set()) + + +# ══════════════════ 参数 / VDB 客户端 ══════════════════ + +async def get_mining_params(sor): + cfg = {} + for k, dv in P.items(): + cfg[k] = await get_param(sor, k, dv) + cache_col = await get_param(sor, "opp_emb_cache_col", CACHE_COL_DEFAULT) + cfg["cache_col"] = cache_col + for ik in ("opp_mine_topk", "opp_mine_min_cluster", "opp_mine_big_split", + "opp_mine_keep_batches", "opp_embed_batch"): + try: + cfg[ik] = int(cfg[ik]) + except (TypeError, ValueError): + cfg[ik] = int(P[ik]) + try: + cfg["opp_mine_tau"] = float(cfg["opp_mine_tau"]) + except (TypeError, ValueError): + cfg["opp_mine_tau"] = float(P["opp_mine_tau"]) + return cfg + + +async def _vdb_base(sor): + recs = await sor.sqlExe("SELECT baseurl FROM upapp WHERE id='rag-vdb'", {}) + await sor.sqlExe("COMMIT", {}) + return (recs[0].baseurl or "").rstrip("/") if recs else "" + + +async def vdb_post(sor, base, path, payload, timeout=40): + """VDB REST 调用(aiohttp)。返回 (ok, dict|错误串)。失败显式报错,禁静默。""" + try: + import aiohttp + async with aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=timeout)) as s: + async with s.post(base + path, json=payload) as r: + txt = await r.text() + try: + d = json.loads(txt) if txt.strip().startswith("{") else {} + except Exception: + d = {} + if r.status != 200 or (d.get("status") and d.get("status") != "SUCCEEDED"): + return False, "VDB %s HTTP%d %s" % (path, r.status, (d.get("error") or txt)[:200]) + return True, d + except Exception as e: + return False, "VDB %s 调用失败: %s" % (path, str(e)[:160]) + + +def _vec_rows(d): + """解析 VDB 响应的 rows(兼容嵌套 data.rows)。""" + data = d.get("data") if isinstance(d, dict) else None + if isinstance(data, dict): + return data.get("rows") or [] + if isinstance(data, list): + return data + return d.get("rows") or [] if isinstance(d, dict) else [] + + +def demand_vec_id(source, src_id): + """需求→稳定向量ID(sha1,≤32字符,可复算幂等)。""" + return hashlib.sha1(("%s:%s" % (source or "", src_id or "")).encode("utf-8")).hexdigest()[:32] + + +# ══════════════════ 数据层:快照拉取 ══════════════════ + +async def pull_demands(sor, batch, ctx, cfg, scope, keyword="", category="", + days=365, sources="", limit=200): + """从爬虫平台拉需求 → 写 opp_demand_snap(按 batch+src_id 去重)。返回 (n, err)。""" + base, token = await get_crawler_config(sor) + if not token: + return 0, "params 缺 tender_api_token(爬虫平台接入未配置)" + org_id = ctx.get("org_id") or "" + bid = batch["id"] + offset = 0 + total = 0 + while True: + params = {"days": days, "limit": limit, "offset": offset, "record_type": "demand"} + if keyword: + params["keyword"] = keyword + if sources: + params["source"] = sources + # 复用爬虫只读端点 + from .opp_common import crawler_get + ok, res = await crawler_get(sor, "/api/demands", params) + if not ok: + return total, "拉取需求失败@offset%d: %s" % (offset, res) + items = (res or {}).get("items") or [] + matched = (res or {}).get("matched", 0) + for it in items: + src_id = str(it.get("id") or it.get("source_id") or "") + title = (it.get("title") or "").strip() + if not src_id or not title: + continue + # 指定类型分析:category 二次过滤(爬虫 item_category 模糊) + if category and category not in (it.get("item_category") or "") and category not in title: + continue + snap_id = new_id() + vec_id = demand_vec_id(it.get("source"), src_id) + await sor.C("opp_demand_snap", { + "id": snap_id, "batch_id": bid, "org_id": org_id, + "src_id": src_id, "source": it.get("source") or "", + "title": title[:500], "item_category": (it.get("item_category") or "")[:128], + "budget_wan": it.get("budget_wan") or 0, "url": (it.get("url") or "")[:500], + "publish_time": (it.get("publish_time") or "")[:32], + "embed_status": "new", + }) + await sor.sqlExe("COMMIT", {}) + # vec_id 暂存到内存映射稍后批量用;这里用 title 去重键 + total += 1 + offset += len(items) + if not items or offset >= matched: + break + return total, "" + + +# ══════════════════ 数据层:embedding 增量缓存 ══════════════════ + +async def _ensure_cache_collection(sor, base, col): + ok, d = await vdb_post(sor, base, "/v1/createcollection", { + "colname": col, "fields": [ + {"name": "id", "type": "str", "is_primary": True, "max_length": 64}, + {"name": "vector", "type": "fvector", "dim": DIM}, + {"name": "text", "type": "str", "max_length": 2000}], + "description": "opp demand embedding cache", "metric": "COSINE"}) + if not ok and "exist" not in str(d).lower(): + return False, d + return True, d + + +async def _cache_existing_ids(sor, base, col, ids): + """批量查缓存 collection 已存在哪些 id(expr in,每批≤200)。返回 set。""" + found = set() + for i in range(0, len(ids), 200): + chunk = ids[i:i + 200] + expr = "id in [%s]" % ",".join('"%s"' % x for x in chunk) + ok, d = await vdb_post(sor, base, "/v1/query", + {"colname": col, "expr": expr, "output_fields": ["id"]}) + if ok: + for r in _vec_rows(d): + if r.get("id"): + found.add(r["id"]) + return found + + +async def embed_snapshots(sor, batch, project_id, cfg, on_progress=None): + """增量 embedding:缓存缺的才调 rag_embed_texts,upsert 进共享缓存。 + 返回 (vec_id→vector dict, err)。一条需求只嵌一次(全机构复用)。""" + bid = batch["id"] + base = await _vdb_base(sor) + if not base: + return None, "upapp 缺 rag-vdb(向量库未配置)" + ok, d = await _ensure_cache_collection(sor, base, cfg["cache_col"]) + if not ok: + return None, "建缓存 collection 失败: %s" % d + + snaps = rows_to_dicts(await sor.sqlExe( + "SELECT id, source, src_id, title FROM opp_demand_snap WHERE batch_id=${b}$", + {"b": bid}), limit=100000) + await sor.sqlExe("COMMIT", {}) + if not snaps: + return None, "批次无需求快照" + + # 稳定 vec_id + 文本 + id_map = {} # snap_id -> vec_id + text_map = {} # vec_id -> text + for s in snaps: + vid = demand_vec_id(s["source"], s["src_id"]) + id_map[s["id"]] = vid + text_map[vid] = (s["title"] or "(无标题)")[:1900] + all_vids = list(set(id_map.values())) + + # 查缓存已有的 + cached = await _cache_existing_ids(sor, base, cfg["cache_col"], all_vids) + to_embed = [v for v in all_vids if v not in cached] + logger.debug("embed batch=%s total=%d cached=%d new=%d", bid, len(all_vids), len(cached), len(to_embed)) + + # 缺的调 rag_embed_texts + vecs = {} + if to_embed: + from pipeline_service import rag_client as rc + bs = cfg["opp_embed_batch"] + texts_to_embed = [text_map[v] for v in to_embed] + got, err = await rc.rag_embed_texts(project_id, texts_to_embed, batch_size=bs) + if err: + return None, "embedding 失败: %s" % err + if len(got) != len(to_embed): + return None, "embedding 数量不符(%d/%d)" % (len(got), len(to_embed)) + # upsert 进缓存 + rows = [{"id": to_embed[i], "vector": got[i], "text": texts_to_embed[i]} + for i in range(len(to_embed))] + for i in range(0, len(rows), 500): + ok, d = await vdb_post(sor, base, "/v1/upsert", + {"colname": cfg["cache_col"], "data": rows[i:i + 500]}) + if not ok: + return None, "缓存 upsert 失败: %s" % d + for i, v in enumerate(to_embed): + vecs[v] = got[i] + + # 缓存命中的也要取回向量(构建工作集用) + if cached: + cached = list(cached) + for i in range(0, len(cached), 200): + chunk = cached[i:i + 200] + expr = "id in [%s]" % ",".join('"%s"' % x for x in chunk) + ok, d = await vdb_post(sor, base, "/v1/query", + {"colname": cfg["cache_col"], "expr": expr, + "output_fields": ["id", "vector"]}) + if ok: + for r in _vec_rows(d): + if r.get("id") and isinstance(r.get("vector"), list): + vecs[r["id"]] = r["vector"] + # 更新快照 embed_status + for sid, vid in id_map.items(): + st = "cached" if vid in vecs else "failed" + await sor.sqlExe("UPDATE opp_demand_snap SET embed_status=${s}$ WHERE id=${i}$", + {"s": st, "i": sid}) + await sor.sqlExe("COMMIT", {}) + return {"id_map": id_map, "text_map": text_map, "vecs": vecs}, "" + + +# ══════════════════ 数据层:批次工作集 + kNN ══════════════════ + +async def build_batch_collection(sor, base, batch, cfg, embed_res): + """从缓存拷向量到批次专属工作集 collection(聚类只在本批内算)。返回 (col, err)。""" + col = "opp_mine_%s" % (batch["id"][:16].replace("-", "")) + ok, d = await vdb_post(sor, base, "/v1/createcollection", { + "colname": col, "fields": [ + {"name": "id", "type": "str", "is_primary": True, "max_length": 64}, + {"name": "vector", "type": "fvector", "dim": DIM}, + {"name": "text", "type": "str", "max_length": 2000}, + {"name": "snap_id", "type": "str", "max_length": 32}], + "description": "opp mining batch workset", "metric": "COSINE"}) + if not ok and "exist" not in str(d).lower(): + return None, "建批次 collection 失败: %s" % d + vecs = embed_res["vecs"] + text_map = embed_res["text_map"] + id_map = embed_res["id_map"] # snap_id -> vec_id + rows = [] + for sid, vid in id_map.items(): + if vid in vecs: + rows.append({"id": vid, "vector": vecs[vid], "text": text_map.get(vid, ""), "snap_id": sid}) + for i in range(0, len(rows), 500): + ok, d = await vdb_post(sor, base, "/v1/upsert", {"colname": col, "data": rows[i:i + 500]}) + if not ok: + return None, "工作集 upsert 失败: %s" % d + return col, "" + + +async def knn_graph(sor, base, col, vecs, topk): + """每条向量查 topk 邻居(排除自身)。返回 vid -> [(score, vid)]。""" + nbr = {} + ids = list(vecs.keys()) + for vid in ids: + ok, d = await vdb_post(sor, base, "/v1/query", { + "colname": col, "vector": vecs[vid], "pagerows": topk, + "output_fields": ["id"]}) + if not ok: + return None, "kNN 失败@%s: %s" % (vid, d) + rows = _vec_rows(d) + nbr[vid] = [(float(r.get("score", 0) or 0), r.get("id")) + for r in rows if r.get("id") and r.get("id") != vid] + return nbr, "" + + +# ══════════════════ 算法层:并查集聚类 ══════════════════ + +def _union_find(ids, nbr, tau): + parent = {v: v for v in ids} + + def find(x): + while parent[x] != x: + parent[x] = parent[parent[x]] + x = parent[x] + return x + + def union(a, b): + ra, rb = find(a), find(b) + if ra != rb: + parent[rb] = ra + for v in ids: + for sc, nid in nbr.get(v, []): + if sc >= tau and nid in parent: + union(v, nid) + comps = {} + for v in ids: + comps.setdefault(find(v), []).append(v) + return sorted(comps.values(), key=len, reverse=True) + + +def _cos(a, b): + return sum(x * y for x, y in zip(a, b)) + + +def _centroid(vecs, members): + n = len(members) + return [sum(vecs[m][k] for m in members) / n for k in range(len(vecs[members[0]]))] + + +def cluster_vectors(vecs, nbr, tau, min_cluster, big_split): + """kNN 图 → τ 并查集 → 大簇二次细分 → 小组质心并入。返回 (final_clusters, other_ids)。 + final_clusters: list[list[vid]],other_ids: 未入簇的 vid。""" + ids = list(vecs.keys()) + comps = _union_find(ids, nbr, tau) + # 超大簇质心二次细分(众包标题模板化重复会过度聚合) + pre = [] + for c in comps: + if len(c) > big_split: + members = set(c) + sub_nbr = {v: [(sc, n) for sc, n in nbr.get(v, []) if n in members] for v in c} + pre.extend(_union_find(c, sub_nbr, min(tau + 0.10, 0.95))) + else: + pre.append(c) + big = [c for c in pre if len(c) >= min_cluster] + small = [c for c in pre if len(c) < min_cluster] + # 小组并入最近大簇(质心相似度 ≥ τ-0.10),否则归其他 + other = [] + if big and small: + cents = [(c, _centroid(vecs, c)) for c in big] + for sc_ in small: + cc = _centroid(vecs, sc_) + best, bs = None, -1 + for bc, bcent in cents: + sim = _cos(cc, bcent) + if sim > bs: + best, bs = bc, sim + if best is not None and bs >= tau - 0.10: + best.extend(sc_) + else: + other.extend(sc_) + else: + for c in small: + other.extend(c) + big.sort(key=len, reverse=True) + return big, other + + +# ══════════════════ 算法层: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 = [] + total = len(vecs) + for ci, members in enumerate(clusters): + cent = _centroid(vecs, members) + sims = sorted(((_cos(vecs[m], cent), m) for m in members), reverse=True) + samples = [] + seen = set() + for _, m in sims: + t = text_map.get(m, "") + if t and t not in seen: + seen.add(t) + 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, + }) + return results + + +def _rule_name(samples): + """规则兜底:取样例标题的最长公共前缀片段或首条去模板词。""" + if not samples: + return "未命名类别" + # 去常见模板前缀 + strip = ["我需要", "需要", "其他", "服务需求", "服务采购", "服务", "需求", "采购", "合作"] + def clean(t): + for s in strip: + t = t.replace(s, "") + return t.strip() + cleaned = [clean(s) for s in samples if clean(s)] + if not cleaned: + return samples[0][:20] + # 取最高频的 2-gram 词 + from collections import Counter + cnt = Counter() + for t in cleaned: + for n in (4, 3, 2): + for i in range(len(t) - n + 1): + cnt[t[i:i + n]] += 1 + if cnt: + top = cnt.most_common(1)[0][0] + return top[:20] + return cleaned[0][:20] diff --git a/pipeline_opportunity/opp_mining_flow.py b/pipeline_opportunity/opp_mining_flow.py new file mode 100644 index 0000000..4d4f108 --- /dev/null +++ b/pipeline_opportunity/opp_mining_flow.py @@ -0,0 +1,253 @@ +# -*- coding:utf-8 -*- +"""挖掘批次编排:状态机驱动的全流程(pulling→embedding→clustering→naming→done/failed)。 + +入口 start_mining() 同步建批次记录后立即返回 batch_id,全流程 asyncio.create_task +后台跑(对齐 executor.py 范式);进度/错误实时落 opp_mining_batches(stats_json/error_msg), +前端轮询 opp_mining_status 可见——失败必须是可行动报错(用户铁律:禁前端挂起干等)。 +""" +import asyncio +import json +import logging +import time + +from .opp_common import get_db, new_id, rows_to_dicts +from .opp_mining import (get_mining_params, pull_demands, embed_snapshots, + build_batch_collection, knn_graph, cluster_vectors, + name_clusters, can_batch_transition, _vdb_base, vdb_post) + +logger = logging.getLogger("pipeline.opp_mining") + +_running = {} # batch_id -> asyncio.Task(防重复起跑,进程内) + + +async def _set_status(sor, batch_id, cur, nxt, **extra): + if not can_batch_transition(cur, nxt): + raise RuntimeError("非法状态迁移 %s→%s" % (cur, nxt)) + sets = ["status=${nxt}$", "updated_at=CURRENT_TIMESTAMP"] + args = {"nxt": nxt, "bid": batch_id} + for k, v in extra.items(): + sets.append("%s=${%s}$" % (k, k)) + args[k] = v + await sor.sqlExe( + "UPDATE opp_mining_batches SET " + ", ".join(sets) + " WHERE id=${bid}$", args) + await sor.sqlExe("COMMIT", {}) + + +async def start_mining(sor, ctx, scope="top", days=365, sources="", keyword="", + category="", created_by=""): + """创建挖掘批次并后台执行。返回 (batch_id, err)。 + + scope=top:全量需求 → 聚类 → TopX 排名。 + scope=targeted:用户指定类型(keyword/category 过滤)→ 范围内聚类找子型。 + """ + org_id = ctx.get("org_id") or "" + if scope not in ("top", "targeted"): + return "", "scope 必须是 top 或 targeted" + if scope == "targeted" and not (keyword or category): + return "", "targeted 模式必须指定 keyword 或 category(如:合同管理)" + bid = new_id() + params_json = json.dumps({"scope": scope, "days": days, "sources": sources, + "keyword": keyword, "category": category}, ensure_ascii=False) + await sor.C("opp_mining_batches", { + "id": bid, "org_id": org_id, "scope": scope, "params_json": params_json, + "status": "pulling", "vdb_col": "", "stats_json": "", "error_msg": "", + "created_by": created_by or ctx.get("user_id") or "", + }) + await sor.sqlExe("COMMIT", {}) + if bid in _running and not _running[bid].done(): + return bid, "" + task = asyncio.create_task(_run_mining(bid, ctx)) + _running[bid] = task + return bid, "" + + +async def _run_mining(batch_id, ctx): + db, DBNAME = get_db() + try: + async with db.sqlorContext(DBNAME) as sor: + await _run_mining_inner(sor, batch_id, ctx) + except Exception as e: + logger.exception("mining batch %s crashed", batch_id) + try: + async with db.sqlorContext(DBNAME) as sor: + recs = await sor.sqlExe( + "SELECT status FROM opp_mining_batches WHERE id=${b}$", {"b": batch_id}) + await sor.sqlExe("COMMIT", {}) + cur = recs[0].status if recs else "pulling" + await _set_status(sor, batch_id, cur, "failed", + error_msg="批次异常: %s" % str(e)[:500]) + except Exception: + pass + + +async def _run_mining_inner(sor, batch_id, ctx): + t0 = time.time() + recs = await sor.sqlExe( + "SELECT * FROM opp_mining_batches WHERE id=${b}$", {"b": batch_id}) + await sor.sqlExe("COMMIT", {}) + if not recs: + raise RuntimeError("批次不存在: %s" % batch_id) + batch = rows_to_dicts(recs, limit=1)[0] + params = json.loads(batch.get("params_json") or "{}") + cfg = await get_mining_params(sor) + + def stats(**kw): + kw["elapsed_s"] = round(time.time() - t0, 1) + return json.dumps(kw, ensure_ascii=False) + + # ── 1. pulling ── + n, err = await pull_demands( + sor, batch, ctx, cfg, params.get("scope", "top"), + keyword=params.get("keyword", ""), category=params.get("category", ""), + days=int(params.get("days") or 365), sources=params.get("sources", "")) + if err: + return await _set_status(sor, batch_id, "pulling", "failed", + error_msg=err, stats_json=stats(demands=n)) + if n == 0: + return await _set_status(sor, batch_id, "pulling", "done", + stats_json=stats(demands=0, clusters=0, + note="范围内无需求数据")) + await _set_status(sor, batch_id, "pulling", "embedding", + stats_json=stats(demands=n)) + + # ── 2. embedding(增量缓存,失败可行动报错)── + project_id = ctx.get("project_id") or "" + embed_res, err = await embed_snapshots(sor, batch, project_id, cfg) + if err: + return await _set_status(sor, batch_id, "embedding", "failed", + error_msg=err, stats_json=stats(demands=n)) + vecs = embed_res["vecs"] + text_map = embed_res["text_map"] + if not vecs: + return await _set_status(sor, batch_id, "embedding", "failed", + error_msg="全部需求向量化失败", stats_json=stats(demands=n)) + await _set_status(sor, batch_id, "embedding", "clustering", + stats_json=stats(demands=n, embedded=len(vecs))) + + # ── 3. clustering(工作集 + kNN + 并查集)── + base = await _vdb_base(sor) + col, err = await build_batch_collection(sor, base, batch, cfg, embed_res) + if err: + return await _set_status(sor, batch_id, "clustering", "failed", + error_msg=err, stats_json=stats(demands=n, embedded=len(vecs))) + await sor.sqlExe("UPDATE opp_mining_batches SET vdb_col=${c}$ WHERE id=${b}$", + {"c": col, "b": batch_id}) + await sor.sqlExe("COMMIT", {}) + nbr, err = await knn_graph(sor, base, col, vecs, cfg["opp_mine_topk"]) + if err: + return await _set_status(sor, batch_id, "clustering", "failed", + error_msg=err, stats_json=stats(demands=n, embedded=len(vecs))) + clusters, other = cluster_vectors(vecs, nbr, cfg["opp_mine_tau"], + cfg["opp_mine_min_cluster"], cfg["opp_mine_big_split"]) + await _set_status(sor, batch_id, "clustering", "naming", + stats_json=stats(demands=n, embedded=len(vecs), + clusters=len(clusters), other=len(other), + tau=cfg["opp_mine_tau"])) + + # ── 4. naming(LLM utility + 规则兜底)── + named = await name_clusters(sor, batch, clusters, vecs, text_map, cfg, ctx) + org_id = ctx.get("org_id") or "" + # vid → snap_id 反查 + vid2snap = {v: s for s, v in embed_res["id_map"].items()} + for cl in named: + cid = new_id() + cent_snap = vid2snap.get(cl["centroid_vid"], "") + await sor.C("opp_clusters", { + "id": cid, "batch_id": batch_id, "org_id": org_id, + "name": cl["name"][:128], "doc_count": cl["size"], "share": cl["share"], + "heat_rank": cl["rank"], "naming_evidence": json.dumps( + {"samples": cl["samples"], "by": cl["evidence"]}, ensure_ascii=False), + "centroid_snap_id": cent_snap, + }) + await sor.sqlExe("COMMIT", {}) + # 回填快照 cluster_id(分批 UPDATE,参数化防注入) + snaps = [vid2snap[v] for v in cl["members"] if v in vid2snap] + for i in range(0, len(snaps), 200): + chunk = snaps[i:i + 200] + ph = ",".join("${s%d}$" % j for j in range(len(chunk))) + args = {"cid": cid} + for j, sv in enumerate(chunk): + args["s%d" % j] = sv + await sor.sqlExe( + "UPDATE opp_demand_snap SET cluster_id=${cid}$ WHERE id IN (%s)" % ph, args) + await sor.sqlExe("COMMIT", {}) + + await _set_status(sor, batch_id, "naming", "done", + stats_json=stats(demands=n, embedded=len(vecs), + clusters=len(named), other=len(other), + tau=cfg["opp_mine_tau"], + top1=(named[0]["name"] if named else ""))) + + # ── 5. 清理:保留最近 N 批工作集,旧批 drop(缓存 collection 永不清)── + await _cleanup_old_collections(sor, base, org_id, cfg) + + +async def _cleanup_old_collections(sor, base, org_id, cfg): + keep = cfg["opp_mine_keep_batches"] + recs = await sor.sqlExe( + "SELECT id, vdb_col FROM opp_mining_batches WHERE org_id=${o}$ " + "ORDER BY created_at DESC LIMIT 200", {"o": org_id}) + await sor.sqlExe("COMMIT", {}) + cols = [r.vdb_col for r in (recs or []) if getattr(r, "vdb_col", "")] + for col in cols[keep:]: + await vdb_post(sor, base, "/v1/dropcollection", {"colname": col}) + await sor.sqlExe("UPDATE opp_mining_batches SET vdb_col='' WHERE vdb_col=${c}$", + {"c": col}) + await sor.sqlExe("COMMIT", {}) + logger.debug("dropped old batch collection %s", col) + + +async def mining_status(sor, ctx, batch_id=""): + """查批次状态(org 隔离)。batch_id 空 = 本机构最近批次列表。""" + org_id = ctx.get("org_id") or "" + if batch_id: + recs = await sor.sqlExe( + "SELECT id, scope, status, stats_json, error_msg, vdb_col, created_at " + "FROM opp_mining_batches WHERE id=${b}$ AND org_id=${o}$", + {"b": batch_id, "o": org_id}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return False, "批次不存在或不属于当前机构" + return True, rows_to_dicts(recs, limit=1)[0] + recs = await sor.sqlExe( + "SELECT id, scope, status, stats_json, error_msg, created_at " + "FROM opp_mining_batches WHERE org_id=${o}$ ORDER BY created_at DESC LIMIT 10", + {"o": org_id}) + await sor.sqlExe("COMMIT", {}) + return True, rows_to_dicts(recs, limit=10) + + +async def list_clusters(sor, ctx, batch_id): + """类别排名表(org 隔离:批次必须属于本机构)。""" + org_id = ctx.get("org_id") or "" + recs = await sor.sqlExe( + "SELECT id FROM opp_mining_batches WHERE id=${b}$ AND org_id=${o}$", + {"b": batch_id, "o": org_id}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return False, "批次不存在或不属于当前机构" + recs = await sor.sqlExe( + "SELECT id, name, doc_count, share, heat_rank, naming_evidence, centroid_snap_id " + "FROM opp_clusters WHERE batch_id=${b}$ AND org_id=${o}$ ORDER BY heat_rank ASC", + {"b": batch_id, "o": org_id}) + await sor.sqlExe("COMMIT", {}) + return True, rows_to_dicts(recs, limit=200) + + +async def cluster_detail(sor, ctx, cluster_id, limit=20): + """类内需求明细(样例,带来源 URL 可查证)。""" + org_id = ctx.get("org_id") or "" + recs = await sor.sqlExe( + "SELECT id, name, doc_count, share, heat_rank, naming_evidence " + "FROM opp_clusters WHERE id=${c}$ AND org_id=${o}$", + {"c": cluster_id, "o": org_id}) + await sor.sqlExe("COMMIT", {}) + if not recs: + return False, "类别不存在或不属于当前机构" + cl = rows_to_dicts(recs, limit=1)[0] + items = rows_to_dicts(await sor.sqlExe( + "SELECT id, title, source, budget_wan, url, publish_time FROM opp_demand_snap " + "WHERE cluster_id=${c}$ ORDER BY budget_wan DESC LIMIT ${n}$", + {"c": cluster_id, "n": int(limit)}), limit=int(limit)) + await sor.sqlExe("COMMIT", {}) + return True, {"cluster": cl, "samples": items} diff --git a/scripts/gen_ddl.sh b/scripts/gen_ddl.sh old mode 100644 new mode 100755 index 408ffc7..c9dd9c7 --- a/scripts/gen_ddl.sh +++ b/scripts/gen_ddl.sh @@ -1,6 +1,16 @@ #!/bin/bash # 生成 pipeline-opportunity 的 MySQL DDL(json2ddl 需要 sqlor + appPublic 在 PYTHONPATH) +# 路径以脚本位置推导(2026-09-11 修:原写死 /home/ymq/work/pipeline 旧路径失效) set -e -R=/home/ymq/work/repos +R=$(cd "$(dirname "$0")/../.." && pwd) export PYTHONPATH="$R/xls2ddl:$R/sqlor:$R/apppublic" -python3 "$R/xls2ddl/gen_opportunity_ddl.py" +python3 - "$R" <<'PYEOF' +import sys +from xls2ddl.json2ddl import model2ddl +R = sys.argv[1] +s = model2ddl(R + '/pipeline-opportunity/models', 'mysql') +out = R + '/pipeline-opportunity/mysql.ddl.sql' +with open(out, 'w', encoding='utf-8') as f: + f.write(s) +print('written', out, len(s), 'chars') +PYEOF