# -*- 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))