Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 12 additions & 3 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ ES_TAG_VECTORS_INDEX=somni_audio_tag_dictionary
SOMNI_ES_NODE=http://localhost:9201
SOMNI_ES_AUDIO_INDEX=somni_audio_materials
SOMNI_ES_TAG_VECTORS_INDEX=somni_audio_tag_dictionary
SOMNI_ES_SEARCH_EVENTS_INDEX=somni_audio_search_events

# MongoDB(功能手板 Fullive)
MONGO_URI=mongodb://user:password@host:27017/Fullive
Expand Down Expand Up @@ -68,10 +69,18 @@ LOG_LEVEL=INFO
LOG_DIR=logs
LOG_RETENTION=7 days

# 音频检索缓存(空 REDIS_URL 表示关闭;TTL 自写入起算,命中不续期)
# 功能手板 / HTTP 检索缓存(空 REDIS_URL 表示关闭;TTL 自写入起算,命中不续期)
REDIS_URL=redis://127.0.0.1:6379/0
# 量产 Redis(与手板隔离)
SOMNI_REDIS_URL=redis://127.0.0.1:6379/1
# 量产 Redis(独立实例;空则 GetHot 关闭,不回退 REDIS_URL)
# 本地可另起:redis-server --port 6380 --save "" --appendonly no
SOMNI_REDIS_URL=redis://127.0.0.1:6380/0
# GetHot 热点排行(Redis ZSET + ES 搜索事件索引)
SOMNI_HOT_ENABLED=true
SOMNI_HOT_TOP_N=10
SOMNI_HOT_REDIS_KEY=somni:audio:hot:v1
SOMNI_REDIS_MAX_CONNECTIONS=128
SOMNI_REDIS_CONNECT_TIMEOUT_SEC=2
SOMNI_REDIS_SOCKET_TIMEOUT_SEC=2
# 连接池大小:需 ≥ HTTP 并发峰值,过小会 Too many connections
REDIS_MAX_CONNECTIONS=512
SEARCH_CACHE_MAX_SIZE=2048
Expand Down
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@

- **核心**:三维度音频检索(ES 召回 + 标签词典向量 + 精排)
- **功能手板**:`MONGO_*` + `ES_NODE` + `REDIS_URL`;HTTP `:8080` + gRPC `:50065`
- **量产**:`SOMNI_MONGO_*` + `SOMNI_ES_*` + `SOMNI_REDIS_URL`;gRPC `:50064`
- **量产**:`SOMNI_MONGO_*` + `SOMNI_ES_*` + `SOMNI_REDIS_URL`(独立 Redis 实例);gRPC `:50064`
- **写路径**:直连 Mongo,再同步本侧 ES(不再调用 BioNode)
- **接口文档**:`docs/功能手板接口文档.md`、`docs/量产接口文档.md`

Expand Down
13 changes: 10 additions & 3 deletions app/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ class Settings(BaseSettings):
somni_es_node: str = ""
somni_es_audio_index: str = "somni_audio_materials"
somni_es_tag_vectors_index: str = "somni_audio_tag_dictionary"
somni_es_search_events_index: str = "somni_audio_search_events"

mongo_uri: str = ""
mongo_db: str = "Fullive"
Expand All @@ -54,8 +55,8 @@ class Settings(BaseSettings):
somni_mongo_answers_collection: str = "somni_quiz_answers"

sim_threshold: float = 0.7 # 内容形态向量模糊命中阈值(规范 §五-2)
# GetAudio query_text 与根标签向量相似度下限
get_audio_root_tag_sim_threshold: float = 0.85
# GetAudio query_text 与内容形态标签(含二级)向量相似度下限
get_audio_root_tag_sim_threshold: float = 0.75
# 多路文本检索厌恶硬剔除阈值;≥ 该值 penalty=1.0 丢弃候选
strong_dislike_sim_threshold: float = 0.85
search_sleep_stage_filter_enabled: bool = True # 检索步骤 1 是否按睡眠阶段过滤
Expand Down Expand Up @@ -89,8 +90,14 @@ class Settings(BaseSettings):

# 功能手板 Redis(空 URL 表示关闭)
redis_url: str = ""
# 量产 Redis(与手板隔离
# 量产 Redis(独立实例;空则 GetHot 关闭,不回退 redis_url
somni_redis_url: str = ""
somni_hot_enabled: bool = True
somni_hot_top_n: int = 10
somni_hot_redis_key: str = "somni:audio:hot:v1"
somni_redis_max_connections: int = 128
somni_redis_connect_timeout_sec: float = 2.0
somni_redis_socket_timeout_sec: float = 2.0
# 连接池需覆盖 HTTP 并发峰值;redis-py 默认仅 100,高并发易 Too many connections
redis_max_connections: int = 512
search_cache_max_size: int = 2048
Expand Down
40 changes: 40 additions & 0 deletions app/core/somni_redis.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
"""量产 Redis 客户端(与功能手板/HTTP Redis 物理隔离)。"""

from __future__ import annotations

from loguru import logger
from redis.asyncio import Redis

from app.core.config import Settings


def resolve_somni_redis_url(settings: Settings) -> str:
return settings.somni_redis_url.strip()


