Skip to content

Commit fdf3728

Browse files
edycursoragent
authored andcommitted
音频检索 Somni: Mongo 差异同步与 somni ES 索引检索适配
将数据源从 comm 扁平表切换为 Mongo Somni 集合,新增 index_mappings 定义 somni 双索引 mapping,同步脚本、检索流水线与配套单测一并更新。 Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 20605f5 commit fdf3728

20 files changed

Lines changed: 1146 additions & 870 deletions

.env.example

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,12 @@ APP_DEBUG=false
55

66
# Elasticsearch
77
ES_NODE=http://localhost:9200
8-
ES_AUDIO_INDEX=audio_materials
9-
ES_TAG_VECTORS_INDEX=tag_vectors
8+
ES_AUDIO_INDEX=somni_audio_materials
9+
ES_TAG_VECTORS_INDEX=somni_audio_tag_dictionary
10+
11+
# MongoDB(Somni 数据同步源)
12+
MONGO_URI=mongodb://user:password@host:27017/Fullive
13+
MONGO_DB=Fullive
1014

1115
# comm-service gRPC
1216
COMM_GRPC_HOST=bionode-test.fulai.tech
@@ -16,21 +20,21 @@ COMM_GRPC_USE_TLS=true
1620
# 检索参数
1721
SIM_THRESHOLD=0.7
1822
DEFAULT_TOP_K=10
19-
# 检索步骤 1 是否按睡眠阶段过滤(暂时关闭可设为 false)
20-
SEARCH_SLEEP_STAGE_FILTER_ENABLED=false
23+
# 检索步骤 1 是否按 sleep_stage_tags 做 nested 精确过滤
24+
SEARCH_SLEEP_STAGE_FILTER_ENABLED=true
2125

2226
# Embedding(bge-small-zh-v1.5,512 维;换模型须重建 ES 并全量重算向量)
2327
EMBEDDING_BACKEND=onnx
2428
EMBEDDING_MODEL=BAAI/bge-small-zh-v1.5
2529
EMBEDDING_DIM=512
2630
EMBEDDING_ONNX_DIR=models/onnx/bge-small-zh-v1.5
2731

28-
# comm → ES 差异同步(服务内每周定时 + 可手动 scripts/sync_es_from_comm.py)
32+
# Mongo → ES 差异同步(服务内每周定时 + 可手动 scripts/sync_es_from_comm.py)
2933
SYNC_ENABLED=true
3034
SYNC_INTERVAL_DAYS=7
3135
SYNC_PAGE_SIZE=100
3236
SYNC_BACKUP_DIR=data/sync_backup
33-
SYNC_BACKUP_FILENAME=audio_materials_backup.json
37+
SYNC_BACKUP_FILENAME=somni_audio_materials_backup.json
3438

3539
# 日志(控制台 + logs/uburnode.log 文件双写,每次请求自动记录)
3640
LOG_LEVEL=INFO

README.md

Lines changed: 131 additions & 80 deletions
Original file line numberDiff line numberDiff line change
@@ -1,83 +1,92 @@
1-
# UburNode
1+
# UburPython
22

3-
BioNode 体系中的**中间层服务****音频检索**为核心能力,同时承担算法端与底层存储之间的**数据结构转换**
3+
BioNode 体系中的 **Somni 音频检索服务**以三维度检索为核心,从 MongoDB 同步 Somni 原料与标签词典至 Elasticsearch,对外提供 HTTP 检索 API
44

