feat(mining): P1聚类质量分析脚本——多τ对比(0.70/0.75/0.82/0.88)+超大簇(>150)质心二次细分(τ+0.10)+向量/kNN文件缓存(省embedding费用)+样例去重;2075条真实数据实测167s,τ=0.75为最优拐点(42簇,top1=8.2%微信小程序,语义纯净无失控巨簇)

This commit is contained in:
yumoqing 2026-09-11 22:49:56 +08:00
parent 4ad5dfda81
commit 091d9db02d

View File

@ -0,0 +1,178 @@
# -*- coding:utf-8 -*-
"""P1 聚类质量分析(多τ对比 + 超大簇二次细分 + 向量/kNN缓存)。
在 pipeline-app 测试机:P1_PROJECT_ID=<真实商机项目> ./py3/bin/python /tmp/p1_cluster_analyze.py
缓存:/tmp/p1_vec_cache.json(向量+文本)、/tmp/p1_knn_cache.json(kNN图)——存在则跳过 embedding/kNN(省费用)。
输出 /tmp/p1_analyze_result.json + 控制台对比。
"""
import asyncio, json, os, sys, time
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
from ahserver.globalEnv import initEnv
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'
initEnv()
TAUS=[0.70,0.75,0.82,0.88]
TOPK=30; MIN_CLUSTER=5; DIM=1024; BIG_SPLIT=150 # 超过此规模的簇二次细分
VEC_CACHE="/tmp/p1_vec_cache.json"; KNN_CACHE="/tmp/p1_knn_cache.json"
TEST_COL="opp_p1_an_%d"%int(time.time())
async def _sql(sql,args=None):
async with DBPools().sqlorContext("pipeline") as sor:
r=await sor.sqlExe(sql,args or {}); await sor.sqlExe("COMMIT",{})
return r or []
def centroid(vecs):
n=len(vecs)
return [sum(v[k] for v in vecs)/n for k in range(DIM)]
def cos(a,b):
return sum(x*y for x,y in zip(a,b))
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[v]:
if sc>=tau: union(v,nid)
comps={}
for v in ids: comps.setdefault(find(v),[]).append(v)
return sorted(comps.values(),key=len,reverse=True)
async def main():
import aiohttp
timeout=aiohttp.ClientTimeout(total=60)
# ── 拉需求 ──
recs=await _sql("SELECT params_name,params_value FROM params WHERE params_name IN ('tender_api_base','tender_api_token')")
pm={r.params_name:r.params_value for r in recs}
base=(pm.get("tender_api_base") or "http://192.168.16.2:9085").rstrip("/"); token=pm.get("tender_api_token") or ""
hdr={"X-API-Token":token}
demands=[]; offset=0
async with aiohttp.ClientSession(timeout=timeout) as s:
while True:
async with s.get(base+"/api/demands",headers=hdr,params={"days":3650,"limit":200,"offset":offset}) as r:
d=await r.json()
items=d.get("items") or []; demands.extend(items); offset+=len(items)
if not items or offset>=d.get("matched",0): break
texts=[(it.get("title") or "").strip() or "(无标题)" for it in demands]
print("需求:",len(texts))
# ── 向量(缓存优先)──
if os.path.exists(VEC_CACHE):
c=json.load(open(VEC_CACHE)); vectors=c["vectors"]; assert len(vectors)==len(texts),"缓存与需求数不符,删缓存重跑"
print("向量: 命中缓存")
else:
from pipeline_service import rag_client as rc
pid=os.environ.get("P1_PROJECT_ID","")
if not pid: print("FATAL: 缺 P1_PROJECT_ID"); sys.exit(1)
vectors=[None]*len(texts)
for i in range(0,len(texts),100):
vs,err=None,None
for a in range(3):
vs,err=await rc.rag_embed_texts(pid,texts[i:i+100],batch_size=10)
if not err and vs: break
vs=None
if a==2: print("FATAL embed:",err); sys.exit(1)
await asyncio.sleep(2**a)
vectors[i:i+len(vs)]=vs
if i%500==0: print(" embed",i,flush=True)
json.dump({"vectors":vectors},open(VEC_CACHE,"w"))
print("向量: 计算完成并缓存")
ids=["d%d"%i for i in range(len(vectors))]
# ── kNN(缓存优先)──
if os.path.exists(KNN_CACHE):
nbr=json.load(open(KNN_CACHE)); nbr={k:[tuple(x) for x in v] for k,v in nbr.items()}
print("kNN: 命中缓存")
else:
recs=await _sql("SELECT baseurl FROM upapp WHERE id='rag-vdb'")
vdb=(recs[0].baseurl or "").rstrip("/")
async with aiohttp.ClientSession(timeout=timeout) as s:
async def vp(path,pl):
r=await s.post(vdb+path,json=pl); return await r.text()
await vp("/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":2000},
{"name":"batch_id","type":"str","max_length":64}],"description":"P1an","metric":"COSINE"})
for i in range(0,len(vectors),500):
rows=[{"id":ids[j],"vector":vectors[j],"text":texts[j][:1900],"batch_id":"p1an"} for j in range(i,min(i+500,len(vectors)))]
await vp("/v1/upsert",{"colname":TEST_COL,"data":rows})
nbr={}
for i,vid in enumerate(ids):
d=json.loads(await vp("/v1/query",{"colname":TEST_COL,"vector":vectors[i],"pagerows":TOPK,"output_fields":["id"]}))
rows=(d.get("data") or {}).get("rows") or []
nbr[vid]=[(float(r.get("score",0)),r.get("id")) for r in rows if r.get("id")!=vid]
if i%500==0: print(" kNN",i,flush=True)
await vp("/v1/dropcollection",{"colname":TEST_COL})
json.dump(nbr,open(KNN_CACHE,"w"))
print("kNN: 完成并缓存")
# ── 多τ + 大簇细分 ──
vec_by={ids[i]:vectors[i] for i in range(len(ids))}
txt_by={ids[i]:texts[i] for i in range(len(ids))}
result={"total":len(ids),"by_tau":[]}
for tau in TAUS:
comps=union_find(ids,nbr,tau)
# 大簇二次细分:簇内子图用 tau+0.10 重聚
final=[]
for c in comps:
if len(c)>BIG_SPLIT:
sub_nbr={v:[(sc,n) for sc,n in nbr[v] if n in set(c)] for v in c}
subs=union_find(c,sub_nbr,min(tau+0.10,0.95))
final.extend(subs)
else:
final.append(c)
big=[c for c in final if len(c)>=MIN_CLUSTER]
# 小组并入最近大簇
small=[c for c in final if len(c)<MIN_CLUSTER]
if big and small:
cents=[(c,centroid([vec_by[m] for m in c])) for c in big]
other=[]
for sc_ in small:
cc=centroid([vec_by[m] for m in sc_]); best,bs=None,-1
for bc,bcent in cents:
sm=cos(cc,bcent)
if sm>bs: best,bs=bc,sm
if best is not None and bs>=tau-0.10: best.extend(sc_)
else: other.extend(sc_)
else:
other=[m for c in small for m in c]
big.sort(key=len,reverse=True)
covered=sum(len(c) for c in big)
entry={"tau":tau,"n_clusters":len(big),"covered":covered,"other":len(other),
"cover_rate":round(covered/len(ids),3),"top":[]}
for ci,c in enumerate(big[:12]):
cent=centroid([vec_by[m] for m in c])
sims=sorted(((cos(vec_by[m],cent),m) for m in c),reverse=True)
# 去重样例(众包标题重复多)
seen=[];
for _,m in sims:
if txt_by[m] not in seen: seen.append(txt_by[m])
if len(seen)>=6: break
entry["top"].append({"rank":ci+1,"size":len(c),"share":round(len(c)/len(ids),3),"samples":seen})
result["by_tau"].append(entry)
print("\n═══ τ=%.2f 簇=%d 覆盖=%d(%.0f%%) 其他=%d ═══"%(tau,len(big),covered,covered/len(ids)*100,len(other)))
for cl in entry["top"][:10]:
print(" #%d %d条(%.1f%%): %s"%(cl["rank"],cl["size"],cl["share"]*100," | ".join(cl["samples"][:4])[:130]))
json.dump(result,open("/tmp/p1_analyze_result.json","w"),ensure_ascii=False,indent=1)
print("\n报告: /tmp/p1_analyze_result.json")
if __name__=="__main__":
t0=time.time(); asyncio.run(main()); print("TOTAL %.0fs"%(time.time()-t0))