diff --git a/scripts/p1_cluster_analyze.py b/scripts/p1_cluster_analyze.py new file mode 100644 index 0000000..933139d --- /dev/null +++ b/scripts/p1_cluster_analyze.py @@ -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)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))