diff --git a/.gitignore b/.gitignore index 9dbbd16..8dbeaa9 100644 --- a/.gitignore +++ b/.gitignore @@ -13,9 +13,13 @@ __pycache__/ # 密钥(勿提交) .dashscope_key +.deepseek_key .env .env.* +# 本地模型权重(体积过大,不纳入版本库) +Qwen3-Embedding-4B-mxfp8/ + # 本地数据库 voc_structured.sqlite voc_embeddings.sqlite @@ -30,4 +34,4 @@ output/ # 常见 --input-dir 原始数据目录(目录名因产品而异,按需追加) reviews_export/ *-voc/ -cat deterrent indoor / +cat deterrent indoor/ diff --git a/README.md b/README.md index 7a8e63b..3e5ff93 100644 --- a/README.md +++ b/README.md @@ -1,202 +1,257 @@ -# VOC LLM 结构化分析 (VOC_LLM结构化) - -> 基于大语言模型(通义千问 / DashScope)的亚马逊 VOC(Voice of Customer)评论分析流水线:合并 CSV → 清洗 → LLM 结构化 → 向量化 → 聚类与词频 → 生成 HTML 分析报告。 - -## 📖 目录 - -- [核心特性](#-核心特性) -- [环境要求](#-环境要求) -- [安装指南](#-安装指南) -- [使用说明](#-使用说明) -- [示例与输出](#-示例与输出) -- [项目结构](#-项目结构) -- [常见问题](#-常见问题) - - -## ✨ 核心特性 - -- **七步全流程编排** — `main_voc分析.py` 一键串联:合并、清洗、结构化、向量化、聚类、词频、HTML 报告 -- **断点续跑** — 支持 `--from-step` / `--only-step`,从任意步骤恢复,调试时节省 API 成本 -- **LLM 结构化提取** — 从评论中抽取受众、痛点、方面、观点、情感等字段(`prompts/schema.yaml` 可配置) -- **语义聚类** — UMAP + HDBSCAN 多阶段聚类,辅以 LLM 评估簇质量自动调参 -- **词频分析** — LLM 归纳专有名词 + spaCy 全量词频统计,报告内嵌词云与六类归类 -- **可编辑 Prompt** — `prompts/` 目录下 Markdown / YAML 热加载,产品运营可直接改话术(见 `prompts/README.md`) -- **并行加速** — 聚类与词频在步骤 5–6 由线程池并行执行 - -## 🛠 环境要求 - -| 依赖 | 说明 | -|------|------| -| Python | >= 3.10(推荐 3.10+) | -| pip / venv | 安装 `requirements.txt` 中的包 | -| spaCy 英文模型 | `python -m spacy download en_core_web_sm`(词频步骤必需) | -| 阿里云 DashScope API Key | 结构化、向量化、聚类评估、词频、报告等步骤均需调用 | - -**API Key 配置方式**(任选其一,勿提交到 Git): - -1. 环境变量 `DASHSCOPE_API_KEY` -2. 环境变量 `DASHSCOPE_API_KEY_FILE` 指向单行密钥文件 -3. 项目根目录 `.dashscope_key`(单行,无引号) - -可选:`DASHSCOPE_MODEL`(默认 `qwen3.6-flash`)。 - -## 📦 安装指南 - -1. 克隆项目到本地: - -```bash -git clone <你的仓库地址> -cd VOC_LLM结构化 -``` - -2. 创建虚拟环境并安装依赖(可选但推荐): - -```bash -python3 -m venv 310py -source 310py/bin/activate # Windows: 310py\Scripts\activate -pip install -r requirements.txt -python -m spacy download en_core_web_sm -``` - -3. 配置 API Key: - -```bash -export DASHSCOPE_API_KEY="sk-xxx" -# 或在项目根创建 .dashscope_key(已被 .gitignore 忽略) -``` - -4. 准备原始评论 CSV 目录(目录内所有 `*.csv` 表头须一致),例如亚马逊导出的 `*_realtime.csv`。 - -## 🚀 使用说明 - -### 全流程分析 - -```bash -python3 main_voc分析.py \ - --input-dir "reviews_export" \ - --product "cat deterrent indoor" \ - --industry "Pet Supplies" -``` - -- `--input-dir`:原始 CSV 目录 -- `--product`:产品名(写入结构化任务与报告路径) -- `--industry`:行业名,默认 `Pet Supplies` -- `--keep-db`:保留已有结构化 `voc_*.sqlite`,不覆盖删除 - -### 断点续跑 - -```bash -# 从向量化起续跑(步骤 4 起可省略 --product,自动读结构化库) -python3 main_voc分析.py --from-step 4 --keep-db - -# 仅重跑词频(复用已有 voc_terms.json) -python3 main_voc分析.py --from-step 6 --skip-wordfreq-llm - -# 仅重新生成 HTML 报告 -python3 main_voc分析.py --only-step 7 -``` - -### 其他常用参数 - -| 参数 | 说明 | -|------|------| -| `--clean-intermediates` | 报告成功后删除中间 csv/sqlite,减少内存占用 | -| `--filter-small-clusters` | 报告仅保留簇内评论占比 ≥ 10% 的簇,仅当评论数量过万时启用 | -| `--save-llm-raw` | 将报告 LLM 原文保存为 `report_llm_raw.txt`,调试时使用 | - -### 程序式调用 - -```python -from pathlib import Path -from main_voc分析 import run_voc_analysis - -result = run_voc_analysis( - input_dir=Path("reviews_export"), - industry="Pet Supplies", - product_name="cat deterrent indoor", - from_step=1, - clean_databases=True, -) -print(result["report_html"]) -``` - -### Prompt 验收(无需 API Key) - -```bash -python3 prompts/smoke.py # 检查 prompt 能否加载 -python3 prompts/smoke.py --live # 联调模型(需 API Key) -``` - -### 可选变体:jieba 词频 - -中文或需 jieba 分词时,可使用 `main_voc分析_jieba.py`(词频走 `词频_jieba.py`,其余步骤与主流程一致)。 - ---- - -更详细的步骤说明、算法与 SQLite 约定见 **[main_voc分析.md](main_voc分析.md)**。 - -## 📸 示例与输出 - -流程结束后,主要产物如下: - -| 路径 | 说明 | -|------|------| -| `merged_reviews.csv` | 多文件合并结果 | -| `merged_reviews_cleaned.csv` | 清洗、去重后的评论 | -| `voc_structured.sqlite` | LLM 结构化结果 | -| `voc_embeddings.sqlite` | 256 维向量 | -| `voc_clustering.sqlite` | 多阶段聚类标签 | -| `output/voc_terms.json` | 专有名词 / 停用词 | -| `output/word_freq.csv` | 全量词频表 | -| `output/{product}/{product}_voc_report.html` | **最终 VOC 分析报告**(词云、词频、分簇、AI 正文) | - -stdout 会打印 JSON 摘要(含 `report_html` 等键)。 - -> 建议在 README 或文档中补充一张 `*_voc_report.html` 在浏览器中打开的截图,便于新成员快速理解交付物形态。 - -## 📂 项目结构 - -```text -VOC_LLM结构化/ -├── main_voc分析.py # 主流程编排入口(七步) -├── main_voc分析_jieba.py # 词频使用 jieba 的变体入口 -├── main_voc分析.md # 流程与算法详细说明 -├── 合并评论数据.py # 步骤 1:多 CSV 合并 -├── content清洗.py # 步骤 2:评论清洗与去重 -├── 结构化_server.py # 步骤 3:LLM 结构化入库 -├── 结构化_Prompt.py # 结构化 prompt 组装 -├── 向量化.py # 步骤 4:Embedding 入库 -├── 聚类.py # 步骤 5:UMAP + HDBSCAN -├── 词频.py / 词频_jieba.py # 步骤 6:术语提取 + 词频 -├── voc_report.py # 步骤 7:HTML 报告生成 -├── prompts/ # 可编辑 prompt、schema、配置 -│ ├── README.md -│ ├── schema.yaml -│ ├── extraction/ report/ word_freq/ -│ └── loader.py -├── requirements.txt -├── output/ # 报告与词频输出(gitignore) -└── README.md # 本文件 -``` - -## ❓ 常见问题 - -**Q:提示缺少 `DASHSCOPE_API_KEY`?** -A:按上文配置环境变量或 `.dashscope_key`,并确认密钥未提交到仓库。 - -**Q:步骤 4 报错找不到 `product`?** -A:从步骤 1–3 完整跑过,或确保 `voc_structured.sqlite` 中已有最新 job;步骤 4 起可省略 `--product`。 - -**Q:词频步骤报 spaCy 模型缺失?** -A:执行 `python -m spacy download en_core_web_sm`。 - -**Q:合并 CSV 失败?** -A:确保 `--input-dir` 下所有 CSV 表头完全一致。 - -**Q:修改 LLM 话术后报告解析失败?** -A:勿修改 `prompts/schema.yaml` 中 `report.markers` 四段标记名;改完运行 `python3 prompts/smoke.py` 验收。 - - ---- - -*README 与 `main_voc分析.py` 七步流程保持一致;深度说明请参阅 `main_voc分析.md`。* +# VOC LLM 结构化分析 (VOC_LLM结构化) + +> 基于大语言模型([DeepSeek](https://api.deepseek.com) OpenAI 兼容 API)与本地 MLX 向量的亚马逊 VOC(Voice of Customer)评论分析流水线:合并 CSV → 清洗 → LLM 结构化 → 向量化 → 聚类与词频 → 生成 HTML 分析报告。 + +## 📖 目录 + +- [核心特性](#-核心特性) +- [环境要求](#-环境要求) +- [安装指南](#-安装指南) +- [使用说明](#-使用说明) +- [示例与输出](#-示例与输出) +- [项目结构](#-项目结构) +- [常见问题](#-常见问题) +- [参与贡献](#-参与贡献) +- [开源协议](#-开源协议) +- [联系方式与鸣谢](#-联系方式与鸣谢) + +## ✨ 核心特性 + +- **七步全流程编排** — `main_voc分析.py` 一键串联:合并、清洗、结构化、向量化、聚类、词频、HTML 报告 +- **断点续跑** — 支持 `--from-step` / `--only-step`,从任意步骤恢复,调试时节省 API 成本 +- **LLM 结构化提取** — 从评论中抽取受众、痛点、方面、观点、情感等字段(`prompts/schema.yaml` 可配置) +- **本地向量化** — Apple Silicon 上运行 `Qwen3-Embedding-4B-mxfp8`(MLX),无需云端 Embedding API +- **语义聚类** — UMAP + HDBSCAN 多阶段聚类,辅以 LLM 评估簇质量自动调参 +- **词频分析** — LLM 归纳专有名词 + spaCy 全量词频统计,报告内嵌词云与六类归类 +- **可编辑 Prompt** — `prompts/` 目录下 Markdown / YAML 热加载,产品运营可直接改话术(见 `prompts/README.md`) +- **并行加速** — 聚类与词频在步骤 5–6 由线程池并行执行;结构化批间并行(默认 8 路) + +## 🛠 环境要求 + +| 依赖 | 说明 | +|------|------| +| Python | >= 3.10(推荐 3.12,项目内 `310py`) | +| pip / uv | 安装 `requirements.txt` 中的包 | +| spaCy 英文模型 | 经 `uv pip` 安装 `en-core-web-sm`(见安装指南,词频步骤必需) | +| DeepSeek API Key | 结构化、聚类评估、词频、报告等 Chat 步骤 | +| 本地 Embedding 模型 | 目录 `Qwen3-Embedding-4B-mxfp8/`(约 4GB,已 gitignore,需自行下载) | +| Apple Silicon | 本地向量化依赖 MLX(M 系列芯片) | + +**Chat API Key**(任选其一,勿提交到 Git): + +1. 环境变量 `DEEPSEEK_API_KEY` +2. 环境变量 `DEEPSEEK_API_KEY_FILE` 指向单行密钥文件 +3. 项目根目录 `.deepseek_key`(单行,无引号) + +可选:`DEEPSEEK_MODEL`(默认 `deepseek-v4-pro`)、`DEEPSEEK_BASE_URL`(默认 `https://api.deepseek.com`)。 + +## 📦 安装指南 + +1. 克隆项目到本地: + +```bash +git clone https://git.onesvm.com/whoops/amz_review_analyse.git +cd amz_review_analyse # 或你的本地目录名 +``` + +2. 创建虚拟环境并安装依赖(推荐): + +```bash +uv venv 310py --python 3.12 +uv pip install --python 310py/bin/python -r requirements.txt +uv pip install --python 310py/bin/python \ + "en-core-web-sm @ https://github.com/explosion/spacy-models/releases/download/en_core_web_sm-3.8.0/en_core_web_sm-3.8.0-py3-none-any.whl" +``` + +3. 配置 DeepSeek API Key: + +```bash +export DEEPSEEK_API_KEY="sk-xxx" +# 或在项目根创建 .deepseek_key(已被 .gitignore 忽略) +``` + +4. 准备本地 Embedding 模型(首次向量化前): + +将 `Qwen3-Embedding-4B-mxfp8` 放到项目根,或设置 `VOC_EMBED_MODEL_PATH` 指向模型目录。可从 [Hugging Face](https://huggingface.co/mlx-community/Qwen3-Embedding-4B-mxfp8) 下载。 + +5. 准备原始评论 CSV 目录(目录内所有 `*.csv` 表头须一致),例如亚马逊导出的 `*_realtime.csv`。 + +## 🚀 使用说明 + +### 全流程分析 + +```bash +./310py/bin/python main_voc分析.py \ + --input-dir "reviews_export" \ + --product "cat deterrent indoor" \ + --industry "Pet Supplies" +``` + +- `--input-dir`:原始 CSV 目录 +- `--product`:产品名(写入结构化任务与报告路径) +- `--industry`:行业名,默认 `-`(可在步骤 3 写入库) +- `--keep-db`:保留已有结构化 `voc_*.sqlite`,不覆盖删除 + +### 加速(批间并行,默认已开启) + +步骤 3 结构化默认多批并行 Chat 请求;步骤 4 向量为本地 MLX 串行批处理(勿对同一模型多线程): + +```bash +# 全流程 +./310py/bin/python main_voc分析.py --input-dir "reviews_export" --product "产品名" + +# 调低结构化并发(遇 429 时) +./310py/bin/python main_voc分析.py --input-dir "reviews_export" --product "产品名" \ + --struct-workers 4 + +# 环境变量:export VOC_STRUCT_WORKERS=8 VOC_EMBED_BATCH_SIZE=16 +``` + +### 断点续跑 + +```bash +# 从向量化起续跑(步骤 4 起可省略 --product,自动读结构化库) +./310py/bin/python main_voc分析.py --from-step 4 --keep-db + +# 仅重跑词频(复用已有 voc_terms.json) +./310py/bin/python main_voc分析.py --from-step 6 --skip-wordfreq-llm + +# 仅重新生成 HTML 报告 +./310py/bin/python main_voc分析.py --only-step 7 +``` + +### 其他常用参数 + +| 参数 | 说明 | +|------|------| +| `--clean-intermediates` | 报告成功后删除中间 csv/sqlite,减少磁盘占用 | +| `--filter-small-clusters` | 报告仅保留簇内评论占比 ≥ 10% 的簇 | +| `--save-llm-raw` | 将报告 LLM 原文保存为 `report_llm_raw.txt`,调试时使用 | + +### 程序式调用 + +```python +from pathlib import Path +from main_voc分析 import run_voc_analysis + +result = run_voc_analysis( + input_dir=Path("reviews_export"), + industry="Pet Supplies", + product_name="cat deterrent indoor", + from_step=1, + clean_databases=True, +) +print(result["report_html"]) +``` + +### Prompt 验收(无需 API Key) + +```bash +./310py/bin/python prompts/smoke.py # 检查 prompt 能否加载 +./310py/bin/python prompts/smoke.py --live # 联调模型(需 DEEPSEEK_API_KEY) +``` + +### 可选变体:jieba 词频 + +中文或需 jieba 分词时,可使用 `main_voc分析_jieba.py`(词频走 `词频_jieba.py`,其余步骤与主流程一致)。 + +--- + +更详细的步骤说明、算法与 SQLite 约定见 **[main_voc分析.md](main_voc分析.md)**。 + +## 📸 示例与输出 + +流程结束后,主要产物如下: + +| 路径 | 说明 | +|------|------| +| `merged_reviews.csv` | 多文件合并结果 | +| `merged_reviews_cleaned.csv` | 清洗、去重后的评论 | +| `voc_structured.sqlite` | LLM 结构化结果 | +| `voc_embeddings.sqlite` | 本地 Qwen3 向量(维度见库内 `dimensions` 字段) | +| `voc_clustering.sqlite` | 多阶段聚类标签 | +| `output/voc_terms.json` | 专有名词 / 停用词 | +| `output/word_freq.csv` | 全量词频表 | +| `output/{product}/{product}_voc_report.html` | **最终 VOC 分析报告**(词云、词频、分簇、AI 正文) | + +stdout 会打印 JSON 摘要(含 `report_html` 等键)。 + +## 📂 项目结构 + +```text +VOC_LLM结构化/ +├── main_voc分析.py # 主流程编排入口(七步) +├── main_voc分析_jieba.py # 词频使用 jieba 的变体入口 +├── main_voc分析.md # 流程与算法详细说明 +├── voc_llm.py # DeepSeek Chat 密钥与客户端 +├── local_embedding.py # 本地 MLX Qwen3 向量化 +├── 合并评论数据.py # 步骤 1:多 CSV 合并 +├── content清洗.py # 步骤 2:评论清洗与去重 +├── 结构化_server.py # 步骤 3:LLM 结构化入库 +├── 结构化_Prompt.py # 结构化 prompt 组装 +├── 向量化.py # 步骤 4:本地 Embedding 入库 +├── 聚类.py # 步骤 5:UMAP + HDBSCAN +├── 词频.py / 词频_jieba.py # 步骤 6:术语提取 + 词频 +├── voc_report.py # 步骤 7:HTML 报告生成 +├── prompts/ # 可编辑 prompt、schema、配置 +├── Qwen3-Embedding-4B-mxfp8/ # 本地模型(gitignore,需自行放置) +├── requirements.txt +├── output/ # 报告与词频输出(gitignore) +└── README.md # 本文件 +``` + +## ❓ 常见问题 + +**Q:提示缺少 `DEEPSEEK_API_KEY`?** +A:按上文配置环境变量或 `.deepseek_key`,并确认密钥未提交到仓库。 + +**Q:只有 `.dashscope_key` 报错?** +A:Chat 已切换为 DeepSeek,DashScope 密钥不能用于 `api.deepseek.com`,请改用 `.deepseek_key`。 + +**Q:步骤 4 向量化失败 / 找不到模型?** +A:确认 `Qwen3-Embedding-4B-mxfp8/` 在项目根,或设置 `VOC_EMBED_MODEL_PATH`;需在 Apple Silicon + Python 3.10+ 环境。 + +**Q:步骤 4 报错找不到 `product`?** +A:从步骤 1–3 完整跑过,或确保 `voc_structured.sqlite` 中已有最新 job;步骤 4 起可省略 `--product`。 + +**Q:词频步骤报 spaCy 模型缺失?** +A:执行 `uv pip install --python 310py/bin/python "en-core-web-sm @ https://github.com/explosion/spacy-models/releases/download/en_core_web_sm-3.8.0/en_core_web_sm-3.8.0-py3-none-any.whl"`。 + +**Q:合并 CSV 失败?** +A:确保 `--input-dir` 下所有 CSV 表头完全一致。 + +**Q:修改 LLM 话术后报告解析失败?** +A:勿修改 `prompts/schema.yaml` 中 `report.markers` 四段标记名;改完运行 `./310py/bin/python prompts/smoke.py` 验收。 + +## 🤝 参与贡献 + +欢迎提交 Issue 与 Pull Request。建议流程: + +1. Fork 本仓库 +2. 创建特性分支:`git checkout -b feature/your-feature` +3. 提交更改:`git commit -m '简要说明变更'` +4. 推送并发起 Pull Request + +修改主流程或 CLI 时,请同步更新 `main_voc分析.md`;修改 `prompts/` 时请遵循 `prompts/README.md` 中的占位符与 schema 约定。 + +**安全提醒**:勿提交 `.deepseek_key`、`.dashscope_key`、`.env`、真实评论 CSV、`*.sqlite`、`output/` 及本地模型目录(见 `.gitignore`)。 + +## 📄 开源协议 + +本项目尚未在仓库中附带 `LICENSE` 文件。若为内部项目,请按组织规范使用;若计划开源,请补充协议文件(如 MIT)并更新本节链接。 + +## ✉️ 联系方式与鸣谢 + +- **详细技术文档**:[main_voc分析.md](main_voc分析.md)、[prompts/README.md](prompts/README.md) +- **项目仓库**:https://git.onesvm.com/whoops/amz_review_analyse + +### 鸣谢 + +- [DeepSeek API](https://api.deepseek.com) — Chat 结构化、聚类评估、词频与报告 +- [mlx-community/Qwen3-Embedding-4B-mxfp8](https://huggingface.co/mlx-community/Qwen3-Embedding-4B-mxfp8) — 本地向量化 +- [UMAP](https://umap-learn.readthedocs.io/)、[HDBSCAN](https://hdbscan.readthedocs.io/) — 聚类管线 +- [spaCy](https://spacy.io/) — 英文词频与 NLP + +--- + +*README 与 `main_voc分析.py` 七步流程保持一致;深度说明请参阅 `main_voc分析.md`。* diff --git a/local_embedding.py b/local_embedding.py new file mode 100644 index 0000000..d1610ff --- /dev/null +++ b/local_embedding.py @@ -0,0 +1,144 @@ +""" +本地 MLX Qwen3 Embedding(Apple Silicon)。 + +默认模型目录:项目根 ``Qwen3-Embedding-4B-mxfp8``,可用 ``VOC_EMBED_MODEL_PATH`` 覆盖。 + +在 M4 / 16GB、约 8GB 可用内存下实测(mxfp8,500 字符/条): + - 加载峰值约 1.5GB + - batch 1–16 稳定;默认 ``VOC_EMBED_BATCH_SIZE=16``,串行推理(勿多线程并行加载同一 MLX 模型) + +需 Python 3.10+ 与 ``mlx-embeddings>=0.1.0``(推荐项目 ``310py`` 虚拟环境)。 +""" +from __future__ import annotations + +import logging +import os +import sys +import types +from pathlib import Path +from typing import List, Sequence + +logger = logging.getLogger("local_embedding") + +PROJECT_ROOT = Path(__file__).resolve().parent +DEFAULT_MODEL_PATH = PROJECT_ROOT / "Qwen3-Embedding-4B-mxfp8" +# 实测短句 batch=32 仍 <0.5GB增量;长句 500 字符 batch=16 约 1.5GB 峰值 +DEFAULT_BATCH_SIZE = 16 +DEFAULT_MAX_TEXT_CHARS = 512 + + +def _apply_hf_hub_shim() -> None: + if "huggingface_hub.utils._errors" in sys.modules: + return + try: + from huggingface_hub.errors import RepositoryNotFoundError + except ImportError: + try: + from huggingface_hub.utils._errors import RepositoryNotFoundError # type: ignore + except ImportError: + RepositoryNotFoundError = Exception # type: ignore[misc, assignment] + mod = types.ModuleType("huggingface_hub.utils._errors") + mod.RepositoryNotFoundError = RepositoryNotFoundError + sys.modules["huggingface_hub.utils._errors"] = mod + + +def _resolve_model_path() -> Path: + raw = os.environ.get("VOC_EMBED_MODEL_PATH", "").strip() + p = Path(raw).expanduser() if raw else DEFAULT_MODEL_PATH + if not p.is_dir(): + raise FileNotFoundError(f"本地 embedding 模型目录不存在: {p}") + return p.resolve() + + +def _resolve_batch_size(explicit: int | None = None) -> int: + if explicit is not None and explicit > 0: + return explicit + env = os.environ.get("VOC_EMBED_BATCH_SIZE", "").strip() + if env.isdigit() and int(env) > 0: + return int(env) + return DEFAULT_BATCH_SIZE + + +def _resolve_max_chars() -> int: + env = os.environ.get("VOC_EMBED_MAX_TEXT_CHARS", "").strip() + if env.isdigit() and int(env) > 0: + return int(env) + return DEFAULT_MAX_TEXT_CHARS + + +def _truncate(text: str, max_chars: int) -> str: + t = (text or "").strip() + if len(t) <= max_chars: + return t + return t[: max_chars - 3] + "..." + + +_MODEL = None +_PROCESSOR = None +_DIMENSIONS: int | None = None + + +def _load(): + global _MODEL, _PROCESSOR + if _MODEL is not None: + return _MODEL, _PROCESSOR + _apply_hf_hub_shim() + from mlx_embeddings import load + + path = _resolve_model_path() + logger.info("加载本地 embedding: %s", path) + _MODEL, _PROCESSOR = load(str(path)) + return _MODEL, _PROCESSOR + + +def embedding_dimensions() -> int: + global _DIMENSIONS + if _DIMENSIONS is not None: + return _DIMENSIONS + vecs = embed_texts(["dimension probe"]) + _DIMENSIONS = len(vecs[0]) + return _DIMENSIONS + + +def embed_texts( + texts: Sequence[str], + *, + batch_size: int | None = None, + max_chars: int | None = None, +) -> List[List[float]]: + """返回与输入等长的浮点向量列表(已 L2 归一化)。""" + if not texts: + return [] + from mlx_embeddings import generate + + model, processor = _load() + bs = _resolve_batch_size(batch_size) + cap = max_chars if max_chars is not None else _resolve_max_chars() + cleaned = [_truncate(t, cap) for t in texts] + + out_all: List[List[float]] = [] + n = len(cleaned) + n_chunks = (n + bs - 1) // bs + log_every = max(1, n_chunks // 20) + + for chunk_idx, i in enumerate(range(0, n, bs)): + chunk = cleaned[i : i + bs] + output = generate(model, processor, texts=chunk) + emb = output.text_embeds + for row in emb: + out_all.append([float(x) for x in row.tolist()]) + + done = chunk_idx + 1 + if done == 1 or done == n_chunks or done % log_every == 0: + rows_done = min(done * bs, n) + logger.info( + "向量化进度 %s/%s 批(%s/%s 条)", + done, + n_chunks, + rows_done, + n, + ) + global _DIMENSIONS + if _DIMENSIONS is None and out_all: + _DIMENSIONS = len(out_all[0]) + return out_all diff --git a/main_voc分析.md b/main_voc分析.md index 014e121..58307e5 100644 --- a/main_voc分析.md +++ b/main_voc分析.md @@ -23,8 +23,8 @@ flowchart LR |------|------|------------------------| | 1 | `合并评论数据.py` | 否 | | 2 | `content清洗.py` | 否 | -| 3 | `结构化_server.py` | **是**(Chat 结构化) | -| 4 | `向量化.py` | **是**(Embedding) | +| 3 | `结构化_server.py` | **是**(Chat 结构化;默认 8 路批间并行) | +| 4 | `向量化.py` | **是**(Embedding;默认 8 路批间并行) | | 5–6 | `聚类.py` ∥ `词频.py` | **是**(聚类调参评估 + 词频术语提取) | | 7 | `voc_report.py` | **是**(报告撰写、词频分类、翻译等) | @@ -50,6 +50,8 @@ flowchart LR | `--skip-wordfreq-llm` | 否 | 关闭 | 词频复用已有 `output/voc_terms.json`,跳过 LLM 术语提取 | | `--filter-small-clusters` | 否 | 关闭 | 报告仅纳入簇内去重评论占比 ≥ 10% 的簇 | | `--save-llm-raw` | 否 | 关闭 | 将报告 LLM 完整原文写入 `{product}_voc_report.html` 同目录 `report_llm_raw.txt` | +| `--struct-workers` | 否 | `8`(`VOC_STRUCT_WORKERS`) | 步骤 3 结构化**批间并行** Chat 请求数 | +| `--embed-workers` | 否 | `8`(`VOC_EMBED_WORKERS`) | 步骤 4 向量化**批间并行** Embedding 请求数 | ### 2.2 输出 @@ -88,16 +90,16 @@ flowchart LR cd "/Users/onesvmwhoops/Cursor_Project/VOC_LLM结构化" # 全流程(默认写入 sqlite 前会清理旧库;加 --keep-db 则保留) -python3 main_voc分析.py --input-dir "某目录" --product "产品名" +./310py/bin/python main_voc分析.py --input-dir "某目录" --product "产品名" # 从步骤 4 续跑(product 可省略) -python3 main_voc分析.py --from-step 4 --keep-db +./310py/bin/python main_voc分析.py --from-step 4 --keep-db # 仅重跑词频 LLM 之前的 spaCy 统计 -python3 main_voc分析.py --from-step 6 --skip-wordfreq-llm +./310py/bin/python main_voc分析.py --from-step 6 --skip-wordfreq-llm # 仅生成报告 -python3 main_voc分析.py --only-step 7 +./310py/bin/python main_voc分析.py --only-step 7 ``` ### 2.4 程序式调用 @@ -138,17 +140,18 @@ result = run_voc_analysis( ### 步骤 3:结构化(`结构化_server.py` + `结构化_Prompt.py` + `prompts/`) - **输入**:清洗 CSV、`industry`、`product_name`。 -- **方法**:DashScope **Chat Completions**(默认 `qwen3.6-flash`),按 token 估算批量调用,从评论中提取 JSON 字段(audience、pain_points、aspect、opinion、category、sentiment 等,以 `prompts/schema.yaml` 为准)。 +- **方法**:DeepSeek **Chat Completions**(默认 `deepseek-v4-pro`,思考关闭),按**输入 token** 与**条数**动态分批(默认单批估算输入 ≤ 200K、最多 **200** 条/批,`VOC_STRUCT_BATCH_MAX_REVIEWS` 可覆盖);从评论中提取 JSON 字段(以 `prompts/schema.yaml` 为准)。 - **输出**:`voc_structured.sqlite`(`analysis_jobs`、`comment_extractions`)。 -### 步骤 4:向量化(`向量化.py`) +### 步骤 4:向量化(`向量化.py` + `local_embedding.py`) - **输入**:最新或指定 `job_id` 的结构化实体;`embed_text` 由 audience / pain_point / aspect / opinion / aspect_opinion 等展开。 - **方法**: - - API:`text-embedding-v4`,**256 维**,余弦相似度空间中的稠密向量; + - 本地 MLX:`Qwen3-Embedding-4B-mxfp8`(默认 **2560 维**,L2 归一化); - 存储:`float32` 打包为 BLOB(`struct.pack`); - - 批大小 ≤ 10/请求,默认 6 线程并行多批。 + - 默认批大小 **16**、串行推理(`VOC_EMBED_BATCH_SIZE`);M4 16GB 实测峰值约 1.5GB。 - **输出**:`voc_embeddings.sqlite`(`embedding_items`)。 +- **运行**:推荐 `./310py/bin/python`(Python 3.10+)。 ### 步骤 5:聚类(`聚类.py`) @@ -198,28 +201,34 @@ result = run_voc_analysis( --- -## 4. API Key 使用说明 +## 4. API Key 与模型 -`main_voc分析.py` **本身不读取、不持有 API Key**;密钥解析在各子模块内统一实现,优先级一致: +### Chat(步骤 3 / 5 / 6 / 7) -1. 环境变量 `DASHSCOPE_API_KEY` -2. 环境变量 `DASHSCOPE_API_KEY_FILE` 指向的单行密钥文件 -3. 项目根文件 `.dashscope_key`(单行,无引号) +统一经 `voc_llm.py`,默认 **`deepseek-v4-pro`**(`chat_extra_body` 关闭思考),`https://api.deepseek.com`。步骤 7 主报告在 `voc_report.py` 单独使用 Pro + `reasoning_effort=max`。 -可选环境变量:`DASHSCOPE_MODEL`(默认各模块为 `qwen3.6-flash`)。 +密钥优先级: -| 模块 | 使用 API Key 的位置 | API 类型 / 用途 | -|------|---------------------|-----------------| -| `结构化_server.py` | `_resolve_dashscope_api_key()` → `OpenAI(...)` | Chat:批量/单条评论结构化 | -| `向量化.py` | `_resolve_api_key()` → `_embed_one_api_batch` | Embeddings:`text-embedding-v4` | -| `聚类.py` | `_resolve_api_key()` → `run_clustering` 内 `OpenAI` | Chat:簇间相似度评估、调 `n_neighbors` | -| `词频.py` | `_resolve_api_key()` → `_step1_extract_terms` / `_call_llm` | Chat:专有名词与停用词提取 | -| `voc_report.py` | `_resolve_api_key()` → `generate_report` 及子函数 | Chat:报告生成、词频分类、翻译等 | -| `prompts/smoke.py` | `--live` 时 | 冒烟测试(非 main 流程) | +1. `DEEPSEEK_API_KEY` +2. `DEEPSEEK_API_KEY_FILE` +3. 项目根 `.deepseek_key` + +可选:`DEEPSEEK_MODEL`、`DEEPSEEK_BASE_URL`。 + +| 模块 | 用途 | +|------|------| +| `结构化_server.py` | 评论结构化 | +| `聚类.py` | 簇质量评估与调参 | +| `词频.py` | 专有名词 / 停用词 | +| `voc_report.py` | 报告、词频分类、翻译 | + +### Embedding(步骤 4) + +**无需 API Key**。本地目录 `Qwen3-Embedding-4B-mxfp8`(或 `VOC_EMBED_MODEL_PATH`)。 **仓库安全规范**(见 `.gitignore`): -- **禁止提交** `.dashscope_key`、`voc_structured.sqlite` 及含真实评论/密钥的敏感导出; +- **禁止提交** `.deepseek_key`、`.dashscope_key`、`voc_structured.sqlite` 及含真实评论/密钥的敏感导出; - 密钥仅通过环境变量或本机未跟踪文件提供; - 文档与代码中勿写入真实 `sk-` 密钥。 @@ -276,10 +285,12 @@ result = run_voc_analysis( ```bash cd "/Users/onesvmwhoops/Cursor_Project/VOC_LLM结构化" -python3 -m venv 310py && source 310py/bin/activate # 可选,与 .gitignore 一致 -pip install -r requirements.txt -python -m spacy download en_core_web_sm # 词频步骤需要 -export DASHSCOPE_API_KEY="sk-xxx" # 或配置 .dashscope_key(勿提交) +uv venv 310py --python 3.12 # 与 .gitignore 中 310py/ 一致 +uv pip install --python 310py/bin/python -r requirements.txt +# uv 虚拟环境无 pip,勿用「python -m spacy download」;直接装模型 wheel: +uv pip install --python 310py/bin/python "en-core-web-sm @ https://github.com/explosion/spacy-models/releases/download/en_core_web_sm-3.8.0/en_core_web_sm-3.8.0-py3-none-any.whl" +export DEEPSEEK_API_KEY="sk-xxx" # 或配置 .deepseek_key(勿提交) +# 向量化推荐:./310py/bin/python main_voc分析.py ... ``` `requirements.txt` 中与数学/ NLP 相关的主要包:`numpy`、`umap-learn`、`hdbscan`、`scikit-learn`、`spacy`、`openai`、`pyyaml`。 @@ -290,7 +301,7 @@ export DASHSCOPE_API_KEY="sk-xxx" # 或配置 .dashscope_key( - 修改 LLM 话术:编辑 `prompts/` 下对应 `.md`,**勿改** `schema.yaml` 中 `report.markers` 四段标记名(见 `prompts/README.md`)。 - 修改主流程步骤顺序或默认路径:改 `main_voc分析.py` 后请同步更新**本文档**。 -- 验收 prompt 加载:`python3 prompts/smoke.py`(无需 Key);联调模型:`python3 prompts/smoke.py --live`。 +- 验收 prompt 加载:`./310py/bin/python prompts/smoke.py`(无需 Key);联调模型:`./310py/bin/python prompts/smoke.py --live`。 --- diff --git a/main_voc分析.py b/main_voc分析.py index cbeedf3..d7bb86b 100644 --- a/main_voc分析.py +++ b/main_voc分析.py @@ -4,17 +4,19 @@ VOC 全流程:合并 → 清洗 → 结构化 → 向量化 →(聚类 ∥ 用法:: # 全流程(--industry 默认 Pet Supplies;写入 sqlite 前默认清理旧库,keep-db不清理) - python3 main_voc分析.py --input-dir 目录 --product "产品名" --keep-db - + ./310py/bin/python main_voc分析.py --input-dir 目录 --product "产品名" --keep-db +--industry "行业名" # 断点续跑(步骤 4 起可省略 --product,自动读 voc_structured.sqlite) - python3 main_voc分析.py --from-step 5 - python3 main_voc分析.py --from-step 6 --skip-wordfreq-llm # 仅重跑词频 - python3 main_voc分析.py --only-step 7 # 仅生成报告 + ./310py/bin/python main_voc分析.py --from-step 5 + ./310py/bin/python main_voc分析.py --from-step 6 --skip-wordfreq-llm # 仅重跑词频 + ./310py/bin/python main_voc分析.py --only-step 7 # 仅生成报告 # 保留已有 voc_structured / voc_embeddings / voc_clustering.sqlite,不清理覆盖 - python3 main_voc分析.py --from-step 4 --keep-db + ./310py/bin/python main_voc分析.py --from-step 4 --keep-db 报告 HTML:output/{product_name}/{product_name}_voc_report.html + +Chat 默认 deepseek-v4-pro、思考关闭(DEEPSEEK_API_KEY);报告主 LLM 见 voc_report(Pro + max);向量化本地 Qwen3-Embedding-4B-mxfp8(推荐 ./310py/bin/python)。 """ from __future__ import annotations @@ -29,9 +31,10 @@ from pathlib import Path from content清洗 import process_reviews, save_cleaned_reviews from 合并评论数据 import merge_csv_directory -from 向量化 import run_embed -from 结构化_server import run_analysis +from 向量化 import EMBED_DEFAULT_WORKERS, run_embed +from 结构化_server import STRUCT_DEFAULT_WORKERS, run_analysis from 聚类 import run_clustering +from voc_llm import require_chat_api_key from voc_report import DEFAULT_CLUSTER_MIN_REVIEW_RATIO, generate_report from 词频 import run as run_wordfreq @@ -53,7 +56,7 @@ WORD_FREQ_CSV = OUTPUT_DIR / "word_freq.csv" DEFAULT_MERGED = PROJECT_ROOT / "merged_reviews.csv" DEFAULT_CLEANED = PROJECT_ROOT / "merged_reviews_cleaned.csv" -DEFAULT_INDUSTRY = "Pet Supplies" +DEFAULT_INDUSTRY = "-" def _safe_product_dir_name(product_name: str) -> str: @@ -124,6 +127,8 @@ def run_voc_analysis( skip_wordfreq_llm: bool = False, min_cluster_review_ratio: float | None = None, save_llm_raw: bool = False, + struct_workers: int | None = None, + embed_workers: int | None = None, ) -> dict: result: dict = {} @@ -140,6 +145,15 @@ def run_voc_analysis( save_cleaned_reviews(df, cleaned_csv) result["cleaned_csv"] = str(cleaned_csv) + need_chat = ( + _should_run(3, from_step, only_step) + or _should_run(5, from_step, only_step) + or _should_run(6, from_step, only_step) + or _should_run(7, from_step, only_step) + ) + if need_chat: + require_chat_api_key() + if _should_run(3, from_step, only_step): logger.info("步骤 3/7:结构化分析") if not cleaned_csv.is_file(): @@ -149,9 +163,13 @@ def run_voc_analysis( product_name=product_name, file_path=str(cleaned_csv), clean_databases=clean_databases, + workers=struct_workers, ) result["structured_job_id"] = ar.get("job_id") result["structured_db"] = str(STRUCTURED_DB) + n_ext = len(ar.get("extractions") or {}) + if n_ext == 0: + raise RuntimeError("步骤 3 结构化无有效结果,已中止(不会继续向量化)") job_id: int | None = None ind, prod = industry, product_name @@ -175,6 +193,7 @@ def run_voc_analysis( structured_db=STRUCTURED_DB, embed_db=EMBED_DB, reset_db=clean_databases, + workers=embed_workers, ) result["embed"] = er job_id = int(er.get("job_id", job_id or 0)) @@ -352,6 +371,26 @@ def main() -> None: action="store_true", help="步骤 7 将报告 LLM 完整原文写入报告同目录下的 report_llm_raw.txt(默认不保存)", ) + parser.add_argument( + "--struct-workers", + type=int, + default=None, + metavar="N", + help=( + f"步骤 3 结构化批间并行数(默认 {STRUCT_DEFAULT_WORKERS};" + "环境变量 VOC_STRUCT_WORKERS)" + ), + ) + parser.add_argument( + "--embed-workers", + type=int, + default=None, + metavar="N", + help=( + f"步骤 4 向量化批间并行数(默认 {EMBED_DEFAULT_WORKERS};" + "环境变量 VOC_EMBED_WORKERS)" + ), + ) args = parser.parse_args() need_input = args.only_step in (None, 1) and args.from_step <= 1 @@ -385,6 +424,8 @@ def main() -> None: else None ), save_llm_raw=args.save_llm_raw, + struct_workers=args.struct_workers, + embed_workers=args.embed_workers, ) print(json.dumps(out, ensure_ascii=False, indent=2)) diff --git a/main_voc分析_jieba.py b/main_voc分析_jieba.py index 7085aa2..6b56743 100644 --- a/main_voc分析_jieba.py +++ b/main_voc分析_jieba.py @@ -3,10 +3,10 @@ VOC 全流程(jieba 分词版):与 main_voc分析.py 相同,步骤 6 词 用法:: - python3 main_voc分析_jieba.py --input-dir reviews_export --industry "Pet supplements" --product "Turkey tail mushroom for dogs" + ./310py/bin/python main_voc分析_jieba.py --input-dir reviews_export --industry "Pet supplements" --product "Turkey tail mushroom for dogs" - python3 main_voc分析_jieba.py --from-step 6 # 仅重跑 jieba 词频(可省略 --industry/--product) - python3 main_voc分析_jieba.py --only-step 7 # 仅生成报告 + ./310py/bin/python main_voc分析_jieba.py --from-step 6 # 仅重跑 jieba 词频(可省略 --industry/--product) + ./310py/bin/python main_voc分析_jieba.py --only-step 7 # 仅生成报告 """ from __future__ import annotations @@ -20,8 +20,8 @@ from pathlib import Path from content清洗 import process_reviews, save_cleaned_reviews from 合并评论数据 import merge_csv_directory -from 向量化 import run_embed -from 结构化_server import run_analysis +from 向量化 import EMBED_DEFAULT_WORKERS, run_embed +from 结构化_server import STRUCT_DEFAULT_WORKERS, run_analysis from 聚类 import run_clustering from voc_report import DEFAULT_CLUSTER_MIN_REVIEW_RATIO, generate_report from 词频_jieba import run as run_wordfreq @@ -98,6 +98,8 @@ def run_voc_analysis( skip_wordfreq_llm: bool = False, min_cluster_review_ratio: float | None = None, save_llm_raw: bool = False, + struct_workers: int | None = None, + embed_workers: int | None = None, ) -> dict: result: dict = {"report_html": str(REPORT_HTML), "tokenizer": "jieba"} @@ -122,6 +124,7 @@ def run_voc_analysis( industry=industry, product_name=product_name, file_path=str(cleaned_csv), + workers=struct_workers, ) result["structured_job_id"] = ar.get("job_id") result["structured_db"] = str(STRUCTURED_DB) @@ -141,6 +144,7 @@ def run_voc_analysis( csv_path=cleaned_csv, structured_db=STRUCTURED_DB, embed_db=EMBED_DB, + workers=embed_workers, ) result["embed"] = er job_id = int(er.get("job_id", job_id or 0)) @@ -302,6 +306,20 @@ def main() -> None: action="store_true", help="步骤 7 将报告 LLM 完整原文写入 output/report_llm_raw.txt(默认不保存)", ) + parser.add_argument( + "--struct-workers", + type=int, + default=None, + metavar="N", + help=f"步骤 3 结构化批间并行数(默认 {STRUCT_DEFAULT_WORKERS})", + ) + parser.add_argument( + "--embed-workers", + type=int, + default=None, + metavar="N", + help=f"步骤 4 向量化批间并行数(默认 {EMBED_DEFAULT_WORKERS})", + ) args = parser.parse_args() need_input = args.only_step in (None, 1) and args.from_step <= 1 @@ -334,6 +352,8 @@ def main() -> None: else None ), save_llm_raw=args.save_llm_raw, + struct_workers=args.struct_workers, + embed_workers=args.embed_workers, ) print(json.dumps(out, ensure_ascii=False, indent=2)) diff --git a/prompts/README.md b/prompts/README.md index 4f34ab9..b5676ea 100644 --- a/prompts/README.md +++ b/prompts/README.md @@ -22,8 +22,8 @@ ```bash cd "/Users/onesvmwhoops/Cursor_Project/VOC_LLM结构化" -python3 prompts/smoke.py # 检查文件能否加载(无需 API Key) -python3 prompts/smoke.py --live # 用 3 条样例评论调模型(需 DASHSCOPE_API_KEY) +./310py/bin/python prompts/smoke.py # 检查文件能否加载(无需 API Key) +./310py/bin/python prompts/smoke.py --live # 用 3 条样例评论调模型(需 DEEPSEEK_API_KEY) ``` ## 注意 diff --git a/prompts/extraction/field_rules.md b/prompts/extraction/field_rules.md index 32e4423..2fac9a7 100644 --- a/prompts/extraction/field_rules.md +++ b/prompts/extraction/field_rules.md @@ -12,7 +12,10 @@ - 将用户对产品的反馈拆解为具体对象与反馈内容,包含以下子字段: - aspect (对象): 提炼成准确的英文简短名词。**务必保留具体的成分、材质或特定属性词**(例如 "chicken flavor" 不能泛化为 "flavor")。 - opinion (反馈内容): 必须是英文简短词组。 - - sentiment (情感): 仅限 {sentiments_literal}。 + - sentiment (情感): 仅限 {sentiments_literal}(大小写不限,输出须为 Positive / Negative / Neutral 三者之一)。 + - **禁止** Mixed、Ambiguous、Both、Balanced 等自创词。 + - 同一条反馈褒贬交织时:选**最主要**倾向;或拆成多条 product_feedback 分别标注。 + - 无法判断倾向时用 Neutral。 - category (类别): 仅限 {categories_literal}(禁止 Value、Cost 等自创词;性价比高/物有所值 归入 Price)。 - 每条 product_feedback 必须同时包含 aspect、opinion、sentiment、category 四个子字段。 - 若无提及 product_feedback,输出空列表 [] diff --git a/prompts/smoke.py b/prompts/smoke.py index 56d82ce..1202a3a 100644 --- a/prompts/smoke.py +++ b/prompts/smoke.py @@ -2,8 +2,8 @@ """ Prompt 验收脚本(修改 prompts/ 后运行)。 - python3 prompts/smoke.py # 仅校验文件加载与渲染(无需 API Key) - python3 prompts/smoke.py --live # 调用模型跑 smoke/reviews.yaml(需 DASHSCOPE_API_KEY) + ./310py/bin/python prompts/smoke.py # 仅校验文件加载与渲染(无需 API Key) + ./310py/bin/python prompts/smoke.py --live # 调用模型跑 smoke/reviews.yaml(需 DEEPSEEK_API_KEY) """ from __future__ import annotations diff --git a/requirements.txt b/requirements.txt index 431d4da..b9fd654 100644 --- a/requirements.txt +++ b/requirements.txt @@ -6,3 +6,6 @@ scikit-learn>=1.3.0 spacy>=3.7.0 pyyaml>=6.0 jieba>=0.42.1 +# 本地 embedding(Python 3.10+,Apple Silicon) +mlx>=0.22.0 +mlx-embeddings>=0.1.0 diff --git a/voc_llm.py b/voc_llm.py new file mode 100644 index 0000000..c76c505 --- /dev/null +++ b/voc_llm.py @@ -0,0 +1,74 @@ +""" +VOC 流水线 Chat 模型配置(DeepSeek OpenAI 兼容 API)。 + +密钥(任选其一):: + export DEEPSEEK_API_KEY="sk-..." + export DEEPSEEK_API_KEY_FILE="/path/to/key.txt" + 项目根单行文件 .deepseek_key + +模型:: + 默认 deepseek-v4-pro;可用 DEEPSEEK_MODEL 覆盖。 + 流水线默认关闭思考(extra_body thinking disabled);报告等单独开启处见 voc_report。 +""" +from __future__ import annotations + +import os +from pathlib import Path +from typing import Any, Dict + +from openai import OpenAI + +PROJECT_ROOT = Path(__file__).resolve().parent + +CHAT_BASE_URL = os.environ.get("DEEPSEEK_BASE_URL", "https://api.deepseek.com").strip() +CHAT_MODEL = os.environ.get("DEEPSEEK_MODEL", "deepseek-v4-pro").strip() + + +def resolve_chat_api_key() -> str: + """环境变量 > DEEPSEEK_API_KEY_FILE > 项目根 .deepseek_key。""" + v = os.environ.get("DEEPSEEK_API_KEY", "").strip() + if v: + return v + fp = os.environ.get("DEEPSEEK_API_KEY_FILE", "").strip() + if fp: + p = Path(fp).expanduser() + if p.is_file(): + return p.read_text(encoding="utf-8").strip().strip('"').strip("'") + local = PROJECT_ROOT / ".deepseek_key" + if local.is_file(): + return local.read_text(encoding="utf-8").strip().strip('"').strip("'") + # 兼容:仅有 .dashscope_key 时提示(DashScope 密钥不能用于 api.deepseek.com) + legacy = PROJECT_ROOT / ".dashscope_key" + if legacy.is_file(): + raise RuntimeError( + "检测到 .dashscope_key,但 Chat 已切换为 DeepSeek。" + "请在项目根创建 .deepseek_key(单行 DEEPSEEK 密钥)," + "或执行 export DEEPSEEK_API_KEY='sk-...'" + ) + return "" + + +def require_chat_api_key() -> str: + key = resolve_chat_api_key() + if not key: + raise RuntimeError( + "缺少 DeepSeek API Key:设置 DEEPSEEK_API_KEY," + "或 DEEPSEEK_API_KEY_FILE,或在项目根创建 .deepseek_key(单行)" + ) + return key + + +def create_chat_client(*, api_key: str | None = None, timeout: float | None = None) -> OpenAI: + kw: Dict[str, Any] = { + "api_key": api_key or require_chat_api_key(), + "base_url": CHAT_BASE_URL, + } + if timeout is not None: + kw["timeout"] = timeout + return OpenAI(**kw) + + +def chat_extra_body(model: str | None = None) -> Dict[str, Any]: + """流水线 Chat 默认关闭思考(Pro 模型 API 默认 otherwise 为 enabled)。""" + _ = (model or CHAT_MODEL).lower() + return {"thinking": {"type": "disabled"}} diff --git a/voc_report.py b/voc_report.py index 84caa38..8553060 100644 --- a/voc_report.py +++ b/voc_report.py @@ -2,12 +2,12 @@ 基于聚类结果与词频 CSV 生成 output/voc_report.html(词云 + 词频 + AI 报告 + 各簇表述)。 用法:: - python3 main_voc分析.py - python3 main_voc分析.py --only-step 7 --industry "..." --product "..." + ./310py/bin/python main_voc分析.py + ./310py/bin/python main_voc分析.py --only-step 7 --industry "..." --product "..." # 仅纳入本 stage 内评论占比 ≥10% 的簇: - python3 main_voc分析.py --only-step 7 --filter-small-clusters ... + ./310py/bin/python main_voc分析.py --only-step 7 --filter-small-clusters ... # 调试:保存报告 LLM 完整原文到 output/report_llm_raw.txt - python3 main_voc分析.py --only-step 7 --save-llm-raw --industry "..." --product "..." + ./310py/bin/python main_voc分析.py --only-step 7 --save-llm-raw --industry "..." --product "..." """ from __future__ import annotations @@ -24,8 +24,6 @@ from dataclasses import dataclass, field from pathlib import Path from typing import Any, Dict, List, Sequence, Set, Tuple -from openai import OpenAI - from prompts.loader import ( build_category_analysis_prompts, build_report_analysis_requirements, @@ -36,14 +34,19 @@ from prompts.loader import ( build_word_assign_prompts, sync_voc_report_constants, ) +from voc_llm import CHAT_MODEL, chat_extra_body, create_chat_client, require_chat_api_key logger = logging.getLogger("voc_report") PROJECT_ROOT = Path(__file__).resolve().parent STRUCTURED_DB = PROJECT_ROOT / "voc_structured.sqlite" STRUCTURED_APPENDIX_SAMPLE_N = 15 -DASHSCOPE_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" -MODEL_NAME = os.environ.get("DASHSCOPE_MODEL", "qwen3.6-flash").strip() +MODEL_NAME = CHAT_MODEL +# 步骤 7 主报告(HTML + JSON 标记块);其它报告内 LLM 仍用 MODEL_NAME +REPORT_MODEL = os.environ.get("DEEPSEEK_REPORT_MODEL", "deepseek-v4-pro").strip() +REPORT_REASONING_EFFORT = os.environ.get( + "DEEPSEEK_REPORT_REASONING_EFFORT", "max" +).strip() or "max" WORDCLOUD_TOP_N = 180 WORD_FREQ_TABLE_N = 180 # 页面词频表、排名展示上限 @@ -150,21 +153,6 @@ class ClusterBundle: embed_texts_zh: List[str] = field(default_factory=list) -def _resolve_api_key() -> str: - v = os.environ.get("DASHSCOPE_API_KEY", "").strip() - if v: - return v - fp = os.environ.get("DASHSCOPE_API_KEY_FILE", "").strip() - if fp: - p = Path(fp).expanduser() - if p.is_file(): - return p.read_text(encoding="utf-8").strip().strip('"').strip("'") - local = PROJECT_ROOT / ".dashscope_key" - if local.is_file(): - return local.read_text(encoding="utf-8").strip().strip('"').strip("'") - return "" - - def _strip_think(text: str) -> str: if not text: return text @@ -179,30 +167,66 @@ def _call_llm_messages( messages: Sequence[Dict[str, str]], api_key: str, *, + model: str | None = None, temperature: float = 0.3, max_tokens: int = 16384, timeout: float = 300.0, + reasoning_effort: str | None = None, + extra_body: Dict[str, Any] | None = None, + report_thinking: bool = False, ) -> str: - client = OpenAI( - api_key=api_key, base_url=DASHSCOPE_BASE_URL, timeout=timeout - ) - extra_body: Dict[str, Any] = {} - if MODEL_NAME.lower().startswith(("qwen3.6", "qwen3.5", "qwen3")): - extra_body["enable_thinking"] = False - resp = client.chat.completions.create( - model=MODEL_NAME, - messages=list(messages), - temperature=temperature, - max_tokens=max_tokens, - **({"extra_body": extra_body} if extra_body else {}), - ) + _ = api_key + client = create_chat_client(timeout=timeout) + use_model = model or MODEL_NAME + create_kw: Dict[str, Any] = { + "model": use_model, + "messages": list(messages), + "max_tokens": max_tokens, + } + if report_thinking: + create_kw["reasoning_effort"] = reasoning_effort or REPORT_REASONING_EFFORT + create_kw["extra_body"] = extra_body or {"thinking": {"type": "enabled"}} + else: + create_kw["temperature"] = temperature + eb = extra_body if extra_body is not None else chat_extra_body(use_model) + if eb: + create_kw["extra_body"] = eb + + resp = client.chat.completions.create(**create_kw) msg = resp.choices[0].message - text = msg.content or getattr(msg, "reasoning_content", None) or "" + if report_thinking: + text = msg.content or "" + if not text.strip(): + fr = getattr(resp.choices[0], "finish_reason", None) + raise RuntimeError( + f"报告 LLM content 为空(model={use_model},finish_reason={fr!r});" + "思考模式下请检查 max_tokens 是否截断" + ) + else: + text = msg.content or getattr(msg, "reasoning_content", None) or "" if not text.strip(): raise RuntimeError("LLM 返回为空") return _strip_think(text) +def _call_report_llm_messages( + messages: Sequence[Dict[str, str]], + api_key: str, + *, + max_tokens: int = REPORT_MAX_OUTPUT_TOKENS, + timeout: float = REPORT_LLM_TIMEOUT_SEC, +) -> str: + """主分析报告:deepseek-v4-pro + 思考模式 max(仅此处启用)。""" + return _call_llm_messages( + messages, + api_key, + model=REPORT_MODEL, + max_tokens=max_tokens, + timeout=timeout, + report_thinking=True, + ) + + def _call_llm( system: str, user: str, @@ -1440,11 +1464,9 @@ def _fetch_report_llm_raw_with_retry( {"role": "system", "content": system}, {"role": "user", "content": user}, ] - raw = _call_llm_messages( + raw = _call_report_llm_messages( messages, api_key, - max_tokens=REPORT_MAX_OUTPUT_TOKENS, - timeout=REPORT_LLM_TIMEOUT_SEC, ) for attempt in range(max_retries + 1): errors = _validate_report_response(raw, expect_cluster_names=expect_cluster_names) @@ -1472,11 +1494,9 @@ def _fetch_report_llm_raw_with_retry( "content": _build_report_correction_user_message(errors, raw), } ) - raw = _call_llm_messages( + raw = _call_report_llm_messages( messages, api_key, - max_tokens=REPORT_MAX_OUTPUT_TOKENS, - timeout=REPORT_LLM_TIMEOUT_SEC, ) return raw @@ -2319,9 +2339,7 @@ def generate_report( llm_raw_path: Path | None = None, ) -> dict: sync_voc_report_constants(sys.modules[__name__]) - api_key = _resolve_api_key() - if not api_key: - raise RuntimeError("缺少 DASHSCOPE_API_KEY 或 .dashscope_key") + api_key = require_chat_api_key() cleaned_csv = cleaned_csv.resolve() word_freq = _load_word_freq( @@ -2339,7 +2357,11 @@ def generate_report( top2_audience = _ordered_step2_audiences(bundles, [])[:2] logger.info("阶段二 top2 受众簇: %s", top2_audience) - logger.info("调用 LLM 生成分析报告(含词频分类、簇命名、译文映射)…") + logger.info( + "调用 LLM 生成分析报告(%s,reasoning_effort=%s)…", + REPORT_MODEL, + REPORT_REASONING_EFFORT, + ) system, user = _build_report_prompt( product_name=product_name, industry=industry, diff --git a/合并评论数据.py b/合并评论数据.py index f123ecf..34e2d13 100644 --- a/合并评论数据.py +++ b/合并评论数据.py @@ -10,7 +10,7 @@ from pathlib import Path import pandas as pd _PROJECT_ROOT = Path(__file__).resolve().parent -DEFAULT_INPUT_DIR = _PROJECT_ROOT / "turkey tail mushroom-voc" +DEFAULT_INPUT_DIR = _PROJECT_ROOT / '/Users/onesvmwhoops/Documents/no scratch spray for cats' DEFAULT_OUTPUT_PATH = _PROJECT_ROOT / "merged_reviews.csv" diff --git a/向量化.py b/向量化.py index 485bf12..bc0474c 100644 --- a/向量化.py +++ b/向量化.py @@ -1,18 +1,18 @@ """ -从 voc_structured.sqlite 最新 job 展开 audience / pain_points / aspect / opinion / -aspect_opinion,调用 DashScope text-embedding-v4(256 维)写入 voc_embeddings.sqlite。 +从 voc_structured.sqlite 最新 job 展开 audience / pain_point / aspect_opinion +(与聚类.py 一致,不向量化单独的 aspect、opinion),使用本地 MLX 写入 voc_embeddings.sqlite。 溯源:source_row 与结构化时一致(CSV 第 1 条数据行=1);对应 merged_reviews_cleaned.csv 物理行号 = source_row + 1(第 1 行为表头),content 取自该数据行。 -用法(项目根目录):: +用法(项目根目录,推荐 310py 虚拟环境 Python 3.10+):: - python3 向量化.py - python3 向量化.py --workers 8 - python3 向量化.py --job-id 4 --csv merged_reviews_cleaned.csv + ./310py/bin/python 向量化.py + ./310py/bin/python 向量化.py --batch-size 16 + ./310py/bin/python 向量化.py --job-id 4 --csv merged_reviews_cleaned.csv -说明:DashScope text-embedding-v4 单次请求最多 10 条文本;脚本按批调用 API, -默认多线程并行多批(--workers),并非逐条请求。 +环境变量:VOC_EMBED_MODEL_PATH、VOC_EMBED_BATCH_SIZE(默认 16)、VOC_EMBED_MAX_TEXT_CHARS(默认 512)。 +本地 MLX 推理串行执行,--workers 仅保留兼容、固定为 1。 """ from __future__ import annotations @@ -23,13 +23,17 @@ import os import sqlite3 import struct import sys -from concurrent.futures import ThreadPoolExecutor, as_completed from dataclasses import dataclass from pathlib import Path from typing import Any, Dict, List, Sequence, Tuple from csv import DictReader -from openai import OpenAI + +from local_embedding import ( + DEFAULT_BATCH_SIZE as LOCAL_DEFAULT_BATCH, + embed_texts, + embedding_dimensions, +) logging.basicConfig( level=logging.INFO, @@ -42,16 +46,21 @@ PROJECT_ROOT = Path(__file__).resolve().parent STRUCTURED_DB = PROJECT_ROOT / "voc_structured.sqlite" EMBED_DB = PROJECT_ROOT / "voc_embeddings.sqlite" DEFAULT_CSV = PROJECT_ROOT / "merged_reviews_cleaned.csv" -DASHSCOPE_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" -EMBEDDING_MODEL = "text-embedding-v4" -EMBEDDING_DIMENSIONS = 256 -EMBED_BATCH_SIZE = 10 # DashScope text-embedding-v4 单请求 input 数组上限 10 -EMBED_DEFAULT_WORKERS = 6 # 并行批次数(每批最多 10 条) +EMBEDDING_MODEL = os.environ.get( + "VOC_EMBED_MODEL_PATH", + str(PROJECT_ROOT / "Qwen3-Embedding-4B-mxfp8"), +).strip() or "Qwen3-Embedding-4B-mxfp8" +EMBED_BATCH_SIZE = LOCAL_DEFAULT_BATCH +EMBED_DEFAULT_WORKERS = 1 # 本地 MLX 模型不可多线程并行推理 + + +def _resolve_embed_workers(explicit: int | None = None) -> int: + if explicit is not None and explicit > 1: + logger.warning("本地 embedding 仅支持串行,--workers 已忽略(使用 1)") + return 1 ENTITY_TYPES = ( "audience", "pain_point", - "aspect", - "opinion", "aspect_opinion", ) @@ -72,21 +81,6 @@ class EmbedTask: content: str -def _resolve_api_key() -> str: - v = os.environ.get("DASHSCOPE_API_KEY", "").strip() - if v: - return v - fp = os.environ.get("DASHSCOPE_API_KEY_FILE", "").strip() - if fp: - p = Path(fp).expanduser() - if p.is_file(): - return p.read_text(encoding="utf-8").strip().strip('"').strip("'") - local = PROJECT_ROOT / ".dashscope_key" - if local.is_file(): - return local.read_text(encoding="utf-8").strip().strip('"').strip("'") - return "" - - def _latest_job_id(conn: sqlite3.Connection) -> int: row = conn.execute( "SELECT id FROM analysis_jobs ORDER BY id DESC LIMIT 1" @@ -242,144 +236,55 @@ def _build_tasks( opinion = str(item.get("opinion", "")).strip() category = str(item.get("category", "")).strip() or None sentiment = str(item.get("sentiment", "")).strip() or None - if not aspect and not opinion: + if not (aspect and opinion): continue - if aspect: - tasks.append( - EmbedTask( - job_id=job_id, - extraction_id=ext_id, - source_row=source_row, - entity_type="aspect", - entity_index=i, - embed_text=aspect, - audience=aud_norm, - aspect=aspect, - opinion=opinion or None, - category=category, - sentiment=sentiment, - content=content, - ) - ) - if opinion: - tasks.append( - EmbedTask( - job_id=job_id, - extraction_id=ext_id, - source_row=source_row, - entity_type="opinion", - entity_index=i, - embed_text=opinion, - audience=aud_norm, - aspect=aspect or None, - opinion=opinion, - category=category, - sentiment=sentiment, - content=content, - ) - ) - if aspect and opinion: - merged = f"{aspect}, {opinion}" - tasks.append( - EmbedTask( - job_id=job_id, - extraction_id=ext_id, - source_row=source_row, - entity_type="aspect_opinion", - entity_index=i, - embed_text=merged, - audience=aud_norm, - aspect=aspect, - opinion=opinion, - category=category, - sentiment=sentiment, - content=content, - ) + merged = f"{aspect}, {opinion}" + tasks.append( + EmbedTask( + job_id=job_id, + extraction_id=ext_id, + source_row=source_row, + entity_type="aspect_opinion", + entity_index=i, + embed_text=merged, + audience=aud_norm, + aspect=aspect, + opinion=opinion, + category=category, + sentiment=sentiment, + content=content, ) + ) return tasks -def _embed_one_api_batch(api_key: str, texts: List[str]) -> List[bytes]: - """单次 API 调用(最多 batch_size 条文本)。""" - client = OpenAI(api_key=api_key, base_url=DASHSCOPE_BASE_URL) - resp = client.embeddings.create( - model=EMBEDDING_MODEL, - input=texts, - dimensions=EMBEDDING_DIMENSIONS, - ) - if len(resp.data) != len(texts): - raise RuntimeError( - f"embedding 返回条数 {len(resp.data)} != 请求 {len(texts)}" - ) - ordered = sorted(resp.data, key=lambda d: d.index) - blobs: List[bytes] = [] - for item in ordered: - vec = item.embedding - if len(vec) != EMBEDDING_DIMENSIONS: - raise RuntimeError(f"维度 {len(vec)} != 期望 {EMBEDDING_DIMENSIONS}") - blobs.append(_pack_embedding(vec)) - return blobs - - def _embed_batches( - api_key: str, tasks: List[EmbedTask], *, batch_size: int = EMBED_BATCH_SIZE, - workers: int = EMBED_DEFAULT_WORKERS, -) -> List[bytes]: - """ - 将任务切成每批最多 batch_size 条,并行请求 DashScope。 - 返回与 tasks 顺序一致的 embedding BLOB 列表。 - """ - if batch_size < 1 or batch_size > 10: - raise ValueError("batch_size 须在 1–10 之间(DashScope 单请求上限 10)") - workers = max(1, workers) - - chunks: List[List[EmbedTask]] = [ - tasks[i : i + batch_size] for i in range(0, len(tasks), batch_size) - ] - n_chunks = len(chunks) - if n_chunks == 0: - return [] +) -> Tuple[List[bytes], int]: + """本地 MLX 串行分批 embedding,返回 BLOB 列表与向量维度。""" + if batch_size < 1: + raise ValueError("batch_size 须 ≥ 1") + if not tasks: + return [], 0 + texts = [t.embed_text for t in tasks] + n_chunks = (len(texts) + batch_size - 1) // batch_size logger.info( - "共 %s 条文本,%s 批(每批≤%s 条),并行 workers=%s", - len(tasks), + "共 %s 条文本,%s 批(每批≤%s 条),本地 MLX 串行", + len(texts), n_chunks, batch_size, - min(workers, n_chunks), ) - # 单线程:逻辑简单,便于限流环境 - if workers == 1: - all_blobs: List[bytes] = [] - for i, chunk in enumerate(chunks, start=1): - blobs = _embed_one_api_batch(api_key, [t.embed_text for t in chunk]) - all_blobs.extend(blobs) - if i == n_chunks or i % 20 == 0: - logger.info("Embedding 进度 %s/%s 批", i, n_chunks) - return all_blobs - - results: List[List[bytes] | None] = [None] * n_chunks - done = 0 - with ThreadPoolExecutor(max_workers=min(workers, n_chunks)) as pool: - future_map = { - pool.submit( - _embed_one_api_batch, - api_key, - [t.embed_text for t in chunk], - ): idx - for idx, chunk in enumerate(chunks) - } - for fut in as_completed(future_map): - idx = future_map[fut] - results[idx] = fut.result() - done += 1 - if done == n_chunks or done % 20 == 0: - logger.info("Embedding 进度 %s/%s 批", done, n_chunks) - - return [blob for batch in results for blob in batch] # type: ignore[union-attr] + vecs = embed_texts(texts, batch_size=batch_size) + dim = len(vecs[0]) if vecs else embedding_dimensions() + blobs = [_pack_embedding(v) for v in vecs] + for i, v in enumerate(vecs): + if len(v) != dim: + raise RuntimeError(f"第 {i} 条维度 {len(v)} != {dim}") + return blobs, dim def run_embed( @@ -389,14 +294,10 @@ def run_embed( structured_db: Path = STRUCTURED_DB, embed_db: Path = EMBED_DB, batch_size: int = EMBED_BATCH_SIZE, - workers: int = EMBED_DEFAULT_WORKERS, + workers: int | None = None, reset_db: bool = True, ) -> Dict[str, Any]: - api_key = _resolve_api_key() - if not api_key: - raise RuntimeError( - "缺少 DASHSCOPE_API_KEY 或项目根 .dashscope_key" - ) + embed_workers = _resolve_embed_workers(workers) csv_path = csv_path.expanduser().resolve() if not csv_path.is_file(): @@ -438,12 +339,7 @@ def run_embed( by_type, ) - blobs = _embed_batches( - api_key, - tasks, - batch_size=batch_size, - workers=workers, - ) + blobs, embed_dim = _embed_batches(tasks, batch_size=batch_size) econn = sqlite3.connect(embed_db) try: @@ -474,7 +370,7 @@ def run_embed( task.category, task.sentiment, task.content, - EMBEDDING_DIMENSIONS, + embed_dim, blob, ) for task, blob in zip(tasks, blobs) @@ -491,7 +387,8 @@ def run_embed( "embed_db": str(embed_db), "csv": str(csv_path), "model": EMBEDDING_MODEL, - "dimensions": EMBEDDING_DIMENSIONS, + "dimensions": embed_dim, + "embed_workers": embed_workers, "extractions": len(ext_rows), "vectors": total, "by_entity_type": by_type, @@ -510,13 +407,13 @@ def main() -> None: "--batch-size", type=int, default=EMBED_BATCH_SIZE, - help="每批 API 请求条数,最大 10(DashScope 限制)", + help=f"每批本地推理条数(默认 {EMBED_BATCH_SIZE};可用 VOC_EMBED_BATCH_SIZE)", ) parser.add_argument( "--workers", type=int, - default=EMBED_DEFAULT_WORKERS, - help="并行批次数;设为 1 则串行", + default=None, + help="保留兼容;本地 MLX 固定串行 workers=1", ) args = parser.parse_args() summary = run_embed( @@ -525,7 +422,7 @@ def main() -> None: structured_db=args.structured_db, embed_db=args.embed_db, batch_size=args.batch_size, - workers=args.workers, + workers=_resolve_embed_workers(args.workers), ) print(json.dumps(summary, ensure_ascii=False, indent=2)) diff --git a/结构化_server.py b/结构化_server.py index 9c08697..395446c 100644 --- a/结构化_server.py +++ b/结构化_server.py @@ -9,22 +9,20 @@ VOC 评论结构化服务(命令行 / 直接调用,无 MCP)。 pip install -r requirements.txt ---------------------------------------------------------------------- -运行前提供密钥(任选其一;勿把密钥写进代码仓库):: +运行前提供 DeepSeek 密钥(任选其一;勿把密钥写进代码仓库):: - export DASHSCOPE_API_KEY="sk-xxx" - # 或项目根单行文件 .dashscope_key - # 或 export DASHSCOPE_API_KEY_FILE="/path/to/key.txt" + export DEEPSEEK_API_KEY="sk-xxx" + # 或项目根单行文件 .deepseek_key + # 或 export DEEPSEEK_API_KEY_FILE="/path/to/key.txt" ---------------------------------------------------------------------- 用法:: - python3 结构化_server.py --industry "Pet supplements" --product "Turkey tail mushroom for dogs" --file merged_reviews_cleaned.csv + ./310py/bin/python 结构化_server.py --industry "Pet supplements" --product "Turkey tail mushroom for dogs" --file merged_reviews_cleaned.csv - python3 结构化_server.py --smoke + ./310py/bin/python 结构化_server.py --smoke -模型:默认 ``qwen3.6-flash``(OpenAI 兼容 Chat API;可通过环境变量 ``DASHSCOPE_MODEL`` 覆盖)。 - -地域:固定中国大陆华北2(北京)``https://dashscope.aliyuncs.com/compatible-mode/v1``,不使用新加坡/国际节点。 +模型:默认 ``deepseek-v4-pro``(思考关闭;DeepSeek OpenAI 兼容 Chat API;可通过 ``DEEPSEEK_MODEL`` 覆盖)。 数据库:项目根目录 ``voc_structured.sqlite``;每次写入前会清理该库及下游 ``voc_embeddings.sqlite``、``voc_clustering.sqlite``(冒烟测试用临时库时不清理项目根文件)。 """ @@ -38,12 +36,14 @@ import os import re import sqlite3 import sys +from concurrent.futures import ThreadPoolExecutor, as_completed from csv import DictReader from datetime import datetime, timezone from pathlib import Path from typing import Any, Dict, List, Tuple from prompts.loader import product_feedback_categories +from voc_llm import CHAT_MODEL, chat_extra_body, create_chat_client, require_chat_api_key logging.basicConfig( level=logging.INFO, @@ -56,9 +56,7 @@ PROJECT_ROOT = Path(__file__).resolve().parent DB_PATH = PROJECT_ROOT / "voc_structured.sqlite" EMBED_DB = PROJECT_ROOT / "voc_embeddings.sqlite" CLUSTER_DB = PROJECT_ROOT / "voc_clustering.sqlite" -# 中国大陆华北2(北京)OpenAI 兼容模式 -DASHSCOPE_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" -MODEL_NAME = os.environ.get("DASHSCOPE_MODEL", "qwen3.6-flash").strip() +MODEL_NAME = CHAT_MODEL REVIEW_COLUMN_NAMES = ("content", "review", "评论") BATCH_FETCH_MAX_ATTEMPTS = 3 BATCH_CORRECTION_MAX_ATTEMPTS = 3 @@ -69,39 +67,41 @@ _VALID_CATEGORIES = product_feedback_categories() # 模块加载时快照;校 def _get_valid_categories() -> frozenset[str]: return product_feedback_categories() -# 模型上下文上限(保守估);单批仍用下方 DEFAULT_MAX_* 控制,保证 JSON 稳定 +# 模型上下文上限;动态分批受 DEFAULT_MAX_BATCH_INPUT_TOKENS 与 DEFAULT_MAX_BATCH_REVIEWS 约束 MODEL_MAX_INPUT_TOKENS = 991_800 MODEL_MAX_OUTPUT_TOKENS = 65_530 CHARS_PER_TOKEN_EST = 3.2 -DEFAULT_MAX_BATCH_INPUT_TOKENS = 32_000 -DEFAULT_MAX_BATCH_OUTPUT_TOKENS = 12_000 +DEFAULT_MAX_BATCH_INPUT_TOKENS = 200_000 +# 仅用于 batch_plan 日志中的输出 token 粗估,不参与分批与 API max_tokens DEFAULT_OUTPUT_TOKENS_PER_REVIEW = 450 BATCH_COUNT_MIN = 1 -BATCH_COUNT_MAX = 50 -# 批量结构化单次请求输出 token(含 JSON 开销,上限不超过模型) -BATCH_OUTPUT_TOKEN_BUFFER = 1024 -BATCH_OUTPUT_TOKEN_FLOOR = 4096 +# 动态分批时单批评论条数上限(避免单请求过大导致输出截断) +DEFAULT_MAX_BATCH_REVIEWS = 100 +# Chat 批间并行;与 embedding 共用账号时不宜过高,避免连带 429 +STRUCT_DEFAULT_WORKERS = 8 -def _batch_max_output_tokens(review_count: int) -> int: - est = review_count * DEFAULT_OUTPUT_TOKENS_PER_REVIEW + BATCH_OUTPUT_TOKEN_BUFFER - return min(MODEL_MAX_OUTPUT_TOKENS, max(BATCH_OUTPUT_TOKEN_FLOOR, est)) +def _resolve_max_batch_reviews(explicit: int | None = None) -> int: + if explicit is not None and explicit > 0: + return explicit + env = os.environ.get("VOC_STRUCT_BATCH_MAX_REVIEWS", "").strip() + if env.isdigit() and int(env) > 0: + return int(env) + return DEFAULT_MAX_BATCH_REVIEWS -def _resolve_dashscope_api_key() -> str: - """环境变量 > DASHSCOPE_API_KEY_FILE > 项目根 .dashscope_key(均仅一行密钥,无引号)。""" - v = os.environ.get("DASHSCOPE_API_KEY", "").strip() - if v: - return v - fp = os.environ.get("DASHSCOPE_API_KEY_FILE", "").strip() - if fp: - p = Path(fp).expanduser() - if p.is_file(): - return p.read_text(encoding="utf-8").strip().strip('"').strip("'") - local = PROJECT_ROOT / ".dashscope_key" - if local.is_file(): - return local.read_text(encoding="utf-8").strip().strip('"').strip("'") - return "" +def _resolve_struct_workers(explicit: int | None = None) -> int: + if explicit is not None and explicit > 0: + return explicit + env = os.environ.get("VOC_STRUCT_WORKERS", "").strip() + if env.isdigit() and int(env) > 0: + return int(env) + return STRUCT_DEFAULT_WORKERS + + +def _api_max_output_tokens() -> int: + """单次 Chat 请求 max_tokens:使用模型上限,不按条数折算。""" + return MODEL_MAX_OUTPUT_TOKENS def _load_prompt_module(): @@ -251,17 +251,15 @@ def chunk_reviews_by_token_budget( product_name: str, *, max_input_tokens: int = DEFAULT_MAX_BATCH_INPUT_TOKENS, - max_output_tokens: int = DEFAULT_MAX_BATCH_OUTPUT_TOKENS, - max_reviews_per_batch: int = BATCH_COUNT_MAX, - min_reviews_per_batch: int = BATCH_COUNT_MIN, + max_reviews_per_batch: int = DEFAULT_MAX_BATCH_REVIEWS, ) -> List[List[Tuple[int, str]]]: """ - 按预估输入/输出 token 与单批条数上限,将评论动态打包。 - 短评可多条约一批,长评自动减少条数(单条超长则单独成批)。 + 按预估输入 token 将评论动态打包。 + 单批条数不超过 max_reviews_per_batch;输入 token 超过 max_input_tokens 时拆批。 + 单条过长则单独成批。 """ max_input_tokens = min(max_input_tokens, MODEL_MAX_INPUT_TOKENS) - max_output_tokens = min(max_output_tokens, MODEL_MAX_OUTPUT_TOKENS) - max_reviews_per_batch = max(min_reviews_per_batch, max_reviews_per_batch) + max_reviews_per_batch = max(BATCH_COUNT_MIN, max_reviews_per_batch) batches: List[List[Tuple[int, str]]] = [] i = 0 @@ -269,13 +267,9 @@ def chunk_reviews_by_token_budget( batch: List[Tuple[int, str]] = [] while i < len(reviews): candidate = batch + [reviews[i]] - if len(candidate) > max_reviews_per_batch: - break - inp = _estimate_batch_input_tokens(industry, product_name, candidate) - out = _estimate_batch_output_tokens(len(candidate)) - if batch and (inp > max_input_tokens or out > max_output_tokens): + if batch and inp > max_input_tokens: break batch = candidate @@ -289,6 +283,9 @@ def chunk_reviews_by_token_budget( ) break + if len(batch) >= max_reviews_per_batch: + break + if not batch: break batches.append(batch) @@ -369,26 +366,15 @@ def _call_dashscope_chat( *, max_tokens: int | None = None, ) -> Any: - api_key = _resolve_dashscope_api_key() - if not api_key: - raise RuntimeError( - "Missing DashScope API key: set DASHSCOPE_API_KEY in the process environment, " - "or DASHSCOPE_API_KEY_FILE to a one-line key file, " - "or create a one-line file at project root: .dashscope_key" - ) + require_chat_api_key() + client = create_chat_client() + extra_body = chat_extra_body(MODEL_NAME) - from openai import OpenAI - - client = OpenAI(api_key=api_key, base_url=DASHSCOPE_BASE_URL) - extra_body: Dict[str, Any] = {} - model_lower = MODEL_NAME.lower() - if model_lower.startswith(("qwen3.6", "qwen3.5", "qwen3")) or model_lower.startswith( - "deepseek-v4" - ): - extra_body["enable_thinking"] = False - - out_tokens = max_tokens if max_tokens is not None else DEFAULT_MAX_BATCH_OUTPUT_TOKENS - out_tokens = min(MODEL_MAX_OUTPUT_TOKENS, max(BATCH_OUTPUT_TOKEN_FLOOR, out_tokens)) + out_tokens = ( + _api_max_output_tokens() + if max_tokens is None + else min(MODEL_MAX_OUTPUT_TOKENS, max_tokens) + ) resp = client.chat.completions.create( model=MODEL_NAME, @@ -483,6 +469,29 @@ _SENTIMENT_CANONICAL = { "neutral": "Neutral", } +# 模型常见误写/自创词 -> 合法 sentiment 键(校验前自动映射,减少无效重试) +_SENTIMENT_ALIASES: Dict[str, str] = { + "mixed": "neutral", + "ambivalent": "neutral", + "ambiguous": "neutral", + "both": "neutral", + "balanced": "neutral", + "unclear": "neutral", + "unsure": "neutral", + "unknown": "neutral", + "pos": "positive", + "positve": "positive", + "positiv": "positive", + "neg": "negative", + "neu": "neutral", + "somewhat positive": "positive", + "slightly positive": "positive", + "mildly positive": "positive", + "somewhat negative": "negative", + "slightly negative": "negative", + "mildly negative": "negative", +} + # 模型常见误写 -> 合法 category(校验前自动映射,减少无效重试) _CATEGORY_ALIASES: Dict[str, str] = { "value": "Price", @@ -505,9 +514,20 @@ _CATEGORY_ALIASES: Dict[str, str] = { } -def _normalize_sentiment(raw: Any) -> str: +def _canonical_sentiment_key(raw: Any) -> str | None: key = str(raw or "").strip().lower() - return _SENTIMENT_CANONICAL.get(key, "Neutral") + if not key: + return None + if key in _SENTIMENT_CANONICAL: + return key + return _SENTIMENT_ALIASES.get(key) + + +def _normalize_sentiment(raw: Any) -> str: + key = _canonical_sentiment_key(raw) + if key is None: + return "Neutral" + return _SENTIMENT_CANONICAL[key] def _normalize_category(raw: Any) -> str | None: @@ -579,11 +599,10 @@ def _validate_extraction_strict(obj: Any, key_label: str = "") -> Dict[str, Any] for sub in ("aspect", "opinion", "sentiment", "category"): if not str(item.get(sub, "")).strip(): raise ValueError(f"{prefix}product_feedback[{i}] 缺少或空的 {sub}") - sent_key = str(item.get("sentiment", "")).strip().lower() - if sent_key not in _SENTIMENT_CANONICAL: + if _canonical_sentiment_key(item.get("sentiment")) is None: raise ValueError( f"{prefix}product_feedback[{i}].sentiment 必须为 " - "Positive、Negative 或 Neutral" + "Positive、Negative 或 Neutral(禁止 Mixed 等自创词)" ) raw_cat = str(item.get("category", "")).strip() cat = _normalize_category(raw_cat) @@ -641,7 +660,8 @@ def _build_batch_correction_message(keys: List[str], errors: List[str]) -> str: f"- 输出一个 JSON 对象,顶层键必须且仅能是:{keys_literal}\n" "- 每个键的值必须包含 audience、pain_points、product_feedback\n" "- product_feedback 每条须含 aspect、opinion、sentiment" - "(Positive/Negative/Neutral)、category(禁止 Value,性价比用 Price)\n" + "(仅 Positive/Negative/Neutral,禁止 Mixed/Ambiguous;褒贬交织选主倾向或拆条)、" + "category(禁止 Value,性价比用 Price)\n" "- 只输出一个 JSON 对象(不要用数组),不要 markdown 代码围栏或解释文字" ) @@ -676,7 +696,7 @@ def _parse_batch_with_model_correction( ) try: raw_parsed = _call_dashscope_chat( - messages, max_tokens=_batch_max_output_tokens(len(sub_keys)) + messages, max_tokens=_api_max_output_tokens() ) raw_obj = _normalize_batch_response(raw_parsed, sub_keys) except (RuntimeError, json.JSONDecodeError, ValueError, TypeError) as e: @@ -823,6 +843,100 @@ def _fetch_batch_extractions( return merged +def _apply_batch_extraction_result( + *, + b_idx: int, + batch: List[Tuple[int, str]], + raw_obj: Dict[str, Any], + extractions_by_row: Dict[str, Dict[str, Any]], + sqlite_rows: List[Tuple[int, int, str, Dict[str, Any]]], + batch_details: List[Dict[str, Any]], + total_batches: int, +) -> None: + _, keys, key_to_row = _format_tagged_batch(batch) + logger.info( + "批次 %s/%s 完成:成功 %s/%s 条", + b_idx + 1, + total_batches, + len(raw_obj), + len(batch), + ) + for ck in keys: + if ck not in raw_obj: + continue + validated = raw_obj[ck] + src_row = key_to_row[ck] + extractions_by_row[str(src_row)] = validated + sqlite_rows.append((src_row, b_idx, ck, validated)) + batch_details.append( + { + "batch_index": b_idx, + "keys": keys, + "source_rows": [key_to_row[k] for k in keys], + "model_output": raw_obj, + } + ) + + +def _run_struct_batches_parallel( + industry: str, + product_name: str, + batches: List[List[Tuple[int, str]]], + workers: int, + *, + extractions_by_row: Dict[str, Dict[str, Any]], + sqlite_rows: List[Tuple[int, int, str, Dict[str, Any]]], + batch_details: List[Dict[str, Any]], +) -> None: + total_batches = len(batches) + workers = max(1, min(workers, total_batches)) + logger.info( + "结构化批间并行:%s 批,workers=%s(模型 %s)", + total_batches, + workers, + MODEL_NAME, + ) + + if workers == 1: + for b_idx, batch in enumerate(batches): + logger.info("批次 %s/%s:处理 %s 条评论…", b_idx + 1, total_batches, len(batch)) + raw_obj = _fetch_batch_extractions(industry, product_name, batch, b_idx) + _apply_batch_extraction_result( + b_idx=b_idx, + batch=batch, + raw_obj=raw_obj, + extractions_by_row=extractions_by_row, + sqlite_rows=sqlite_rows, + batch_details=batch_details, + total_batches=total_batches, + ) + return + + done = 0 + with ThreadPoolExecutor(max_workers=workers) as pool: + future_map = { + pool.submit( + _fetch_batch_extractions, industry, product_name, batch, b_idx + ): (b_idx, batch) + for b_idx, batch in enumerate(batches) + } + for fut in as_completed(future_map): + b_idx, batch = future_map[fut] + raw_obj = fut.result() + done += 1 + if done == total_batches or done % 5 == 0: + logger.info("结构化进度 %s/%s 批", done, total_batches) + _apply_batch_extraction_result( + b_idx=b_idx, + batch=batch, + raw_obj=raw_obj, + extractions_by_row=extractions_by_row, + sqlite_rows=sqlite_rows, + batch_details=batch_details, + total_batches=total_batches, + ) + + def run_analysis( industry: str, product_name: str, @@ -831,8 +945,8 @@ def run_analysis( *, clean_databases: bool = True, max_batch_input_tokens: int = DEFAULT_MAX_BATCH_INPUT_TOKENS, - max_batch_output_tokens: int = DEFAULT_MAX_BATCH_OUTPUT_TOKENS, - max_reviews_per_batch: int = BATCH_COUNT_MAX, + max_batch_reviews: int | None = None, + workers: int | None = None, ) -> Dict[str, Any]: src = str(Path(file_path).expanduser().resolve()) reviews = load_reviews_from_file(src) @@ -846,24 +960,30 @@ def run_analysis( batching_mode = "fixed" batch_size_record = batch_size else: + cap_reviews = _resolve_max_batch_reviews(max_batch_reviews) batches = chunk_reviews_by_token_budget( reviews, industry, product_name, max_input_tokens=max_batch_input_tokens, - max_output_tokens=max_batch_output_tokens, - max_reviews_per_batch=max_reviews_per_batch, + max_reviews_per_batch=cap_reviews, ) batching_mode = "dynamic" - batch_size_record = max_reviews_per_batch + batch_size_record = 0 batch_plan = _summarize_batch_plan(batches, industry, product_name) total_batches = len(batches) counts = [p["review_count"] for p in batch_plan] + if batching_mode == "dynamic" and counts: + batch_size_record = max(counts) + cap_note = "" + if batching_mode == "dynamic": + cap_note = f",单批≤{_resolve_max_batch_reviews(max_batch_reviews)}条" logger.info( - "开始结构化:共 %s 条评论,模式=%s,共 %s 批,每批条数 min/med/max=%s/%s/%s", + "开始结构化:共 %s 条评论,模式=%s%s,共 %s 批,每批条数 min/med/max=%s/%s/%s", len(reviews), batching_mode, + cap_note, total_batches, min(counts) if counts else 0, sorted(counts)[len(counts) // 2] if counts else 0, @@ -881,34 +1001,25 @@ def run_analysis( extractions_by_row: Dict[str, Dict[str, Any]] = {} batch_details: List[Dict[str, Any]] = [] sqlite_rows: List[Tuple[int, int, str, Dict[str, Any]]] = [] + struct_workers = _resolve_struct_workers(workers) - for b_idx, batch in enumerate(batches): - logger.info("批次 %s/%s:处理 %s 条评论…", b_idx + 1, total_batches, len(batch)) - _, keys, key_to_row = _format_tagged_batch(batch) - raw_obj = _fetch_batch_extractions(industry, product_name, batch, b_idx) - logger.info( - "批次 %s/%s 完成:成功 %s/%s 条", - b_idx + 1, - total_batches, - len(raw_obj), - len(batch), - ) + _run_struct_batches_parallel( + industry, + product_name, + batches, + struct_workers, + extractions_by_row=extractions_by_row, + sqlite_rows=sqlite_rows, + batch_details=batch_details, + ) + batch_details.sort(key=lambda x: x["batch_index"]) - for ck in keys: - if ck not in raw_obj: - continue - validated = raw_obj[ck] - src_row = key_to_row[ck] - extractions_by_row[str(src_row)] = validated - sqlite_rows.append((src_row, b_idx, ck, validated)) - - batch_details.append( - { - "batch_index": b_idx, - "keys": keys, - "source_rows": [key_to_row[k] for k in keys], - "model_output": raw_obj, - } + if not sqlite_rows: + require_chat_api_key() # 无结果时若缺 key 则给出明确错误 + raise RuntimeError( + f"结构化 0/{len(reviews)} 条成功。" + "请检查 DEEPSEEK_API_KEY / .deepseek_key 与网络;" + "查看上方 WARNING 中的具体原因。" ) created = datetime.now(timezone.utc).isoformat() @@ -962,7 +1073,8 @@ def run_analysis( "batch_size": batch_size_record, "batch_plan": batch_plan, "max_batch_input_tokens": max_batch_input_tokens, - "max_batch_output_tokens": max_batch_output_tokens, + "api_max_output_tokens": MODEL_MAX_OUTPUT_TOKENS, + "struct_workers": struct_workers, "model": MODEL_NAME, "extractions": extractions_by_row, "batches": batch_details, @@ -1036,19 +1148,29 @@ def main() -> None: "--max-batch-input-tokens", type=int, default=DEFAULT_MAX_BATCH_INPUT_TOKENS, - help=f"单批最大输入 token 估算上限(模型上限约 {MODEL_MAX_INPUT_TOKENS})", + help=( + f"动态分批:单批最大输入 token 估算上限(默认 {DEFAULT_MAX_BATCH_INPUT_TOKENS};" + f"模型上限约 {MODEL_MAX_INPUT_TOKENS});" + f"单批最多 {DEFAULT_MAX_BATCH_REVIEWS} 条(环境变量 VOC_STRUCT_BATCH_MAX_REVIEWS)。" + ), ) parser.add_argument( - "--max-batch-output-tokens", + "--max-batch-reviews", type=int, - default=DEFAULT_MAX_BATCH_OUTPUT_TOKENS, - help=f"单批最大输出 token 估算上限(模型上限约 {MODEL_MAX_OUTPUT_TOKENS})", + default=None, + help=( + f"动态分批单批最多评论条数(默认 {DEFAULT_MAX_BATCH_REVIEWS};" + "环境变量 VOC_STRUCT_BATCH_MAX_REVIEWS)" + ), ) parser.add_argument( - "--max-reviews-per-batch", + "--workers", type=int, - default=BATCH_COUNT_MAX, - help=f"动态模式下每批最多条数(短评可接近此值,默认 {BATCH_COUNT_MAX})", + default=None, + help=( + f"批间并行请求数(默认 {STRUCT_DEFAULT_WORKERS};" + "1=串行;可用环境变量 VOC_STRUCT_WORKERS)" + ), ) parser.add_argument( "-o", @@ -1065,8 +1187,8 @@ def main() -> None: file_path=str(args.file), batch_size=args.batch_size, max_batch_input_tokens=args.max_batch_input_tokens, - max_batch_output_tokens=args.max_batch_output_tokens, - max_reviews_per_batch=args.max_reviews_per_batch, + max_batch_reviews=args.max_batch_reviews, + workers=args.workers, ) text = json.dumps(result, ensure_ascii=False, indent=2) if args.output_json: diff --git a/聚类.py b/聚类.py index 3f0a34f..d3eb18e 100644 --- a/聚类.py +++ b/聚类.py @@ -16,9 +16,9 @@ 用法:: - python3 聚类.py + ./310py/bin/python 聚类.py # 默认自动使用 voc_structured.sqlite 中最新 analysis_jobs.id(须已向量化) - python3 聚类.py --job-id 4 # 可选:手动指定 + ./310py/bin/python 聚类.py --job-id 4 # 可选:手动指定 """ from __future__ import annotations @@ -46,6 +46,8 @@ 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", @@ -71,8 +73,7 @@ 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" -DASHSCOPE_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" -LLM_MODEL = "qwen3.6-flash" +LLM_MODEL = CHAT_MODEL UMAP_N_COMPONENTS = 30 UMAP_MIN_DIST = 0.1 @@ -129,21 +130,6 @@ class EmbedRow: embedding: np.ndarray -def _resolve_api_key() -> str: - v = os.environ.get("DASHSCOPE_API_KEY", "").strip() - if v: - return v - fp = os.environ.get("DASHSCOPE_API_KEY_FILE", "").strip() - if fp: - p = Path(fp).expanduser() - if p.is_file(): - return p.read_text(encoding="utf-8").strip().strip('"').strip("'") - local = PROJECT_ROOT / ".dashscope_key" - if local.is_file(): - return local.read_text(encoding="utf-8").strip().strip('"').strip("'") - return "" - - def _unpack_embedding(blob: bytes, dimensions: int) -> np.ndarray: n = dimensions expected = n * 4 @@ -581,9 +567,13 @@ def _ai_evaluate_cluster_samples( ], max_tokens=1500, temperature=0.0, - extra_body={"enable_thinking": False}, + response_format={"type": "json_object"}, + # 关闭思考,避免 token 耗在 reasoning_content 导致 content 为空且无 JSON + extra_body=chat_extra_body(LLM_MODEL), ) - data = _parse_json_from_llm(resp.choices[0].message.content) + 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 @@ -1079,9 +1069,7 @@ def run_clustering( cluster_db: Path = CLUSTER_DB, reset_db: bool = True, ) -> dict: - api_key = _resolve_api_key() - if not api_key: - raise RuntimeError("缺少 DASHSCOPE_API_KEY 或 .dashscope_key") + require_chat_api_key() econn = sqlite3.connect(embed_db) try: @@ -1092,7 +1080,7 @@ def run_clustering( finally: econn.close() - client = OpenAI(api_key=api_key, base_url=DASHSCOPE_BASE_URL) + client = create_chat_client() cconn = _reset_cluster_db(cluster_db, reset=reset_db) created = datetime.now(timezone.utc).isoformat() try: diff --git a/词频.py b/词频.py index 058b29b..abd3642 100644 --- a/词频.py +++ b/词频.py @@ -6,8 +6,8 @@ 用法:: - python3 词频词云.py - python3 词频词云.py --skip-llm # 复用 output/voc_terms.json,仅跑第 2 步 + ./310py/bin/python 词频.py + ./310py/bin/python 词频.py --skip-llm # 复用 output/voc_terms.json,仅跑第 2 步 """ from __future__ import annotations @@ -24,9 +24,10 @@ from collections import Counter from pathlib import Path from typing import Any, Dict, Iterable, List, Sequence, Set, Tuple -from openai import OpenAI from spacy.lang.en.stop_words import STOP_WORDS as EN_STOP_WORDS +from voc_llm import CHAT_MODEL, chat_extra_body, create_chat_client, require_chat_api_key + logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s", @@ -40,28 +41,12 @@ OUTPUT_DIR = PROJECT_ROOT / "output" TERMS_JSON = OUTPUT_DIR / "voc_terms.json" WORD_FREQ_CSV = OUTPUT_DIR / "word_freq.csv" -DASHSCOPE_BASE_URL = "https://dashscope.aliyuncs.com/compatible-mode/v1" -MODEL_NAME = os.environ.get("DASHSCOPE_MODEL", "qwen3.6-flash").strip() +MODEL_NAME = CHAT_MODEL SAMPLE_SIZE = 25 SAMPLE_SEED = 42 -def _resolve_api_key() -> str: - v = os.environ.get("DASHSCOPE_API_KEY", "").strip() - if v: - return v - fp = os.environ.get("DASHSCOPE_API_KEY_FILE", "").strip() - if fp: - p = Path(fp).expanduser() - if p.is_file(): - return p.read_text(encoding="utf-8").strip().strip('"').strip("'") - local = PROJECT_ROOT / ".dashscope_key" - if local.is_file(): - return local.read_text(encoding="utf-8").strip().strip('"').strip("'") - return "" - - def _strip_think(text: str) -> str: if not text: return text @@ -156,10 +141,9 @@ def _build_terms_prompt( def _call_llm(system: str, user: str, api_key: str) -> str: - client = OpenAI(api_key=api_key, base_url=DASHSCOPE_BASE_URL) - extra_body: Dict[str, Any] = {} - if MODEL_NAME.lower().startswith(("qwen3.6", "qwen3.5", "qwen3")): - extra_body["enable_thinking"] = False + _ = api_key + client = create_chat_client() + extra_body = chat_extra_body(MODEL_NAME) resp = client.chat.completions.create( model=MODEL_NAME, messages=[ @@ -167,7 +151,7 @@ def _call_llm(system: str, user: str, api_key: str) -> str: {"role": "user", "content": user}, ], temperature=0.3, - **({"extra_body": extra_body} if extra_body else {}), + extra_body=extra_body, ) msg = resp.choices[0].message text = msg.content or getattr(msg, "reasoning_content", None) or "" @@ -349,7 +333,9 @@ def _load_spacy(): return spacy.load("en_core_web_sm", disable=["ner", "parser"]) except OSError as e: raise RuntimeError( - "未安装 spaCy 英文模型,请执行: python3 -m spacy download en_core_web_sm" + "未安装 spaCy 英文模型,请执行: uv pip install --python 310py/bin/python " + "'en-core-web-sm @ https://github.com/explosion/spacy-models/releases/download/" + "en_core_web_sm-3.8.0/en_core_web_sm-3.8.0-py3-none-any.whl'" ) from e @@ -560,9 +546,8 @@ def run(*, skip_llm: bool = False) -> dict: ) logger.info("第 1 步跳过,复用 %s", TERMS_JSON) else: - api_key = _resolve_api_key() - if not api_key: - raise RuntimeError("缺少 DASHSCOPE_API_KEY 或 .dashscope_key") + require_chat_api_key() + api_key = "" samples = _sample_reviews(rows, n=SAMPLE_SIZE, seed=SAMPLE_SEED) meta = _step1_extract_terms( job_id=job_id, diff --git a/词频_jieba.py b/词频_jieba.py index af597e0..5a33f1a 100644 --- a/词频_jieba.py +++ b/词频_jieba.py @@ -6,8 +6,8 @@ 用法:: - python3 词频_jieba.py - python3 词频_jieba.py --skip-llm + ./310py/bin/python 词频_jieba.py + ./310py/bin/python 词频_jieba.py --skip-llm """ from __future__ import annotations @@ -36,7 +36,7 @@ from 词频 import ( _load_terms_json, _merge_word_forms, _overlaps_span, - _resolve_api_key, + require_chat_api_key, _resolve_source_path, _sample_reviews, _save_word_freq, @@ -133,9 +133,8 @@ def run(*, skip_llm: bool = False) -> dict: ) logger.info("第 1 步跳过,复用 %s", TERMS_JSON) else: - api_key = _resolve_api_key() - if not api_key: - raise RuntimeError("缺少 DASHSCOPE_API_KEY 或 .dashscope_key") + require_chat_api_key() + api_key = "" samples = _sample_reviews(rows, n=SAMPLE_SIZE, seed=SAMPLE_SEED) meta = _step1_extract_terms( job_id=job_id,