pipeline-opportunity/scripts/p1_cluster_analyze.py

179 lines
8.1 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# -*- 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.jsonkNN图——存在则跳过 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))