async def create_somni_redis(settings: Settings) -> Redis | None:
"""仅按 SOMNI_REDIS_URL 建连接;启用热点时配置错误立即阻止启动。"""
if not settings.somni_hot_enabled:
return None
url = resolve_somni_redis_url(settings)
if not url:
logger.warning("未配置 SOMNI_REDIS_URL,量产 GetHot 热点排行不可用")
return None
client = Redis.from_url(
url,
decode_responses=True,
max_connections=max(1, settings.somni_redis_max_connections),
socket_connect_timeout=max(0.1, settings.somni_redis_connect_timeout_sec),
socket_timeout=max(0.1, settings.somni_redis_socket_timeout_sec),
health_check_interval=30,
)
try:
await client.ping()
except Exception:
await client.aclose()
raise
logger.info(
"已连接量产独立 Redis,max_connections={}",
settings.somni_redis_max_connections,
)
return client
95 changes: 95 additions & 0 deletions app/es/search_events.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
"""量产音频搜索事件 ES 明细。"""

from __future__ import annotations

import asyncio
from datetime import UTC, datetime
from typing import Any

from loguru import logger

from app.core.config import Settings

_MAPPING = {
"mappings": {
"properties": {
"keyword": {"type": "keyword"},
"raw_query": {"type": "keyword"},
"created_at": {"type": "date"},
"hit_count": {"type": "integer"},
"request_id": {"type": "keyword"},
}
}
}


def _build_event_doc(
*,
keyword: str,
raw_query: str,
hit_count: int,
request_id: str,
) -> dict[str, Any]:
return {
"keyword": keyword,
"raw_query": raw_query,
"created_at": datetime.now(UTC).isoformat(),
"hit_count": int(hit_count),
"request_id": request_id,
}


class SearchEventsStore:
def __init__(self, client: Any, settings: Settings) -> None:
self._client = client
self._index = settings.somni_es_search_events_index
self._ensure_lock = asyncio.Lock()
self._index_ready = False

async def ensure_index(self) -> None:
if self._index_ready:
return
async with self._ensure_lock:
if self._index_ready:
return
if await self._client.indices.exists(index=self._index):
self._index_ready = True
return
try:
await self._client.indices.create(index=self._index, body=_MAPPING)
logger.info("已创建 ES 索引:{}", self._index)
except Exception as exc:
if not _is_already_exists_error(exc):
raise
self._index_ready = True

async def index_event(
self,
*,
keyword: str,
raw_query: str,
hit_count: int,
request_id: str = "",
) -> None:
await self.ensure_index()
doc = _build_event_doc(
keyword=keyword,
raw_query=raw_query,
hit_count=hit_count,
request_id=request_id,
)
await self._client.index(index=self._index, document=doc)


def _is_already_exists_error(exc: Exception) -> bool:
details = (
str(exc),
str(getattr(exc, "error", "")),
str(getattr(exc, "body", "")),
str(getattr(exc, "info", "")),
)
return any(
marker in detail
for detail in details
for marker in ("resource_already_exists_exception", "index_already_exists_exception")
)
11 changes: 11 additions & 0 deletions app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -55,15 +55,18 @@ def _bootstrap_dev_entry() -> None:
from app.core.config import Settings, get_settings
from app.core.exception_handlers import register_exception_handlers
from app.core.logging import setup_logging
from app.core.somni_redis import create_somni_redis
from app.embedding.encoder import Encoder, create_encoder
from app.es.client import create_es_client
from app.es.search import EsSearch
from app.es.search_events import SearchEventsStore
from app.es.sync import EsSync
from app.middleware.request_log import register_request_log_middleware
from app.server.bootstrap import GrpcServers, start_grpc_servers, stop_grpc_servers
from app.server.handboard.audio.service import AudioService
from app.server.handboard.audio.store import MaterialsStore, create_materials_store
from app.server.somni.audio.catalog import AudioCatalogService as SomniAudioService
from app.server.somni.audio.hot import HotTracker
from app.server.somni.quiz.service import QuizService as SomniQuizService
from app.server.somni.report.service import ReportService as SomniReportService
from app.services.retrieval import RetrievalService
Expand Down Expand Up @@ -189,11 +192,15 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]:
audio_index=settings.somni_es_audio_index,
tag_dictionary_index=settings.somni_es_tag_vectors_index,
)
somni_redis = await create_somni_redis(settings)
events_store = SearchEventsStore(somni_es_client, settings)
hot_tracker = HotTracker(somni_redis, events_store, settings)
_app_state.somni_audio_service = SomniAudioService(
somni_mongo,
settings,
es_search=somni_es_search,
encoder=encoder,
hot=hot_tracker,
)

start_sync_scheduler(_app_state, settings)
Expand All @@ -206,11 +213,15 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]:
await stop_grpc_servers(_app_state.grpc_servers)
_app_state.grpc_servers = None
shutdown_sync_scheduler()
if _app_state.somni_audio_service is not None:
await _app_state.somni_audio_service.drain_hot_tasks()
await shutdown_audio_search_cache(search_cache)
if materials_store is not None:
materials_store.close()
if somni_mongo is not None:
somni_mongo.close()
if somni_redis is not None:
await somni_redis.aclose()
await somni_es_client.close()
await es_client.close()

Expand Down
Loading