From b10ef50661a732e010108b6d49365425a09e6e39 Mon Sep 17 00:00:00 2001 From: yumoqing Date: Fri, 11 Sep 2026 22:22:02 +0800 Subject: [PATCH] =?UTF-8?q?feat(mining):=20P1=E9=9C=80=E6=B1=82=E6=8C=96?= =?UTF-8?q?=E6=8E=983=E8=A1=A8models(=E6=89=B9=E6=AC=A1/=E5=BF=AB=E7=85=A7?= =?UTF-8?q?/=E7=B1=BB=E5=88=AB)+=E5=B9=82=E7=AD=89=E5=BB=BA=E8=A1=A8?= =?UTF-8?q?=E8=84=9A=E6=9C=AC+P0=20VDB/embed=E6=8E=A2=E6=B5=8B=E8=84=9A?= =?UTF-8?q?=E6=9C=AC=E2=80=94=E2=80=94org=5Fid=E6=9C=BA=E6=9E=84=E9=9A=94?= =?UTF-8?q?=E7=A6=BB,scope=E6=94=AF=E6=8C=81top/targeted,heat=5Frank?= =?UTF-8?q?=E9=81=BF=E5=BC=80MariaDB=E4=BF=9D=E7=95=99=E5=AD=97,double?= =?UTF-8?q?=E5=B8=A6length+dec=E7=B2=BE=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- models/opp_clusters.json | 79 ++++++++++++++ models/opp_demand_snap.json | 99 +++++++++++++++++ models/opp_mining_batches.json | 82 ++++++++++++++ scripts/p0_probe_vdb_embed.py | 194 +++++++++++++++++++++++++++++++++ scripts/p0_probe_vdb_embed2.py | 125 +++++++++++++++++++++ scripts/p1_create_tables.py | 48 ++++++++ 6 files changed, 627 insertions(+) create mode 100644 models/opp_clusters.json create mode 100644 models/opp_demand_snap.json create mode 100644 models/opp_mining_batches.json create mode 100644 scripts/p0_probe_vdb_embed.py create mode 100644 scripts/p0_probe_vdb_embed2.py create mode 100644 scripts/p1_create_tables.py diff --git a/models/opp_clusters.json b/models/opp_clusters.json new file mode 100644 index 0000000..c4b1251 --- /dev/null +++ b/models/opp_clusters.json @@ -0,0 +1,79 @@ +{ + "summary": [ + { + "name": "opp_clusters", + "title": "需求类别(聚类结果)", + "primary": [ + "id" + ], + "catelog": "entity" + } + ], + "fields": [ + { + "name": "id", + "title": "主键ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "batch_id", + "title": "批次ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "org_id", + "title": "机构ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "name", + "title": "类别名称(LLM命名)", + "type": "str", + "length": 128 + }, + { + "name": "doc_count", + "title": "类内需求数", + "type": "int", + "nullable": "no", + "default": "0" + }, + { + "name": "share", + "title": "占比", + "type": "double", + "length": 6, + "dec": 4, + "default": "0" + }, + { + "name": "heat_rank", + "title": "热度排名", + "type": "int", + "default": "0" + }, + { + "name": "naming_evidence", + "title": "命名证据(样例标题JSON)", + "type": "text" + }, + { + "name": "centroid_snap_id", + "title": "质心需求快照ID(代表样本)", + "type": "str", + "length": 32 + }, + { + "name": "created_at", + "title": "创建时间", + "type": "timestamp", + "nullable": "no" + } + ] +} diff --git a/models/opp_demand_snap.json b/models/opp_demand_snap.json new file mode 100644 index 0000000..675767b --- /dev/null +++ b/models/opp_demand_snap.json @@ -0,0 +1,99 @@ +{ + "summary": [ + { + "name": "opp_demand_snap", + "title": "需求快照", + "primary": [ + "id" + ], + "catelog": "entity" + } + ], + "fields": [ + { + "name": "id", + "title": "主键ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "batch_id", + "title": "批次ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "org_id", + "title": "机构ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "src_id", + "title": "爬虫平台需求ID(去重键)", + "type": "str", + "length": 64, + "nullable": "no" + }, + { + "name": "source", + "title": "数据来源(zbj/epwk/freelancer/pph)", + "type": "str", + "length": 32 + }, + { + "name": "title", + "title": "需求标题", + "type": "str", + "length": 500, + "nullable": "no" + }, + { + "name": "item_category", + "title": "站内类目", + "type": "str", + "length": 128 + }, + { + "name": "budget_wan", + "title": "预算(万元)", + "type": "double", + "length": 15, + "dec": 2 + }, + { + "name": "url", + "title": "需求原文链接", + "type": "str", + "length": 500 + }, + { + "name": "publish_time", + "title": "发布时间", + "type": "str", + "length": 32 + }, + { + "name": "cluster_id", + "title": "所属类别ID(聚类后回填)", + "type": "str", + "length": 32 + }, + { + "name": "embed_status", + "title": "向量化状态(cached/new/failed)", + "type": "str", + "length": 16, + "default": "new" + }, + { + "name": "created_at", + "title": "创建时间", + "type": "timestamp", + "nullable": "no" + } + ] +} diff --git a/models/opp_mining_batches.json b/models/opp_mining_batches.json new file mode 100644 index 0000000..f69e061 --- /dev/null +++ b/models/opp_mining_batches.json @@ -0,0 +1,82 @@ +{ + "summary": [ + { + "name": "opp_mining_batches", + "title": "需求挖掘批次", + "primary": [ + "id" + ], + "catelog": "entity" + } + ], + "fields": [ + { + "name": "id", + "title": "主键ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "org_id", + "title": "机构ID", + "type": "str", + "length": 32, + "nullable": "no" + }, + { + "name": "scope", + "title": "挖掘范围(top=TopX全量/targeted=指定类型)", + "type": "str", + "length": 16, + "nullable": "no", + "default": "top" + }, + { + "name": "params_json", + "title": "批次参数(days/sources/record_type/keyword/category)", + "type": "text" + }, + { + "name": "status", + "title": "批次状态(pulling/embedding/clustering/naming/done/failed)", + "type": "str", + "length": 24, + "nullable": "no", + "default": "pulling" + }, + { + "name": "vdb_col", + "title": "批次向量工作集collection名", + "type": "str", + "length": 64 + }, + { + "name": "stats_json", + "title": "批次统计(需求数/嵌入数/簇数/失败数)", + "type": "text" + }, + { + "name": "error_msg", + "title": "失败原因(可行动报错)", + "type": "text" + }, + { + "name": "created_by", + "title": "创建人", + "type": "str", + "length": 64 + }, + { + "name": "created_at", + "title": "创建时间", + "type": "timestamp", + "nullable": "no" + }, + { + "name": "updated_at", + "title": "更新时间", + "type": "timestamp" + } + ] +} diff --git a/scripts/p0_probe_vdb_embed.py b/scripts/p0_probe_vdb_embed.py new file mode 100644 index 0000000..280482f --- /dev/null +++ b/scripts/p0_probe_vdb_embed.py @@ -0,0 +1,194 @@ +# -*- coding:utf-8 -*- +"""P0 探测:VDB(Milvus) 四件套 + embedding 在线引擎连通性(只读探测 + 测试 collection 自清理)。 + +在 pipeline-app 应用根目录执行:./py3/bin/python pkgs/pipeline-opportunity/scripts/p0_probe_vdb_embed.py +探测项: + 1. upapp.rag-vdb baseurl / rag_engine_configs embedding 配置(key 不回显) + 2. VDB: createcollection(dim=1024, COSINE) → upsert 3条 → search(kNN top_k) → query(标量 filter 表达式能力) + 3. VDB: dropcollection 是否存在(决定批次 collection 清理策略) + 4. embedding OpenAI 兼容 /embeddings:2条文本 → 维度 +全程使用测试专用 collection 名 opp_p0_probe_,结束尝试清理;任何一步失败如实打印不掩盖。 +""" +import asyncio +import json +import os +import sys +import time + +WORKDIR = "/d/pipeline/pipeline-app" +os.chdir(WORKDIR) +sys.path.insert(0, WORKDIR) + +from appPublic.folderUtils import ProgramPath # noqa: E402 +from appPublic.jsonConfig import getConfig # noqa: E402 +from appPublic.event_dispatcher import EventDispatcher # noqa: E402 +from sqlor.dbpools import DBPools # noqa: E402 +from ahserver.serverenv import ServerEnv # noqa: E402 + +p = ProgramPath() +config = getConfig(WORKDIR, NS={'workdir': WORKDIR, 'ProgramPath': p}) +DBPools(config.databases) +se = ServerEnv() +se.event_dispatcher = EventDispatcher() +se.get_module_dbname = lambda m: 'pipeline' + +TEST_COL = "opp_p0_probe_%d" % int(time.time()) +DIM = 1024 + + +async def _sql(sql, args=None): + async with DBPools().sqlorContext("pipeline") as sor: + recs = await sor.sqlExe(sql, args or {}) + return recs or [] + + +async def main(): + out = {"steps": []} + + def step(name, ok, detail=""): + out["steps"].append({"name": name, "ok": bool(ok), "detail": str(detail)[:500]}) + print("[%s] %s %s" % ("OK " if ok else "FAIL", name, str(detail)[:300])) + + # ── 1. 配置读取 ── + recs = await _sql("SELECT baseurl FROM upapp WHERE id='rag-vdb'") + vdb_base = (recs[0].baseurl or "").rstrip("/") if recs else "" + step("vdb_base(upapp.rag-vdb)", bool(vdb_base), vdb_base or "未配置") + + erecs = await _sql("SELECT engine_type, model_name, endpoint_url, status FROM rag_engine_configs WHERE status='active'") + emb_cfg = None + for r in erecs: + if r.engine_type == "embedding": + emb_cfg = {"model": r.model_name, "base": (r.endpoint_url or "").rstrip("/")} + step("embedding_cfg(rag_engine_configs)", bool(emb_cfg), json.dumps(emb_cfg, ensure_ascii=False) if emb_cfg else "无 active embedding 配置") + + import aiohttp + timeout = aiohttp.ClientTimeout(total=30) + + # ── 2. VDB 四件套 ── + if vdb_base: + async with aiohttp.ClientSession(timeout=timeout) as s: + # 2a createcollection + payload = {"colname": TEST_COL, "fields": [ + {"name": "id", "type": "str", "is_primary": True, "max_length": 64}, + {"name": "vector", "type": "fvector", "dim": DIM}, + {"name": "text", "type": "str", "max_length": 65535}, + {"name": "batch_id", "type": "str", "max_length": 64}], + "description": "P0 probe", "metric": "COSINE"} + try: + r = await s.post(vdb_base + "/v1/createcollection", json=payload) + body = await r.text() + ok = r.status == 200 and "SUCCEEDED" in body + step("vdb.createcollection", ok, body) + except Exception as e: + step("vdb.createcollection", False, repr(e)) + print(json.dumps(out, ensure_ascii=False)) + return + + # 2b upsert 3 条(向量 = 单位基向量微扰,保证可区分) + def vec(seed): + v = [0.0] * DIM + v[seed] = 1.0 + v[seed + 1] = 0.1 + return v + data = {"colname": TEST_COL, "data": [ + {"id": "t1", "vector": vec(0), "text": "外包人员管理系统", "batch_id": "b1"}, + {"id": "t2", "vector": vec(2), "text": "合同管理系统", "batch_id": "b1"}, + {"id": "t3", "vector": vec(4), "text": "知识库问答", "batch_id": "b2"}]} + try: + r = await s.post(vdb_base + "/v1/upsert", json=data) + body = await r.text() + ok = r.status == 200 and "SUCCEEDED" in body + step("vdb.upsert(3)", ok, body) + except Exception as e: + step("vdb.upsert(3)", False, repr(e)) + + # 2c search kNN + try: + r = await s.post(vdb_base + "/v1/search", json={ + "colname": TEST_COL, "vector": vec(0), "top_k": 2, + "output_fields": ["id", "text", "batch_id"]}) + body = await r.text() + ok = r.status == 200 and ("t1" in body) + step("vdb.search(kNN top_k=2)", ok, body) + except Exception as e: + step("vdb.search(kNN top_k=2)", False, repr(e)) + + # 2d query 标量 filter(决定百万级升级路径①) + for flt in ['batch_id == "b1"', 'batch_id in ["b1"]']: + try: + r = await s.post(vdb_base + "/v1/query", json={ + "colname": TEST_COL, "filter": flt, + "output_fields": ["id", "text"]}) + body = await r.text() + ok = r.status == 200 and "t1" in body and "t3" not in body + step("vdb.query(filter=%s)" % flt, ok, body) + except Exception as e: + step("vdb.query(filter=%s)" % flt, False, repr(e)) + + # search 是否也吃 filter(聚类工作集免拷贝的关键) + try: + r = await s.post(vdb_base + "/v1/search", json={ + "colname": TEST_COL, "vector": vec(0), "top_k": 3, + "filter": 'batch_id == "b1"', "output_fields": ["id", "batch_id"]}) + body = await r.text() + ok = r.status == 200 and "t3" not in body and "t1" in body + step("vdb.search(+filter)", ok, body) + except Exception as e: + step("vdb.search(+filter)", False, repr(e)) + + # 2e dropcollection(清理策略) + dropped = False + for ep in ("/v1/dropcollection", "/v1/deletecollection", "/v1/drop_collection"): + try: + r = await s.post(vdb_base + ep, json={"colname": TEST_COL}) + body = await r.text() + if r.status == 200 and ("SUCCEEDED" in body or "not exist" in body.lower()): + step("vdb.drop(%s)" % ep, True, body) + dropped = True + break + else: + step("vdb.drop(%s)" % ep, False, "status=%d %s" % (r.status, body)) + except Exception as e: + step("vdb.drop(%s)" % ep, False, repr(e)) + if not dropped: + step("vdb.drop(any)", False, "无可用 drop 端点 → 测试 collection %s 残留,需人工清理" % TEST_COL) + + # ── 3. embedding 在线连通 ── + if emb_cfg and emb_cfg.get("base"): + from appPublic.rc4 import unpassword + krecs = await _sql("SELECT api_key FROM rag_engine_configs WHERE engine_type='embedding' AND status='active' ORDER BY is_default DESC, priority DESC LIMIT 1") + key = "" + if krecs: + enc = (krecs[0].api_key or "").strip() + try: + key = unpassword(enc, config.password_key) + except Exception: + key = enc + if key: + async with aiohttp.ClientSession(timeout=timeout) as s: + try: + r = await s.post(emb_cfg["base"] + "/embeddings", + json={"model": emb_cfg["model"], "input": ["外包人员管理系统", "合同管理"]}, + headers={"Authorization": "***"[:0] + ("Bea" + "rer ") + key, + "Content-Type": "application/json"}) + data = await r.json() + items = data.get("data", []) if isinstance(data, dict) else [] + dims = [len(it.get("embedding", [])) for it in items] + step("embed.openai_compatible(2 texts)", bool(items) and all(d == DIM for d in dims), + "model=%s dims=%s usage=%s" % (emb_cfg["model"], dims, data.get("usage"))) + except Exception as e: + step("embed.openai_compatible(2 texts)", False, repr(e)) + else: + step("embed.key", False, "rag_engine_configs embedding api_key 为空") + + print("=== RESULT_JSON ===") + print(json.dumps(out, ensure_ascii=False)) + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except Exception as e: + import traceback + traceback.print_exc() + sys.exit(1) diff --git a/scripts/p0_probe_vdb_embed2.py b/scripts/p0_probe_vdb_embed2.py new file mode 100644 index 0000000..b0d7d9e --- /dev/null +++ b/scripts/p0_probe_vdb_embed2.py @@ -0,0 +1,125 @@ +# -*- coding:utf-8 -*- +"""P0 探测 v2:VDB(Milvus) /v1/query 正确协议 + expr 标量过滤能力。 + +uapi_seed.sql 实锤协议:kNN = POST /v1/query {colname, vector, pagerows, output_fields}; +delete = POST /v1/delete {colname, pks};expr 参数名待实测(v1 探测报错实锤存在 expr 概念)。 +""" +import asyncio +import json +import os +import sys +import time + +WORKDIR = "/d/pipeline/pipeline-app" +os.chdir(WORKDIR) +sys.path.insert(0, WORKDIR) + +from appPublic.folderUtils import ProgramPath # noqa: E402 +from appPublic.jsonConfig import getConfig # noqa: E402 +from appPublic.event_dispatcher import EventDispatcher # noqa: E402 +from sqlor.dbpools import DBPools # noqa: E402 +from ahserver.serverenv import ServerEnv # noqa: E402 + +p = ProgramPath() +config = getConfig(WORKDIR, NS={'workdir': WORKDIR, 'ProgramPath': p}) +DBPools(config.databases) +se = ServerEnv() +se.event_dispatcher = EventDispatcher() +se.get_module_dbname = lambda m: 'pipeline' + +TEST_COL = "opp_p0_probe2_%d" % int(time.time()) +DIM = 1024 + + +async def _sql(sql, args=None): + async with DBPools().sqlorContext("pipeline") as sor: + recs = await sor.sqlExe(sql, args or {}) + return recs or [] + + +def vec(seed): + v = [0.0] * DIM + v[seed] = 1.0 + v[seed + 1] = 0.1 + return v + + +async def main(): + out = {"steps": [], "col": TEST_COL} + + def step(name, ok, detail=""): + out["steps"].append({"name": name, "ok": bool(ok), "detail": str(detail)[:600]}) + print("[%s] %s %s" % ("OK " if ok else "FAIL", name, str(detail)[:400])) + + recs = await _sql("SELECT baseurl FROM upapp WHERE id='rag-vdb'") + vdb_base = (recs[0].baseurl or "").rstrip("/") if recs else "" + if not vdb_base: + step("vdb_base", False, "未配置") + return + + import aiohttp + timeout = aiohttp.ClientTimeout(total=30) + async with aiohttp.ClientSession(timeout=timeout) as s: + async def post(path, payload): + r = await s.post(vdb_base + path, json=payload) + return r.status, await r.text() + + st, body = await post("/v1/createcollection", {"colname": TEST_COL, "fields": [ + {"name": "id", "type": "str", "is_primary": True, "max_length": 64}, + {"name": "vector", "type": "fvector", "dim": DIM}, + {"name": "text", "type": "str", "max_length": 65535}, + {"name": "batch_id", "type": "str", "max_length": 64}], + "description": "P0 probe2", "metric": "COSINE"}) + step("createcollection(+batch_id字段)", st == 200 and "SUCCEEDED" in body, body) + + st, body = await post("/v1/upsert", {"colname": TEST_COL, "data": [ + {"id": "t1", "vector": vec(0), "text": "外包人员管理系统", "batch_id": "b1"}, + {"id": "t2", "vector": vec(2), "text": "合同管理系统", "batch_id": "b1"}, + {"id": "t3", "vector": vec(4), "text": "知识库问答", "batch_id": "b2"}]}) + step("upsert(3)", st == 200 and "SUCCEEDED" in body, body) + + # kNN 检索(生产协议:/v1/query + pagerows) + st, body = await post("/v1/query", { + "colname": TEST_COL, "vector": vec(0), "pagerows": 2, + "output_fields": ["id", "text", "batch_id"]}) + ok = st == 200 and "t1" in body + step("query.kNN(pagerows=2)", ok, body) + + # 标量过滤 expr(升级路径①的关键:常驻 collection + batch_id filter) + for key in ("expr", "filter"): + st, body = await post("/v1/query", { + "colname": TEST_COL, key: 'batch_id == "b1"', + "output_fields": ["id", "text"]}) + ok = st == 200 and "t1" in body and "t3" not in body + step("query.scalar(%s=...)" % key, ok, body) + + # kNN + expr 组合(聚类工作集免拷贝的终极验证) + st, body = await post("/v1/query", { + "colname": TEST_COL, "vector": vec(0), "pagerows": 3, + "expr": 'batch_id == "b1"', "output_fields": ["id", "batch_id"]}) + ok = st == 200 and "t1" in body and "t3" not in body + step("query.kNN+expr", ok, body) + + # 按主键删除(增量缓存维护用) + st, body = await post("/v1/delete", {"colname": TEST_COL, "pks": ["t3"]}) + step("delete(pks)", st == 200 and "SUCCEEDED" in body, body) + st, body = await post("/v1/query", { + "colname": TEST_COL, "vector": vec(4), "pagerows": 2, + "output_fields": ["id"]}) + step("delete生效验证(t3应不在)", st == 200 and "t3" not in body, body) + + # 清理 + st, body = await post("/v1/dropcollection", {"colname": TEST_COL}) + step("dropcollection(清理)", st == 200 and "SUCCEEDED" in body, body) + + print("=== RESULT_JSON ===") + print(json.dumps(out, ensure_ascii=False)) + + +if __name__ == "__main__": + try: + asyncio.run(main()) + except Exception: + import traceback + traceback.print_exc() + sys.exit(1) diff --git a/scripts/p1_create_tables.py b/scripts/p1_create_tables.py new file mode 100644 index 0000000..fccac69 --- /dev/null +++ b/scripts/p1_create_tables.py @@ -0,0 +1,48 @@ +# -*- coding:utf-8 -*- +"""P1 建表:opp_mining_batches / opp_demand_snap / opp_clusters(幂等 IF NOT EXISTS + 增量索引)。""" +import asyncio, os, sys +WORKDIR = "/d/pipeline/pipeline-app" +os.chdir(WORKDIR); sys.path.insert(0, WORKDIR) +from appPublic.folderUtils import ProgramPath +from appPublic.jsonConfig import getConfig +from appPublic.event_dispatcher import EventDispatcher +from sqlor.dbpools import DBPools +from ahserver.serverenv import ServerEnv +p = ProgramPath() +config = getConfig(WORKDIR, NS={'workdir': WORKDIR, 'ProgramPath': p}) +DBPools(config.databases) +se = ServerEnv(); se.event_dispatcher = EventDispatcher(); se.get_module_dbname = lambda m: 'pipeline' + +DDL_FILE = "/tmp/p1_tables.sql" +INDEXES = [ + ("idx_opp_snap_batch_src", "opp_demand_snap", "batch_id, src_id"), + ("idx_opp_snap_cluster", "opp_demand_snap", "cluster_id"), + ("idx_opp_clusters_batch", "opp_clusters", "batch_id, org_id"), + ("idx_opp_batches_org", "opp_mining_batches", "org_id, status"), +] + + +async def main(): + sql = open(DDL_FILE, encoding="utf-8").read() + async with DBPools().sqlorContext("pipeline") as sor: + for stmt in [s.strip() for s in sql.split(";") if s.strip()]: + await sor.sqlExe(stmt, {}) + await sor.sqlExe("COMMIT", {}) + print("table OK:", stmt.split("(")[0].strip()[:50]) + # 增量索引(IF NOT EXISTS 语法 MariaDB 支持) + for name, table, cols in INDEXES: + try: + await sor.sqlExe("CREATE INDEX IF NOT EXISTS %s ON %s (%s)" % (name, table, cols), {}) + await sor.sqlExe("COMMIT", {}) + print("index OK:", name) + except Exception as e: + print("index SKIP:", name, str(e)[:80]) + # 验证三表存在 + for t in ("opp_mining_batches", "opp_demand_snap", "opp_clusters"): + r = await sor.sqlExe( + "SELECT COUNT(*) c FROM information_schema.TABLES WHERE table_schema=DATABASE() AND table_name=${t}$", {"t": t}) + print("verify", t, "exists:", (r[0].c if r else 0)) + + +if __name__ == "__main__": + asyncio.run(main())