179 lines
8.1 KiB
Python
179 lines
8.1 KiB
Python
# -*- 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))
|