feat(mining): P1需求挖掘3表models(批次/快照/类别)+幂等建表脚本+P0 VDB/embed探测脚本——org_id机构隔离,scope支持top/targeted,heat_rank避开MariaDB保留字,double带length+dec精度

This commit is contained in:
yumoqing 2026-09-11 22:22:02 +08:00
parent 5c37341f4c
commit b10ef50661
6 changed files with 627 additions and 0 deletions

79
models/opp_clusters.json Normal file
View File

@ -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"
}
]
}

View File

@ -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"
}
]
}

View File

@ -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"
}
]
}

View File

@ -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_<ts>,结束尝试清理;任何一步失败如实打印不掩盖。
"""
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)

View File

@ -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)

View File

@ -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())