amz_review_analyse/聚类.py
OnesvmWhoops 91e6c47fc0 迁移 DeepSeek Chat API,并支持本地 Qwen3 Embedding 向量化。
统一 voc_llm 密钥解析与默认模型;向量化改为本地 mlx 模型;更新 README、gitignore 与流水线文档。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-06-04 16:15:03 +08:00

1214 lines
40 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.

"""
从 voc_embeddings.sqlite 读取向量,按五段流程 UMAP+HDBSCAN 聚类,结果写入 voc_clustering.sqlite。
流程:
1. audience(LLM 自动调参 n_neighbors)
2a. 前两 audience 簇各自独立:簇内 pain_point(LLM 自动调参)
2b. 前两 audience 簇各自独立:簇内 aspect_opinion 按 Positive/Negative/Neutral 分桶聚类
3a. 全量 pain_point(LLM 自动调参)
3b. 全量 aspect_opinion 按 Positive/Negative/Neutral 分桶聚类(各自 LLM 自动调参)
规则:
- 聚类单元:每行 embedding_items
- n < 10:不跑 HDBSCAN,逐条独立簇标签(0..n-1)
- 步骤 1/2 的 -1:写入 SQL,不参与后续(步骤 1 的 -1 不进入步骤 2)
- 步骤 3 的 -1:写入 SQL,参与后续统计
用法::
./310py/bin/python 聚类.py
# 默认自动使用 voc_structured.sqlite 中最新 analysis_jobs.id(须已向量化)
./310py/bin/python 聚类.py --job-id 4 # 可选:手动指定
"""
from __future__ import annotations
import argparse
import csv
import json
import logging
import os
import random
import re
import sqlite3
import struct
import sys
import time
import warnings
from collections import defaultdict
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Dict, List, Optional, Sequence, Tuple
import hdbscan
import numpy as np
import umap
from openai import OpenAI
from sklearn.metrics import silhouette_score
from voc_llm import CHAT_MODEL, chat_extra_body, create_chat_client, require_chat_api_key
# UMAP 设 random_state 时的单线程提示;轮廓系数在 extmath 中的 matmul 数值告警
warnings.filterwarnings(
"ignore",
message=r"n_jobs value .* overridden .* by setting random_state",
category=UserWarning,
module=r"umap\.umap_",
)
warnings.filterwarnings(
"ignore",
category=RuntimeWarning,
module=r"sklearn\.utils\.extmath",
)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(message)s",
stream=sys.stderr,
)
logger = logging.getLogger("voc_cluster")
PROJECT_ROOT = Path(__file__).resolve().parent
STRUCTURED_DB = PROJECT_ROOT / "voc_structured.sqlite"
EMBED_DB = PROJECT_ROOT / "voc_embeddings.sqlite"
CLUSTER_DB = PROJECT_ROOT / "voc_clustering.sqlite"
LLM_MODEL = CHAT_MODEL
UMAP_N_COMPONENTS = 30
UMAP_MIN_DIST = 0.1
UMAP_METRIC = "cosine"
UMAP_RANDOM_STATE = 42
MIN_POINTS_FOR_HDBSCAN = 10
# HDBSCAN 目标:非离群簇数不超过该值(通过提高 min_cluster_size 达成)
MAX_CLUSTERS = 20
VALID_SENTIMENTS: Tuple[str, ...] = ("Positive", "Negative", "Neutral")
_SENTIMENT_STAGE_SUFFIX: Dict[str, str] = {
"Positive": "positive",
"Negative": "negative",
"Neutral": "neutral",
}
HDBSCAN_BASE = dict(
min_samples=1,
cluster_selection_method="eom",
prediction_data=True,
)
INITIAL_N_NEIGHBORS = 10
MAX_N_NEIGHBORS = 45
CROSS_SIMILAR_RATIO_THRESHOLD = 0.10
SILHOUETTE_STOP_THRESHOLD = 0.6
SILHOUETTE_DECLINE_WINDOW = 8 # 连续 8 个轮廓值:后 7 个均小于第 1 个则停止
SAMPLE_CAP = 30
SAMPLE_RATIO = 0.6
# 步骤 2 未启用 LLM 调参时的保守默认(当前步骤 2 已启用自动调参)
STEP2_N_NEIGHBORS = 8
# 簇内去重评论数 / 本步骤参与聚类的去重评论总数 < 该比例则不写入库、不进入报告
CLUSTER_MIN_REVIEW_RATIO = 0.10
@dataclass
class EmbedRow:
id: int
job_id: int
extraction_id: int
source_row: int
entity_type: str
entity_index: int
embed_text: str
audience: str | None
aspect: str | None
opinion: str | None
category: str | None
sentiment: str | None
content: str
dimensions: int
embedding: np.ndarray
def _unpack_embedding(blob: bytes, dimensions: int) -> np.ndarray:
n = dimensions
expected = n * 4
if len(blob) != expected:
raise ValueError(f"BLOB 长度 {len(blob)} != {expected}(dim={n})")
return np.array(struct.unpack(f"{n}f", blob), dtype=np.float32)
def _latest_job_id_from_structured(structured_db: Path) -> int | None:
"""与 向量化.py 一致:analysis_jobs 表 id 最大者为最新 job。"""
if not structured_db.is_file():
return None
conn = sqlite3.connect(structured_db)
try:
row = conn.execute(
"SELECT id FROM analysis_jobs ORDER BY id DESC LIMIT 1"
).fetchone()
return int(row[0]) if row else None
finally:
conn.close()
def _latest_job_id_from_embeddings(conn: sqlite3.Connection) -> int | None:
row = conn.execute(
"SELECT job_id FROM embedding_items ORDER BY job_id DESC LIMIT 1"
).fetchone()
return int(row[0]) if row else None
def _count_embeddings(conn: sqlite3.Connection, job_id: int) -> int:
row = conn.execute(
"SELECT COUNT(*) FROM embedding_items WHERE job_id = ?", (job_id,)
).fetchone()
return int(row[0]) if row else 0
def _resolve_job_id(
job_id: int | None,
embed_conn: sqlite3.Connection,
structured_db: Path = STRUCTURED_DB,
) -> int:
"""未传 job_id 时自动选最新结构化 job,并确认向量库中已有数据。"""
if job_id is not None:
jid = job_id
if _count_embeddings(embed_conn, jid) == 0:
raise RuntimeError(
f"job_id={jid} 在 {EMBED_DB.name} 中无向量,请先运行 向量化.py"
)
logger.info("使用指定 job_id=%s(%s 条向量)", jid, _count_embeddings(embed_conn, jid))
return jid
jid = _latest_job_id_from_structured(structured_db)
source = "voc_structured.analysis_jobs"
if jid is None:
jid = _latest_job_id_from_embeddings(embed_conn)
source = "voc_embeddings.embedding_items"
if jid is None:
raise RuntimeError("无可用 job_id:请先运行 结构化_server.py 与 向量化.py")
n = _count_embeddings(embed_conn, jid)
if n == 0:
latest_embed = _latest_job_id_from_embeddings(embed_conn)
if latest_embed is not None and latest_embed != jid:
logger.warning(
"结构化最新 job_id=%s 尚无向量,改用向量库最新 job_id=%s",
jid,
latest_embed,
)
jid = latest_embed
n = _count_embeddings(embed_conn, jid)
source = "voc_embeddings.embedding_items(回退)"
else:
raise RuntimeError(
f"最新 job_id={jid} 在向量库中无数据,请先对应该 job 运行 向量化.py"
)
logger.info("自动选用最新 job_id=%s(来源: %s,%s 条向量)", jid, source, n)
return jid
def _load_embed_rows(
conn: sqlite3.Connection, job_id: int, entity_type: str
) -> List[EmbedRow]:
cur = conn.execute(
"""
SELECT id, job_id, extraction_id, source_row, entity_type, entity_index,
embed_text, audience, aspect, opinion, category, sentiment,
content, dimensions, embedding
FROM embedding_items
WHERE job_id = ? AND entity_type = ?
ORDER BY source_row, entity_index, id
""",
(job_id, entity_type),
)
out: List[EmbedRow] = []
for r in cur.fetchall():
dims = int(r[13])
out.append(
EmbedRow(
id=int(r[0]),
job_id=int(r[1]),
extraction_id=int(r[2]),
source_row=int(r[3]),
entity_type=str(r[4]),
entity_index=int(r[5]),
embed_text=str(r[6]),
audience=r[7],
aspect=r[8],
opinion=r[9],
category=r[10],
sentiment=r[11],
content=str(r[12] or ""),
dimensions=dims,
embedding=_unpack_embedding(r[14], dims),
)
)
return out
def _filter_by_source_rows(
rows: List[EmbedRow], allowed: set[int]
) -> List[EmbedRow]:
return [r for r in rows if r.source_row in allowed]
def _strict_sentiment(raw: str | None) -> str | None:
"""仅接受 Positive / Negative / Neutral,其余丢弃。"""
if raw is None:
return None
s = str(raw).strip()
if s in VALID_SENTIMENTS:
return s
return None
def _filter_by_sentiment(rows: List[EmbedRow], sentiment: str) -> List[EmbedRow]:
return [r for r in rows if _strict_sentiment(r.sentiment) == sentiment]
def _stage_3b(sentiment: str) -> str:
return f"3b_aspect_opinion_{_SENTIMENT_STAGE_SUFFIX[sentiment]}"
def _stage_2b(sentiment: str, audience_cluster: int) -> str:
return f"2b_aspect_opinion_{_SENTIMENT_STAGE_SUFFIX[sentiment]}_audience_c{audience_cluster}"
def _cluster_aspect_opinion_sentiment_stages(
cconn: sqlite3.Connection,
run_id: int,
ao_rows: List[EmbedRow],
*,
stage_for: Any,
outlier_participates: bool,
parent_audience_cluster: int | None,
skipped_reason_prefix: str,
client: OpenAI,
) -> Dict[str, int]:
"""对 aspect_opinion 按三档情感各跑一 stage,返回 stage -> 条数。"""
counts: Dict[str, int] = {}
for sentiment in VALID_SENTIMENTS:
stage = stage_for(sentiment)
sub = _filter_by_sentiment(ao_rows, sentiment)
logger.info(
"[%s] %s 条 aspect_opinion(情感=%s)",
stage,
len(sub),
sentiment,
)
if not sub:
_save_stage_meta(
cconn,
run_id,
stage,
{
"skipped": True,
"reason": f"{skipped_reason_prefix}_{_SENTIMENT_STAGE_SUFFIX[sentiment]}",
"sentiment": sentiment,
},
)
counts[stage] = 0
continue
labels, meta = _cluster_stage(
sub, stage=stage, use_llm_tune=True, client=client
)
parent_clusters = None
if parent_audience_cluster is not None:
parent_clusters = [parent_audience_cluster] * len(sub)
filt = _save_assignments(
cconn,
run_id,
stage,
sub,
labels,
n_neighbors=meta.get("n_neighbors"),
outlier_participates=outlier_participates,
parent_clusters=parent_clusters,
)
meta["cluster_filter"] = filt
meta["sentiment"] = sentiment
_save_stage_meta(cconn, run_id, stage, meta)
if meta.get("tuning_log"):
_save_tuning_logs(cconn, run_id, stage, meta["tuning_log"])
counts[stage] = len(sub)
return counts
def _source_rows_for_audience_cluster(
row_to_aud: Dict[int, int], audience_cluster: int
) -> set[int]:
return {sr for sr, lab in row_to_aud.items() if lab == audience_cluster}
def _stage_name_step2_pain(audience_cluster: int) -> str:
return f"2a_pain_audience_c{audience_cluster}"
def _cluster_step2_for_audience(
cconn: sqlite3.Connection,
run_id: int,
*,
audience_cluster: int,
pain_rows: List[EmbedRow],
ao_rows: List[EmbedRow],
row_to_aud: Dict[int, int],
client: OpenAI,
) -> Dict[str, int]:
"""对单个 audience 簇分别聚类 pain 与 aspect_opinion,返回各 stage 条数。"""
allowed = _source_rows_for_audience_cluster(row_to_aud, audience_cluster)
counts: Dict[str, int] = {}
pain_sub = _filter_by_source_rows(pain_rows, allowed)
ao_sub = _filter_by_source_rows(ao_rows, allowed)
stage_pain = _stage_name_step2_pain(audience_cluster)
logger.info(
"[%s] audience 簇 %s:%s 条 pain(%s 条评论)",
stage_pain,
audience_cluster,
len(pain_sub),
len(allowed),
)
if pain_sub:
labels_pain, meta_pain = _cluster_stage(
pain_sub, stage=stage_pain, use_llm_tune=True, client=client
)
filt = _save_assignments(
cconn,
run_id,
stage_pain,
pain_sub,
labels_pain,
n_neighbors=meta_pain.get("n_neighbors"),
outlier_participates=False,
parent_clusters=[audience_cluster] * len(pain_sub),
)
meta_pain["cluster_filter"] = filt
_save_stage_meta(cconn, run_id, stage_pain, meta_pain)
if meta_pain.get("tuning_log"):
_save_tuning_logs(cconn, run_id, stage_pain, meta_pain["tuning_log"])
counts[stage_pain] = len(pain_sub)
else:
_save_stage_meta(
cconn,
run_id,
stage_pain,
{"skipped": True, "reason": "no_pain_points", "audience_cluster": audience_cluster},
)
counts[stage_pain] = 0
ao_counts = _cluster_aspect_opinion_sentiment_stages(
cconn,
run_id,
ao_sub,
stage_for=lambda s: _stage_2b(s, audience_cluster),
outlier_participates=False,
parent_audience_cluster=audience_cluster,
skipped_reason_prefix="no_aspect_opinion",
client=client,
)
counts.update(ao_counts)
return counts
def _adaptive_min_cluster_size(n: int) -> int:
"""初始 min_cluster_size = max(2, n // MAX_CLUSTERS),无上限;簇数仍 > MAX_CLUSTERS 时再迭代增大。"""
if n < MIN_POINTS_FOR_HDBSCAN:
return 2
return max(2, n // MAX_CLUSTERS)
def _n_clusters_from_labels(labels: np.ndarray) -> int:
labs = {int(x) for x in labels.tolist()}
labs.discard(-1)
return len(labs)
def _stack_embeddings(rows: Sequence[EmbedRow]) -> np.ndarray:
return np.vstack([r.embedding for r in rows]).astype(np.float32)
def _fit_umap(embeddings: np.ndarray, n_neighbors: int) -> np.ndarray:
n = len(embeddings)
n_neighbors = min(n_neighbors, max(2, n - 1))
reducer = umap.UMAP(
n_components=min(UMAP_N_COMPONENTS, max(2, n - 2)),
n_neighbors=n_neighbors,
min_dist=UMAP_MIN_DIST,
metric=UMAP_METRIC,
random_state=UMAP_RANDOM_STATE,
)
return reducer.fit_transform(embeddings)
def _fit_hdbscan(umap_emb: np.ndarray, min_cluster_size: int) -> np.ndarray:
clusterer = hdbscan.HDBSCAN(
min_cluster_size=min_cluster_size,
**HDBSCAN_BASE,
)
return clusterer.fit_predict(umap_emb)
def _run_umap_hdbscan_capped(
embeddings: np.ndarray,
n_neighbors: int,
*,
max_clusters: int = MAX_CLUSTERS,
) -> Tuple[np.ndarray, np.ndarray, int]:
"""UMAP + HDBSCAN;若簇数 > max_clusters 则增大 min_cluster_size 直至满足或无法再增。"""
n = len(embeddings)
umap_emb = _fit_umap(embeddings, n_neighbors)
min_cs = _adaptive_min_cluster_size(n)
labels = _fit_hdbscan(umap_emb, min_cs)
n_clusters = _n_clusters_from_labels(labels)
while n_clusters > max_clusters and min_cs < n:
step = max(1, (n_clusters - max_clusters + 1) // 2)
next_cs = min(min_cs + step, n)
if next_cs <= min_cs:
break
logger.info(
"簇数 %s > %s,min_cluster_size %s -> %s",
n_clusters,
max_clusters,
min_cs,
next_cs,
)
min_cs = next_cs
labels = _fit_hdbscan(umap_emb, min_cs)
n_clusters = _n_clusters_from_labels(labels)
if n_clusters > max_clusters:
logger.warning(
"簇数 %s 仍 > %s(min_cluster_size=%s, n=%s)",
n_clusters,
max_clusters,
min_cs,
n,
)
return labels, umap_emb, min_cs
def _labels_all_zero(n: int) -> np.ndarray:
return np.zeros(n, dtype=int)
def _strip_think(text: str) -> str:
if not text:
return text
text = re.sub(
r"<think>[\s\S]*?</think>", "", text, flags=re.IGNORECASE
)
text = re.sub(r"</?think>", "", text, flags=re.IGNORECASE)
return text.strip()
def _parse_json_from_llm(text: str) -> dict:
text = _strip_think(text or "")
text = text.replace("```json", "").replace("```JSON", "").replace("```", "").strip()
start = text.find("{")
end = text.rfind("}")
if start == -1 or end <= start:
raise ValueError("未找到 JSON 对象")
return json.loads(text[start : end + 1])
def _sample_from_clusters(
sentences: List[str], cluster_labels: np.ndarray, rng: random.Random
) -> Dict[int, List[str]]:
groups: Dict[int, List[str]] = {}
for sent, lab in zip(sentences, cluster_labels.tolist()):
groups.setdefault(int(lab), []).append(sent)
samples: Dict[int, List[str]] = {}
for lab, sents in groups.items():
if lab == -1:
continue
n = len(sents)
k = min(SAMPLE_CAP, max(1, int(n * SAMPLE_RATIO)))
k = min(k, n)
samples[lab] = rng.sample(sents, k) if n > k else list(sents)
return samples
def _ai_evaluate_cluster_samples(
client: OpenAI, cluster_samples: Dict[int, List[str]], max_retries: int = 2
) -> dict:
if not cluster_samples:
return {"cross_similar_count": 0, "total_sentences": 0}
lines = []
total = 0
for lab in sorted(cluster_samples.keys()):
lines.append(f"【聚类{lab}】")
for i, sent in enumerate(cluster_samples[lab], 1):
lines.append(f" {i}. {sent}")
total += 1
sample_text = "\n".join(lines)
prompt = f"""你是 VOC 评论短语聚类质量评估助手。以下是多个聚类类别的抽样短语。
{sample_text}
请统计 cross_similar_count:不同聚类类别之间、语义相似的短语条数(每句最多计 1)。
只输出 JSON:
{{
"cross_similar_count": 整数,
"total_sentences": {total}
}}"""
for attempt in range(max_retries):
try:
resp = client.chat.completions.create(
model=LLM_MODEL,
messages=[
{"role": "system", "content": "只输出合法 JSON。"},
{"role": "user", "content": prompt},
],
max_tokens=1500,
temperature=0.0,
response_format={"type": "json_object"},
# 关闭思考,避免 token 耗在 reasoning_content 导致 content 为空且无 JSON
extra_body=chat_extra_body(LLM_MODEL),
)
msg = resp.choices[0].message
raw = msg.content or getattr(msg, "reasoning_content", None) or ""
data = _parse_json_from_llm(raw)
data["total_sentences"] = int(data.get("total_sentences", total) or total)
data["cross_similar_count"] = int(data.get("cross_similar_count", 0))
return data
except Exception as e:
logger.warning("AI 评估失败(第%s次): %s", attempt + 1, e)
if attempt < max_retries - 1:
time.sleep(1)
return {"cross_similar_count": total, "total_sentences": total}
def _silhouette_decline_should_stop(scores: List[float]) -> bool:
"""最近 SILHOUETTE_DECLINE_WINDOW 个有效轮廓系数中,后 N-1 个是否都严格小于第一个。"""
if len(scores) < SILHOUETTE_DECLINE_WINDOW:
return False
window = scores[-SILHOUETTE_DECLINE_WINDOW:]
first = window[0]
return all(s < first for s in window[1:])
def _best_silhouette_in_window(
snapshots: List[Tuple[int, np.ndarray, float]],
) -> Tuple[int, np.ndarray, float] | None:
"""在最近轮廓窗口内取轮廓系数最高的一轮;无快照时返回 None。"""
if not snapshots:
return None
pool = snapshots[-SILHOUETTE_DECLINE_WINDOW:]
return max(pool, key=lambda x: x[2])
def _auto_tune_n_neighbors(
embeddings: np.ndarray,
sentences: List[str],
client: OpenAI,
stage: str,
) -> Tuple[np.ndarray, int, List[dict], Optional[float], int]:
rng = random.Random(UMAP_RANDOM_STATE)
n_neighbors = INITIAL_N_NEIGHBORS
tuning_log: List[dict] = []
labels: np.ndarray | None = None
min_cs = 2
silhouette_avg: float | None = None
silhouette_history: List[float] = []
# (n_neighbors, labels, silhouette) 仅在有有效轮廓时入栈,供早停回退最优轮次
silhouette_snapshots: List[Tuple[int, np.ndarray, float]] = []
n = len(embeddings)
while True:
logger.info("[%s] 调参轮次 n_neighbors=%s", stage, n_neighbors)
labels, umap_emb, min_cs = _run_umap_hdbscan_capped(embeddings, n_neighbors)
n_clusters = _n_clusters_from_labels(labels)
n_noise = int((labels == -1).sum())
silhouette_avg = None
if n_clusters > 1:
mask = labels != -1
if int(mask.sum()) > 1 and len(set(labels[mask].tolist())) > 1:
silhouette_avg = float(
silhouette_score(umap_emb[mask], labels[mask])
)
logger.info(
"[%s] 簇数=%s 离群=%s 轮廓=%s",
stage,
n_clusters,
n_noise,
f"{silhouette_avg:.4f}" if silhouette_avg is not None else "N/A",
)
samples = _sample_from_clusters(sentences, labels, rng)
total_sample = sum(len(v) for v in samples.values())
log_entry: dict = {
"stage": stage,
"n_neighbors": n_neighbors,
"min_cluster_size": min_cs,
"n_clusters": n_clusters,
"n_noise": n_noise,
"silhouette": silhouette_avg,
"total_sample": total_sample,
}
if silhouette_avg is not None:
silhouette_history.append(silhouette_avg)
silhouette_snapshots.append((n_neighbors, labels.copy(), silhouette_avg))
if _silhouette_decline_should_stop(silhouette_history):
window = silhouette_history[-SILHOUETTE_DECLINE_WINDOW:]
best = _best_silhouette_in_window(silhouette_snapshots)
assert best is not None
best_nn, best_labels, best_sil = best
n_neighbors = best_nn
labels = best_labels
silhouette_avg = best_sil
log_entry["silhouette_window"] = [round(s, 4) for s in window]
log_entry["stop_reason"] = (
f"连续{SILHOUETTE_DECLINE_WINDOW}轮轮廓:后5个均低于窗口首值"
f"{window[0]:.4f},回退至 n_neighbors={best_nn}(轮廓{best_sil:.4f})"
)
tuning_log.append(log_entry)
logger.info(
"[%s] 轮廓窗口 %s,早停并回退 n_neighbors=%s 轮廓=%.4f",
stage,
log_entry["silhouette_window"],
best_nn,
best_sil,
)
break
if n_neighbors > MAX_N_NEIGHBORS:
best = _best_silhouette_in_window(silhouette_snapshots)
if best is not None:
best_nn, best_labels, best_sil = best
n_neighbors = best_nn
labels = best_labels
silhouette_avg = best_sil
window = silhouette_history[-SILHOUETTE_DECLINE_WINDOW:]
log_entry["silhouette_window"] = [round(s, 4) for s in window]
log_entry["stop_reason"] = (
f"n_neighbors 超过上限 {MAX_N_NEIGHBORS},"
f"回退至 n_neighbors={best_nn}(轮廓{best_sil:.4f})"
)
logger.info(
"[%s] n_neighbors 超上限,回退 n_neighbors=%s 轮廓=%.4f",
stage,
best_nn,
best_sil,
)
else:
log_entry["stop_reason"] = (
f"n_neighbors 超过上限 {MAX_N_NEIGHBORS}(无有效轮廓快照,保留当前轮)"
)
tuning_log.append(log_entry)
break
ai_result = _ai_evaluate_cluster_samples(client, samples)
cross_count = ai_result["cross_similar_count"]
total = ai_result["total_sentences"] or total_sample or 1
ratio = cross_count / total
log_entry.update(
{
"cross_similar_count": cross_count,
"cross_similar_ratio": round(ratio, 4),
}
)
sil_ok = silhouette_avg is not None and silhouette_avg > SILHOUETTE_STOP_THRESHOLD
logger.info(
"[%s] AI: 跨类相似 %s/%s (%.1f%%,阈值≤%.0f%%) 轮廓=%s (停止需>%.1f)",
stage,
cross_count,
total,
ratio * 100,
CROSS_SIMILAR_RATIO_THRESHOLD * 100,
f"{silhouette_avg:.4f}" if silhouette_avg is not None else "N/A",
SILHOUETTE_STOP_THRESHOLD,
)
if ratio > CROSS_SIMILAR_RATIO_THRESHOLD:
log_entry["stop_reason"] = "跨类相似率>10%,继续调参"
tuning_log.append(log_entry)
n_neighbors += 1
time.sleep(0.3)
continue
if not sil_ok:
reason = (
"轮廓系数不可用,继续调参"
if silhouette_avg is None
else f"轮廓系数{silhouette_avg:.4f}≤{SILHOUETTE_STOP_THRESHOLD},继续调参"
)
log_entry["stop_reason"] = reason
tuning_log.append(log_entry)
n_neighbors += 1
time.sleep(0.3)
continue
log_entry["stop_reason"] = "跨类相似率≤10%且轮廓>0.6,停止调参"
tuning_log.append(log_entry)
break
assert labels is not None
return labels, n_neighbors, tuning_log, silhouette_avg, min_cs
def _cluster_fixed(
embeddings: np.ndarray, n_neighbors: int = STEP2_N_NEIGHBORS
) -> Tuple[np.ndarray, int, int]:
nn = min(n_neighbors, max(2, len(embeddings) - 1))
labels, _, min_cs = _run_umap_hdbscan_capped(embeddings, nn)
return labels, nn, min_cs
def _cluster_stage(
rows: List[EmbedRow],
*,
stage: str,
use_llm_tune: bool,
client: OpenAI | None,
) -> Tuple[np.ndarray, dict]:
n = len(rows)
meta: dict = {
"stage": stage,
"n_items": n,
"used_hdbscan": False,
"n_neighbors": None,
"tuning_log": [],
"silhouette": None,
}
if n == 0:
return np.array([], dtype=int), meta
if n < MIN_POINTS_FOR_HDBSCAN:
logger.info(
"[%s] n=%s < %s,逐条独立簇标签,不跑 HDBSCAN",
stage,
n,
MIN_POINTS_FOR_HDBSCAN,
)
meta["fallback"] = "per_item_labels_no_hdbscan"
return np.arange(n, dtype=int), meta
embeddings = _stack_embeddings(rows)
sentences = [r.embed_text for r in rows]
meta["used_hdbscan"] = True
if use_llm_tune:
if client is None:
raise RuntimeError(f"{stage} 需要 LLM 客户端")
labels, nn, tlog, sil, min_cs = _auto_tune_n_neighbors(
embeddings, sentences, client, stage
)
meta["n_neighbors"] = nn
meta["min_cluster_size"] = min_cs
meta["max_clusters"] = MAX_CLUSTERS
meta["tuning_log"] = tlog
meta["silhouette"] = sil
else:
labels, nn, min_cs = _cluster_fixed(embeddings)
meta["n_neighbors"] = nn
meta["min_cluster_size"] = min_cs
meta["max_clusters"] = MAX_CLUSTERS
meta["fallback"] = "fixed_params"
return labels, meta
def _init_cluster_db(conn: sqlite3.Connection) -> None:
conn.executescript(
"""
CREATE TABLE IF NOT EXISTS cluster_runs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
created_at TEXT NOT NULL,
job_id INTEGER NOT NULL,
embed_db TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS cluster_assignments (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
stage TEXT NOT NULL,
embedding_item_id INTEGER NOT NULL,
cluster_label INTEGER NOT NULL,
is_outlier INTEGER NOT NULL DEFAULT 0,
participates_downstream INTEGER NOT NULL DEFAULT 1,
parent_audience_cluster INTEGER,
n_neighbors INTEGER,
embed_text TEXT NOT NULL,
entity_type TEXT NOT NULL,
source_row INTEGER NOT NULL,
extraction_id INTEGER NOT NULL,
entity_index INTEGER NOT NULL,
audience TEXT,
aspect TEXT,
opinion TEXT,
category TEXT,
sentiment TEXT,
content TEXT NOT NULL,
FOREIGN KEY (run_id) REFERENCES cluster_runs(id)
);
CREATE TABLE IF NOT EXISTS cluster_stage_meta (
run_id INTEGER NOT NULL,
stage TEXT NOT NULL,
meta_json TEXT NOT NULL,
PRIMARY KEY (run_id, stage)
);
CREATE TABLE IF NOT EXISTS cluster_tuning_log (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id INTEGER NOT NULL,
stage TEXT NOT NULL,
round_index INTEGER NOT NULL,
log_json TEXT NOT NULL,
FOREIGN KEY (run_id) REFERENCES cluster_runs(id)
);
CREATE INDEX IF NOT EXISTS idx_assign_run_stage ON cluster_assignments(run_id, stage);
CREATE INDEX IF NOT EXISTS idx_assign_cluster ON cluster_assignments(run_id, stage, cluster_label);
CREATE INDEX IF NOT EXISTS idx_assign_source ON cluster_assignments(run_id, source_row);
"""
)
def _resolve_source_csv(structured_db: Path, job_id: int) -> Path:
conn = sqlite3.connect(structured_db)
try:
row = conn.execute(
"SELECT source_file FROM analysis_jobs WHERE id = ?", (job_id,)
).fetchone()
finally:
conn.close()
if not row:
raise RuntimeError(f"analysis_jobs 无 job_id={job_id}")
p = Path(str(row[0]))
if p.is_file():
return p.resolve()
cand = PROJECT_ROOT / p
if cand.is_file():
return cand.resolve()
raise FileNotFoundError(f"找不到结构化来源 CSV: {row[0]}")
def _count_csv_reviews(csv_path: Path) -> int:
n = 0
with csv_path.open(encoding="utf-8-sig", newline="") as f:
reader = csv.DictReader(f)
for row in reader:
if (row.get("content") or "").strip():
n += 1
return n
def _keep_cluster_labels(
labels: np.ndarray,
source_rows: Sequence[int],
total_reviews: int,
*,
min_ratio: float = CLUSTER_MIN_REVIEW_RATIO,
) -> Tuple[set[int], Dict[int, int]]:
"""返回 (保留的簇标签, 被丢弃簇标签 -> 去重评论数)。"""
by_lab: Dict[int, set[int]] = defaultdict(set)
for lab, sr in zip(labels.tolist(), source_rows):
by_lab[int(lab)].add(int(sr))
if total_reviews <= 0:
return set(by_lab.keys()), {}
kept: set[int] = set()
dropped: Dict[int, int] = {}
for lab, srs in by_lab.items():
cnt = len(srs)
if cnt / total_reviews >= min_ratio:
kept.add(lab)
else:
dropped[lab] = cnt
return kept, dropped
def _reset_cluster_db(path: Path, *, reset: bool = True) -> sqlite3.Connection:
if reset and path.is_file():
path.unlink()
conn = sqlite3.connect(path)
_init_cluster_db(conn)
return conn
def _save_assignments(
conn: sqlite3.Connection,
run_id: int,
stage: str,
rows: List[EmbedRow],
labels: np.ndarray,
*,
n_neighbors: int | None,
outlier_participates: bool,
parent_clusters: List[int | None] | None = None,
) -> dict:
if parent_clusters is not None and len(parent_clusters) != len(rows):
raise ValueError("parent_clusters 长度与 rows 不一致")
source_rows = [r.source_row for r in rows]
stage_total_reviews = len(set(source_rows))
kept, dropped = _keep_cluster_labels(labels, source_rows, stage_total_reviews)
if dropped:
logger.info(
"[%s] 过滤小簇(去重评论占比 < %.0f%%,本步骤评论总数 %s):%s",
stage,
CLUSTER_MIN_REVIEW_RATIO * 100,
stage_total_reviews,
{
lab: f"{cnt}条({cnt / stage_total_reviews * 100:.1f}%)"
for lab, cnt in sorted(dropped.items())
if stage_total_reviews > 0
},
)
sql = """
INSERT INTO cluster_assignments (
run_id, stage, embedding_item_id, cluster_label, is_outlier,
participates_downstream, parent_audience_cluster, n_neighbors,
embed_text, entity_type, source_row, extraction_id, entity_index,
audience, aspect, opinion, category, sentiment, content
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
"""
saved = 0
for i, row in enumerate(rows):
lab = int(labels[i])
if lab not in kept:
continue
is_out = 1 if lab == -1 else 0
if is_out:
participates = 1 if outlier_participates else 0
else:
participates = 1
parent = None
if parent_clusters is not None:
parent = parent_clusters[i]
conn.execute(
sql,
(
run_id,
stage,
row.id,
lab,
is_out,
participates,
parent,
n_neighbors,
row.embed_text,
row.entity_type,
row.source_row,
row.extraction_id,
row.entity_index,
row.audience,
row.aspect,
row.opinion,
row.category,
row.sentiment,
row.content,
),
)
saved += 1
return {
"kept_clusters": sorted(kept),
"dropped_clusters": dropped,
"saved_rows": saved,
"min_review_ratio": CLUSTER_MIN_REVIEW_RATIO,
"stage_total_reviews": stage_total_reviews,
}
def _save_stage_meta(conn: sqlite3.Connection, run_id: int, stage: str, meta: dict) -> None:
conn.execute(
"""
INSERT OR REPLACE INTO cluster_stage_meta (run_id, stage, meta_json)
VALUES (?, ?, ?)
""",
(run_id, stage, json.dumps(meta, ensure_ascii=False)),
)
def _save_tuning_logs(
conn: sqlite3.Connection, run_id: int, stage: str, tuning_log: List[dict]
) -> None:
for i, entry in enumerate(tuning_log, start=1):
conn.execute(
"""
INSERT INTO cluster_tuning_log (run_id, stage, round_index, log_json)
VALUES (?, ?, ?, ?)
""",
(run_id, stage, i, json.dumps(entry, ensure_ascii=False)),
)
def _top2_audience_clusters(
conn: sqlite3.Connection, run_id: int
) -> Tuple[List[int], Dict[int, int]]:
"""返回 (前两簇标签列表, source_row -> audience簇标签)。"""
cur = conn.execute(
"""
SELECT cluster_label, source_row
FROM cluster_assignments
WHERE run_id = ? AND stage = '1_audience'
AND is_outlier = 0 AND participates_downstream = 1
""",
(run_id,),
)
row_to_cluster: Dict[int, int] = {}
cluster_rows: Dict[int, set[int]] = {}
for lab, sr in cur.fetchall():
lab = int(lab)
sr = int(sr)
row_to_cluster[sr] = lab
cluster_rows.setdefault(lab, set()).add(sr)
ranked = sorted(
cluster_rows.items(), key=lambda x: len(x[1]), reverse=True
)
top2 = [lab for lab, _ in ranked[:2]]
return top2, row_to_cluster
def run_clustering(
*,
job_id: int | None = None,
structured_db: Path = STRUCTURED_DB,
embed_db: Path = EMBED_DB,
cluster_db: Path = CLUSTER_DB,
reset_db: bool = True,
) -> dict:
require_chat_api_key()
econn = sqlite3.connect(embed_db)
try:
jid = _resolve_job_id(job_id, econn, structured_db)
audience_rows = _load_embed_rows(econn, jid, "audience")
pain_rows = _load_embed_rows(econn, jid, "pain_point")
ao_rows = _load_embed_rows(econn, jid, "aspect_opinion")
finally:
econn.close()
client = create_chat_client()
cconn = _reset_cluster_db(cluster_db, reset=reset_db)
created = datetime.now(timezone.utc).isoformat()
try:
cur = cconn.execute(
"INSERT INTO cluster_runs (created_at, job_id, embed_db) VALUES (?,?,?)",
(created, jid, str(embed_db.resolve())),
)
run_id = int(cur.lastrowid)
csv_path = _resolve_source_csv(structured_db, jid)
total_reviews = _count_csv_reviews(csv_path)
logger.info("清洗后评论总数: %s", total_reviews)
# --- 1 audience ---
labels1, meta1 = _cluster_stage(
audience_rows, stage="1_audience", use_llm_tune=True, client=client
)
filt1 = _save_assignments(
cconn,
run_id,
"1_audience",
audience_rows,
labels1,
n_neighbors=meta1.get("n_neighbors"),
outlier_participates=False,
)
meta1["cluster_filter"] = filt1
_save_stage_meta(cconn, run_id, "1_audience", meta1)
if meta1.get("tuning_log"):
_save_tuning_logs(cconn, run_id, "1_audience", meta1["tuning_log"])
top2, row_to_aud = _top2_audience_clusters(cconn, run_id)
meta_top2 = {
"top2_audience_clusters": top2,
"per_cluster_source_rows": {
str(lab): len(_source_rows_for_audience_cluster(row_to_aud, lab))
for lab in top2
},
"step2_mode": "per_audience_cluster_llm_auto_tune",
}
_save_stage_meta(cconn, run_id, "step2_filter", meta_top2)
logger.info("Audience 前两簇(将分别聚类): %s", top2)
step2_counts: Dict[str, int] = {}
for aud_lab in top2:
step2_counts.update(
_cluster_step2_for_audience(
cconn,
run_id,
audience_cluster=aud_lab,
pain_rows=pain_rows,
ao_rows=ao_rows,
row_to_aud=row_to_aud,
client=client,
)
)
# --- 3a pain global ---
labels3a, meta3a = _cluster_stage(
pain_rows, stage="3a_pain_global", use_llm_tune=True, client=client
)
filt3a = _save_assignments(
cconn,
run_id,
"3a_pain_global",
pain_rows,
labels3a,
n_neighbors=meta3a.get("n_neighbors"),
outlier_participates=True,
)
meta3a["cluster_filter"] = filt3a
_save_stage_meta(cconn, run_id, "3a_pain_global", meta3a)
if meta3a.get("tuning_log"):
_save_tuning_logs(cconn, run_id, "3a_pain_global", meta3a["tuning_log"])
# --- 3b aspect_opinion by sentiment ---
step3b_counts = _cluster_aspect_opinion_sentiment_stages(
cconn,
run_id,
ao_rows,
stage_for=_stage_3b,
outlier_participates=True,
parent_audience_cluster=None,
skipped_reason_prefix="no_aspect_opinion",
client=client,
)
cconn.commit()
summary = {
"run_id": run_id,
"job_id": jid,
"total_reviews": total_reviews,
"min_cluster_review_ratio": CLUSTER_MIN_REVIEW_RATIO,
"cluster_db": str(cluster_db.resolve()),
"top2_audience_clusters": top2,
"counts": {
"1_audience": len(audience_rows),
**step2_counts,
"3a_pain_global": len(pain_rows),
**step3b_counts,
},
}
logger.info("聚类完成: %s", summary)
return summary
finally:
cconn.close()
def main() -> None:
parser = argparse.ArgumentParser(description="VOC 向量聚类")
parser.add_argument(
"--job-id",
type=int,
default=None,
help="不指定则自动使用 voc_structured.sqlite 中最新 analysis_jobs.id",
)
parser.add_argument("--structured-db", type=Path, default=STRUCTURED_DB)
parser.add_argument("--embed-db", type=Path, default=EMBED_DB)
parser.add_argument("--cluster-db", type=Path, default=CLUSTER_DB)
args = parser.parse_args()
summary = run_clustering(
job_id=args.job_id,
structured_db=args.structured_db,
embed_db=args.embed_db,
cluster_db=args.cluster_db,
)
print(json.dumps(summary, ensure_ascii=False, indent=2))
if __name__ == "__main__":
main()