5-
- **核心**:三维度音频检索(ES 召回 + 进程内 Embedding + 四步精排流水线)
6-
- **中间层**:统一对外 HTTP/Pydantic 契约,对内经 gRPC 访问 comm-service;在 HTTP、Mongo(扁平标签)、ES(六维结构 + 向量)之间做形态互转,避免算法端直连 Mongo 或各自维护多套字段约定
5+
- **核心**:三维度音频检索(`somni_audio_materials` ES 召回 + 标签词典向量 + 四步精排流水线)
6+
- **数据源**:MongoDB `Fullive` 库(`somni_audio_materials``somni_audio_tag_dictionary`
7+
- **索引**:Elasticsearch `somni_audio_materials``somni_audio_tag_dictionary`(字段含义见 mapping `meta.description`
8+
- **写路径(遗留)**:HTTP CUD 仍经 comm-service gRPC;Somni 索引由同步脚本维护
79

810
## 架构
911

1012
```text
1113
算法端 / 调用方
1214
1315
14-
对外 HTTP (FastAPI + Pydantic) ← 中间层:契约统一 + 数据结构转换
16+
对外 HTTP (FastAPI + Pydantic)
1517
16-
├──读──► Elasticsearch + Embedding + RetrievalService
18+
├──读(检索)──► somni_audio_materials (ES)
19+
│ + somni_audio_tag_dictionary (ES 向量)
20+
│ + 进程内 Embedding (bge-small-zh-v1.5)
1721
18-
└──写──► comm-service (gRPC) ──► MongoDB(真值库,扁平 tags)
19-
└── EsSync ──► Elasticsearch(索引副本,六维 tags + vector_id)
22+
├──写(CUD,遗留)──► comm-service (gRPC) ──► MongoDB(旧 comm 表)
23+
24+
└──同步(定时/手动)──► MongoDB Somni 集合 ──► Elasticsearch
2025
```
2126

22-
## 中间层数据转换
23-
24-
UburNode 不只做检索,还在各存储边界维持**单一对外契约**并完成形态映射:
27+
## 数据流
2528

26-
| 边界 | 入站形态 | 出站形态 | 负责模块 |
27-
|------|----------|----------|----------|
28-
| HTTP 写(CUD) | 六维标签对象 `AudioTagsInput` | comm/Mongo 扁平 `string[]`(带维度前缀) | `app/schemas/audio.py``app/core/tags.py` |
29-
| HTTP 写响应 | comm `AudioMaterialInfo`(gRPC) | HTTP `AudioMaterialData` | `app/schemas/audio.py` |
30-
| ES 同步 | 扁平 tags + 业务字段 | ES 六维 `tags` + `tag_vectors` embedding | `app/es/sync.py` |
31-
| HTTP 读(检索| ES 六维 `TagItem`(含 `vector_id`| 出参六维 label 字符串 | `app/schemas/audio.py``app/services/retrieval.py` |
29+
| 环节 | 来源 | 目标 | 模块 |
30+
|------|------|------|------|
31+
| Mongo → ES 同步 | `somni_audio_materials``somni_audio_tag_dictionary` | 同名 ES 索引 | `scripts/sync_es_from_comm.py` |
32+
| 标签向量 | 词典 `name` / `name_en` | `name_vector` / `name_en_vector` | 同步脚本 + `app/embedding/` |
33+
| HTTP 检索 | ES 原料文档 | `data.materials[]` 原样返回 | `app/services/retrieval.py` |
34+
| HTTP CUD(遗留| 六维标签入参 | comm 扁平 tags | `app/services/audio.py` |
3235

33-
字段命名全链路 **snake_case**;对外以 Pydantic + OpenAPI 为唯一 HTTP 契约,对内 comm 调用走同源 `bionode_comm.proto`
36+
字段命名全链路 **snake_case**。Somni 表结构详见仓库内 `音频表结构.md`
3437

3538
## 目录结构
3639

3740
```text
38-
UburNode/
41+
UburPython/
3942
├── app/
40-
│ ├── main.py # FastAPI 入口 + lifespan
41-
│ ├── core/ # 配置、日志
42-
│ ├── api/audio.py # 4 个 HTTP 端点
43-
│ ├── schemas/audio.py # Pydantic 模型
44-
│ ├── services/ # AudioService、RetrievalService
45-
│ ├── es/ # EsSearch、EsSync
46-
│ ├── embedding/encoder.py # bge-small-zh-v1.5 向量编码
47-
│ ├── bionode_grpc_clients/ # BioNode 外部微服务 gRPC 客户端
48-
│ │ └── comm/ # comm-service(client.py + grpc_gen/)
49-
├── proto/ # bionode_comm.proto(唯一真源)
50-
├── scripts/gen_proto.sh # 生成 gRPC stub
43+
│ ├── main.py # FastAPI 入口 + lifespan
44+
│ ├── core/ # 配置、日志、标签转换
45+
│ ├── api/audio.py # 4 个 HTTP 端点
46+
│ ├── schemas/audio.py # Pydantic 模型
47+
│ ├── services/ # AudioService、RetrievalService
48+
│ ├── es/
49+
│ │ ├── search.py # EsSearch 读路径
50+
│ │ ├── sync.py # EsSync 写路径(CUD 跳过 upsert)
51+
│ │ └── index_mappings.py # ES 索引 mapping + 字段注释
52+
│ ├── embedding/encoder.py # bge-small-zh-v1.5 向量编码
53+
│ └── bionode_grpc_clients/ # comm-service gRPC 客户端
54+
├── scripts/
55+
│ ├── sync_es_from_comm.py # Mongo → ES 差异同步
56+
│ └── gen_proto.sh # 生成 gRPC stub
57+
├── proto/ # bionode_comm.proto
5158
├── tests/
52-
├── .cursor/skills/ # 项目 Skill
5359
├── pyproject.toml
5460
└── .env.example
5561
```
5662

5763
## 快速开始
5864

5965
```bash
60-
# 1. 创建虚拟环境并安装依赖
61-
python3.12 -m venv .venv
62-
source .venv/bin/activate
63-
pip install -e ".[dev]"
66+
# 1. 安装依赖(推荐 uv)
67+
uv sync --extra dev
6468

65-
# 2. 生成 comm gRPC stub
69+
# 2. 生成 comm gRPC stub(CUD 接口需要)
6670
chmod +x scripts/gen_proto.sh
6771
./scripts/gen_proto.sh
6872

69-
# 3. 本地 Elasticsearch(向量索引,需先就绪)
70-
# 未安装 Docker 时:brew install --cask docker-desktop,打开 Docker Desktop 等待就绪
73+
# 3. 本地 Elasticsearch
7174
docker compose -f docker-compose.es.yml up -d
72-
curl -s http://localhost:9200 # 应返回 cluster 信息
73-
# .env 默认 ES_NODE=http://localhost:9200;索引由应用启动时 ensure_indices 自动创建
75+
curl -s http://localhost:9200
7476

7577
# 4. 配置环境变量
7678
cp .env.example .env
77-
# 编辑 ES_NODE、COMM_GRPC_HOST 等
79+
# 编辑 ES_NODE、MONGO_URI、EMBEDDING_ONNX_DIR、COMM_GRPC_* 等
80+
81+
# 5. 导出 ONNX 模型(若 models/ 目录尚无模型)
82+
# 见 scripts/export_onnx_model.py
83+
84+
# 6. Mongo → ES 全量同步
85+
uv run python scripts/sync_es_from_comm.py --dry-run
86+
uv run python scripts/sync_es_from_comm.py
7887

79-
# 5. 启动服务
80-
uvicorn app.main:app --host 0.0.0.0 --port 8080 --reload
88+
# 7. 启动服务
89+
uv run uvicorn app.main:app --host 0.0.0.0 --port 8080 --reload
8190
```
8291

8392
开发模式(`APP_DEBUG=true`)跳过 Embedding 模型加载,便于本地调试 HTTP 路由。
@@ -86,20 +95,56 @@ uvicorn app.main:app --host 0.0.0.0 --port 8080 --reload
8695

8796
| 端点 | 方法 | 说明 |
8897
|------|------|------|
89-
| `/api/audio` | POST | 创建音频并同步 ES |
90-
| `/api/audio/{id}` | PUT | 更新音频并同步 ES |
91-
| `/api/audio/{id}` | DELETE | 删除音频并同步 ES |
98+
| `/api/audio` | POST | 创建音频(comm gRPC,遗留) |
99+
| `/api/audio/{id}` | PUT | 更新音频(comm gRPC,遗留) |
100+
| `/api/audio/{id}` | DELETE | 删除音频 |
92101
| `/api/audio/search` | POST | 三维度检索 |
93102

94103
OpenAPI 文档:启动后访问 `http://localhost:8080/docs`
95104

96-
### comm-service gRPC 连通探测
105+
### 检索接口
97106

98-
```bash
99-
source .venv/bin/activate # 或: .venv/bin/python scripts/test_grpc_connect.py
100-
python scripts/test_grpc_connect.py
101-
# 或集成测试(需可达的 COMM_GRPC_HOST)
102-
COMM_GRPC_INTEGRATION=1 pytest tests/test_comm_grpc.py -v
107+
**请求** `POST /api/audio/search`
108+
109+
```json
110+
{
111+
"sleep_stage_tags": ["放松"],
112+
"content_tags": ["慢钢琴"],
113+
"disliked_tags": [],
114+
"top_k": 10
115+
}
116+
```
117+
118+
**响应** `data.materials` 为命中条目的 `somni_audio_materials` 索引文档(含 `id`,字段与 ES/Mongo 一致,暂不做裁剪):
119+
120+
```json
121+
{
122+
"code": 200,
123+
"msg": "检索成功",
124+
"data": {
125+
"materials": [
126+
{
127+
"id": "6a33a7928030d4cf420efeb6",
128+
"audio_name": "专属冥想南极 助眠解压舒缓情绪",
129+
"description": "...",
130+
"status": true,
131+
"audio_url": "https://cdn.fulai.tech/comm/audio/xxx.mp3",
132+
"operation_type": 0,
133+
"created_by": "qwen3.5-omni-plus",
134+
"updated_by": "qwen3.5-omni-plus",
135+
"sleep_stage_tags": [{ "tag_id": "...", "code": "unwind", "name": "放松" }],
136+
"content_form_tags": [],
137+
"mechanism_tags": [],
138+
"audio_engineering_tags": [],
139+
"medical_risk_tags": [],
140+
"evidence_level_tags": [{ "tag_id": "...", "code": "B", "name": "中等证据" }],
141+
"created_at": "2026-06-18T00:00:00.000Z",
142+
"updated_at": "2026-06-18T00:00:00.000Z"
143+
}
144+
]
145+
},
146+
"timestamp": "..."
147+
}
103148
```
104149

105150
## 检索流水线
@@ -108,17 +153,41 @@ COMM_GRPC_INTEGRATION=1 pytest tests/test_comm_grpc.py -v
108153
睡眠阶段精确过滤 → 内容形态准入 → 厌恶剔除 + 粗排 → 精排
109154
```
110155

111-
## Proto 变更
156+
| 步骤 | 说明 |
157+
|------|------|
158+
| 1 | `sleep_stage_tags.name` nested 精确匹配(可配置跳过) |
159+
| 2 | `content_tags` 与内容/机制/工程标签精确或向量模糊命中 |
160+
| 3 | `disliked_tags` 向量相似则剔除 |
161+
| 4 |`match_count` 降序,`top_k` 截断 |
112162

113-
`comm-service` 修改 `proto/bionode_comm.proto` 后须重新生成 stub:
163+
## Mongo → ES 同步
114164

115165
```bash
116-
./scripts/gen_proto.sh
166+
uv run python scripts/sync_es_from_comm.py # 正式同步
167+
uv run python scripts/sync_es_from_comm.py --dry-run # 仅比对统计
117168
```
118169

119-
## 日志
170+
- 先同步 `somni_audio_tag_dictionary`(写入 `name_vector``name_en_vector`
171+
- 再同步 `somni_audio_materials`(1:1 镜像 Mongo 文档)
172+
- 启动时删除旧索引 `audio_materials``tag_vectors`
173+
174+
服务内按 `SYNC_INTERVAL_DAYS` 定时执行(需配置 `MONGO_URI`)。
175+
176+
## 环境变量(节选)
120177

121-
每次 HTTP 请求自动写入日志文件(`RequestLogMiddleware`),包含 method、path、status、耗时、`request_id`
178+
| 变量 | 默认值 | 说明 |
179+
|------|--------|------|
180+
| `ES_NODE` | `http://localhost:9200` | Elasticsearch 地址 |
181+
| `ES_AUDIO_INDEX` | `somni_audio_materials` | 原料索引名 |
182+
| `ES_TAG_VECTORS_INDEX` | `somni_audio_tag_dictionary` | 标签词典索引名 |
183+
| `MONGO_URI` || MongoDB 连接串(同步必填) |
184+
| `MONGO_DB` | `Fullive` | 数据库名 |
185+
| `SIM_THRESHOLD` | `0.7` | 向量模糊命中阈值 |
186+
| `EMBEDDING_ONNX_DIR` | `models/onnx/bge-small-zh-v1.5` | ONNX 模型目录 |
187+
188+
完整列表见 [`.env.example`](.env.example)
189+
190+
## 日志
122191

123192
| 配置项 | 默认值 | 说明 |
124193
|--------|--------|------|
@@ -127,36 +196,18 @@ COMM_GRPC_INTEGRATION=1 pytest tests/test_comm_grpc.py -v
127196
| `LOG_ROTATION` | `10 MB` | 单文件滚动大小 |
128197
| `LOG_RETENTION` | `7 days` | 历史日志保留 |
129198

130-
日志同时输出到控制台和 `logs/uburnode.log`。响应头会回传 `X-Request-Id` 便于链路追踪。
131-
132-
## Docker 部署(服务器)
199+
响应头回传 `X-Request-Id` 便于链路追踪。
133200

134-
`intelligent_reimbursement` 相同:**服务器 git pull + docker compose build**,不依赖 GHCR。
201+
## Docker 部署
135202

136203
```bash
137-
# 1. 本机一键写入 GitHub Secrets(SSH_HOST / SSH_USER / SSH_PRIVATE_KEY)
138-
chmod +x scripts/setup_github_secrets.sh
139-
./scripts/setup_github_secrets.sh
140-
141-
# 2. 服务器一次性准备
142-
# - 将公钥写入 ~/.ssh/authorized_keys
143-
# - 已安装 Docker 与 Compose;/etc/docker/daemon.json 镜像加速自行维护
144-
# - 克隆仓库并配置 .env:
145-
mkdir -p /opt/uburnode
146-
git clone -b dev https://github.com/dwqnidq/UburNode.git /opt/uburnode
147-
cp /opt/uburnode/.env.example /opt/uburnode/.env # 编辑 COMM_GRPC_* 等
148-
149-
# 3. 首次手动启动(build 约 20~40 分钟)
150-
cd /opt/uburnode && docker compose up -d --build
151-
152-
# 4. 以后:GitHub → Actions → Deploy UburNode → Run workflow
153-
# 或 push 到 dev 分支自动部署(git pull + compose up --build)
204+
cd /opt/uburpython && docker compose up -d --build
154205
```
155206

156-
生产访问:`http://<服务器IP>:8001/docs`(nginx 映射宿主机 8001 → 容器 80,避免与宝塔 80 冲突)。
207+
生产访问:`http://<服务器IP>:8001/docs`(nginx 映射宿主机 8001 → 容器 80)。
157208

158209
## 测试
159210

160211
```bash
161-
pytest
212+
uv run pytest
162213
```

app/api/audio.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,6 @@ async def search_audio(
6464
body: SearchAudioRequest,
6565
service: AudioService = Depends(get_audio_service),
6666
) -> ApiResponse:
67-
"""三维度检索:只读 ES,不写 Mongo / ES(规范红线)。"""
67+
"""三维度检索:只读 ES,返回 somni_audio_materials 索引文档列表。"""
6868
result = await service.search_audio(body)
6969
return success(data=result.model_dump(), msg="检索成功")

app/core/config.py

Lines changed: 15 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -21,15 +21,20 @@ class Settings(BaseSettings):
2121
app_debug: bool = False
2222

2323
es_node: str = "http://localhost:9200"
24-
es_audio_index: str = "audio_materials"
25-
es_tag_vectors_index: str = "tag_vectors"
24+
es_audio_index: str = "somni_audio_materials"
25+
es_tag_vectors_index: str = "somni_audio_tag_dictionary"
26+
27+
mongo_uri: str = ""
28+
mongo_db: str = "Fullive"
29+
mongo_materials_collection: str = "somni_audio_materials"
30+
mongo_tag_dictionary_collection: str = "somni_audio_tag_dictionary"
2631

2732
comm_grpc_host: str = "bionode-test.fulai.tech"
2833
comm_grpc_port: int = 443
2934
comm_grpc_use_tls: bool = True # 443 走 TLS;内网明文可设 false
3035

3136
sim_threshold: float = 0.7 # 内容形态向量模糊命中阈值(规范 §五-2)
32-
search_sleep_stage_filter_enabled: bool = False # 检索步骤 1 是否按睡眠阶段过滤
37+
search_sleep_stage_filter_enabled: bool = True # 检索步骤 1 是否按睡眠阶段过滤
3338

3439
embedding_backend: str = "onnx" # onnx(生产)| torch(对比/回退,需 sentence-transformers)
3540
embedding_model: str = "BAAI/bge-small-zh-v1.5"
@@ -42,12 +47,13 @@ class Settings(BaseSettings):
4247
log_rotation: str = "10 MB"
4348
log_retention: str = "7 days"
4449

45-
# comm → ES 差异同步(服务内定时 + scripts/sync_es_from_comm.py 手动)
50+
# Mongo → ES 差异同步(服务内定时 + scripts/sync_es_from_comm.py 手动)
4651
sync_enabled: bool = True
4752
sync_interval_days: int = 7
4853
sync_page_size: int = 100
4954
sync_backup_dir: str = "data/sync_backup"
50-
sync_backup_filename: str = "audio_materials_backup.json"
55+
sync_backup_filename: str = "somni_audio_materials_backup.json"
56+
sync_tag_dictionary_backup_filename: str = "somni_audio_tag_dictionary_backup.json"
5157

5258
@property
5359
def embedding_onnx_path(self) -> Path:
@@ -61,6 +67,10 @@ def embedding_tokenizer_dir(self) -> Path:
6167
def sync_backup_path(self) -> Path:
6268
return Path(self.sync_backup_dir) / self.sync_backup_filename
6369

70+
@property
71+
def sync_tag_dictionary_backup_path(self) -> Path:
72+
return Path(self.sync_backup_dir) / self.sync_tag_dictionary_backup_filename
73+
6474
@property
6575
def comm_grpc_target(self) -> str:
6676
return f"{self.comm_grpc_host}:{self.comm_grpc_port}"

0 commit comments

Comments
 (0)