欢迎光临
我们一直在努力

搜索系统的AI升级:从BM25到语义搜索的平滑迁移方案与工程实践

搜索系统的AI升级:从BM25到语义搜索的平滑迁移方案与工程实践

一、搜索系统升级的工程困境:不能停服的飞机换引擎

搜索是电商、内容平台最核心的用户入口。日均千万级查询的搜索系统,不能因为技术升级而中断服务。从BM25(词频-逆文档频率的经典检索模型)升级到语义搜索(基于embedding的向量检索),面临的工程约束与高空更换飞机引擎无异:系统不能停服、用户体验不能劣化、索引一致性必须保证。

BM25的优势是确定性——给定查询词,排分结果是确定的、可解释的。劣势是无法处理同义词、近义词、语言歧义("苹果"是水果还是手机)。语义搜索通过embedding向量捕捉词语的语义距离,解决了这些问题,但有两个新挑战:一是embedding模型精度直接影响搜索质量(一个差劲的模型会让搜索结果全面劣化);二是向量检索的延迟比倒排索引高1-2个数量级(需要ANN近似最近邻算法加速)。本文从混合检索架构、灰度迁移策略、生产级代码实现三个维度,提供完整的平滑迁移方案。

二、混合检索架构:BM25+语义搜索的双路融合

混合检索采用RRF(Reciprocal Rank Fusion)算法融合BM25和语义搜索的检索结果。RRF不依赖分数的绝对值(BM25和语义相似度的量纲不同),而是基于排名的倒数进行融合:RRF_score(d) = Σ 1/(k + rank_i(d)),其中k是平滑参数(通常为60),rank_i是文档在各路检索中的排名。RRF的优势在于不需要对各路检索分数做归一化,对不同检索器的分数分布不敏感。

三、生产级代码实现:混合检索与灰度迁移引擎

# search_migration_engine.py
# 搜索系统AI升级:BM25+语义混合检索与灰度迁移引擎

import abc
import hashlib
import math
from dataclasses import dataclass, field
from typing import Optional
from enum import Enum
from collections import OrderedDict
import random

class SearchMode(Enum):
"""搜索模式"""
BM25_ONLY = "bm25_only" # 纯BM25
HYBRID = "hybrid" # 混合检索
SEMANTIC_ONLY = "semantic_only" # 纯语义

@dataclass
class SearchResult:
"""搜索结果项"""
doc_id: str
title: str
bm25_score: float = 0.0
embedding_score: float = 0.0
final_score: float = 0.0
rank_bm25: int = 0
rank_embedding: int = 0
final_rank: int = 0
search_mode: SearchMode = SearchMode.BM25_ONLY

@dataclass
class SearchMetrics:
"""搜索效果指标"""
query_id: str
query_text: str
mode: SearchMode
num_results: int
latency_ms: float
p99_latency_ms: float
click_through_rate: float # 模拟
query_understanding_time_ms: float

# ==================== BM25检索器 ====================

class BM25Retriever:
"""BM25检索器(简化实现)"""

def __init__(self, k1: float = 1.2,
b: float = 0.75):
self.k1 = k1
self.b = b
self.documents: dict[str, dict] = {}
self.avg_doc_length: float = 0.0
self.inverted_index: dict[
str, dict[str, int]
] = {}
self.doc_freq: dict[str, int] = {}
self.total_docs: int = 0

def index_documents(
self, docs: list[dict]
) -> None:
"""构建倒排索引"""
self.total_docs = len(docs)
total_length = 0

for doc in docs:
doc_id = doc["id"]
words = self._tokenize(
doc.get("title", "")
+ " "
+ doc.get("content", "")
)
self.documents[doc_id] = {
"title": doc.get("title", ""),
"content": doc.get("content", ""),
"word_count": len(words),
}
total_length += len(words)

# 构建倒排索引
word_counts = {}
for word in words:
word_counts[word] = (
word_counts.get(word, 0) + 1
)

for word, count in word_counts.items():
if word not in self.inverted_index:
self.inverted_index[word] = {}
self.inverted_index[word][doc_id] = count
self.doc_freq[word] = (
self.doc_freq.get(word, 0) + 1
)

self.avg_doc_length = (
total_length / self.total_docs
if self.total_docs > 0 else 0.0
)

def search(self, query: str, top_k: int = 10
) -> list[SearchResult]:
"""BM25检索"""
query_terms = self._tokenize(query)
scores = {}

for term in query_terms:
posting_list = self.inverted_index.get(
term, {}
)
df = self.doc_freq.get(term, 0)

if df == 0:
continue

idf = math.log(
(self.total_docs – df + 0.5)
/ (df + 0.5)
+ 1.0
)

for doc_id, tf in posting_list.items():
doc_len = self.documents[doc_id][
"word_count"
]
tf_component = (
tf * (self.k1 + 1)
) / (
tf + self.k1 * (
1 – self.b + self.b
* doc_len / (
self.avg_doc_length or 1
)
)
)
scores[doc_id] = (
scores.get(doc_id, 0.0)
+ idf * tf_component
)

# 排序
sorted_docs = sorted(
scores.items(),
key=lambda x: x[1], reverse=True
)[:top_k]

results = []
for rank, (doc_id, score) in enumerate(
sorted_docs, 1
):
doc = self.documents[doc_id]
results.append(SearchResult(
doc_id=doc_id,
title=doc["title"],
bm25_score=score,
rank_bm25=rank,
search_mode=SearchMode.BM25_ONLY,
))

return results

def _tokenize(self, text: str) -> list[str]:
"""简单的分词(生产环境使用jieba/ik-analyzer)"""
# 简化的分词:按空格和标点分割
import re
words = re.findall(r'\\w+', text.lower())
return words

# ==================== 语义检索器(模拟) ====================

class SemanticRetriever:
"""语义检索器(基于embedding)"""

def __init__(self, embedding_dim: int = 768):
self.embedding_dim = embedding_dim
self.doc_embeddings: dict[
str, list[float]
] = {}
self.documents: dict[str, dict] = {}

def index_documents(
self, docs: list[dict]
) -> None:
"""构建向量索引(生产环境使用FAISS/Milvus)"""
for doc in docs:
doc_id = doc["id"]
self.documents[doc_id] = {
"title": doc.get("title", ""),
"content": doc.get("content", ""),
}
# 模拟embedding生成(生产用BGE/M3E模型)
self.doc_embeddings[doc_id] = (
self._simulate_embedding(
doc.get("title", "")
+ " "
+ doc.get("content", "")
)
)

def _simulate_embedding(
self, text: str
) -> list[float]:
"""模拟embedding(生产环境调用BGE-M3等模型)"""
# 用hash值生成确定性的模拟向量
seed = int(
hashlib.md5(text.encode()).hexdigest()[:8],
16
)
random.seed(seed)
return [
random.uniform(-1, 1)
for _ in range(self.embedding_dim)
]

def encode_query(
self, query: str
) -> list[float]:
"""将查询转换为embedding向量"""
return self._simulate_embedding(query)

def search(
self, query: str, top_k: int = 10
) -> list[SearchResult]:
"""语义检索(模拟ANN近似最近邻搜索)"""
query_vec = self.encode_query(query)
scores = {}

for doc_id, doc_vec in (
self.doc_embeddings.items()
):
# 余弦相似度
similarity = self._cosine_similarity(
query_vec, doc_vec
)
scores[doc_id] = similarity

sorted_docs = sorted(
scores.items(),
key=lambda x: x[1], reverse=True
)[:top_k]

results = []
for rank, (doc_id, score) in enumerate(
sorted_docs, 1
):
doc = self.documents[doc_id]
results.append(SearchResult(
doc_id=doc_id,
title=doc["title"],
embedding_score=score,
rank_embedding=rank,
search_mode=SearchMode.SEMANTIC_ONLY,
))

return results

def _cosine_similarity(
self, vec_a: list[float],
vec_b: list[float]
) -> float:
"""余弦相似度"""
dot = sum(a * b for a, b in zip(vec_a, vec_b))
norm_a = math.sqrt(
sum(a * a for a in vec_a)
)
norm_b = math.sqrt(
sum(b * b for b in vec_b)
)
if norm_a == 0 or norm_b == 0:
return 0.0
return dot / (norm_a * norm_b)

# ==================== RRF融合与重排序 ====================

class HybridSearchEngine:
"""混合检索引擎:BM25 + 语义 + RRF融合 + 重排序"""

RRF_K = 60 # RRF平滑参数

def __init__(self, bm25: BM25Retriever,
semantic: SemanticRetriever):
self.bm25 = bm25
self.semantic = semantic

def search(self, query: str, top_k: int = 10,
mode: SearchMode = SearchMode.HYBRID
) -> list[SearchResult]:
"""根据搜索模式执行检索"""
if mode == SearchMode.BM25_ONLY:
return self.bm25.search(query, top_k)
elif mode == SearchMode.SEMANTIC_ONLY:
return self.semantic.search(query, top_k)
else:
# 混合检索
return self._hybrid_search(query, top_k)

def _hybrid_search(
self, query: str, top_k: int
) -> list[SearchResult]:
"""RRF融合BM25和语义搜索结果"""
bm25_results = self.bm25.search(
query, top_k * 2
)
semantic_results = self.semantic.search(
query, top_k * 2
)

# RRF融合
fused_scores: dict[str, float] = {}
merged_docs: dict[str, SearchResult] = {}

# BM25路
for i, result in enumerate(bm25_results):
rrf_score = 1.0 / (self.RRF_K + i + 1)
fused_scores[result.doc_id] = rrf_score
merged_docs[result.doc_id] = result
result.rank_bm25 = i + 1

# 语义路
for i, result in enumerate(
semantic_results
):
rrf_score = 1.0 / (self.RRF_K + i + 1)
if result.doc_id in fused_scores:
fused_scores[result.doc_id] += rrf_score
else:
fused_scores[result.doc_id] = rrf_score
merged_docs[result.doc_id] = result
result.rank_embedding = i + 1

# 重排序
final_results = []
sorted_ids = sorted(
fused_scores.items(),
key=lambda x: x[1], reverse=True
)[:top_k]

for rank, (doc_id, rrf_score) in enumerate(
sorted_ids, 1
):
result = merged_docs[doc_id]
result.final_score = rrf_score
result.final_rank = rank
result.search_mode = SearchMode.HYBRID
final_results.append(result)

return final_results

# ==================== 灰度迁移引擎 ====================

class CanaryMigrationEngine:
"""搜索系统灰度迁移引擎"""

def __init__(self, hybrid_engine: HybridSearchEngine):
self.engine = hybrid_engine
self.metrics: list[SearchMetrics] = []

# 灰度配置
self.traffic_distribution = {
SearchMode.BM25_ONLY: 0.90,
SearchMode.HYBRID: 0.08,
SearchMode.SEMANTIC_ONLY: 0.02,
}

# 降级开关
self.semantic_enabled = True
self.fallback_to_bm25 = False
self.p99_latency_threshold_ms = 200.0

def update_distribution(
self, new_dist: dict[SearchMode, float]
) -> None:
"""更新流量分配比例(渐进式放量)"""
total = sum(new_dist.values())
if abs(total – 1.0) > 0.001:
raise ValueError("流量比例之和必须为1.0")
self.traffic_distribution = new_dist

def _get_search_mode(self) -> SearchMode:
"""根据流量分配决定搜索模式"""
rand = random.random()
cumulative = 0.0

for mode, proportion in (
self.traffic_distribution.items()
):
cumulative += proportion
if rand <= cumulative:
return mode

return SearchMode.BM25_ONLY

def search(self, query: str,
user_id: str = "",
top_k: int = 10
) -> list[SearchResult]:
"""带灰度路由的搜索接口"""

# 检查语义搜索是否已降级
if self.fallback_to_bm25:
results = self.engine.search(
query, top_k, SearchMode.BM25_ONLY
)
self._record_metrics(
query, SearchMode.BM25_ONLY, 0.0
)
return results

mode = self._get_search_mode()

import time
start = time.time()
results = self.engine.search(
query, top_k, mode
)
elapsed = (time.time() – start) * 1000

# 延迟超过阈值,触发降级
if (mode != SearchMode.BM25_ONLY
and elapsed > (
self.p99_latency_threshold_ms
)):
self._handle_degradation(
f"语义搜索P99延迟超标: {elapsed:.0f}ms"
)

self._record_metrics(query, mode, elapsed)
return results

def _handle_degradation(self, reason: str) -> None:
"""触发降级:停止语义搜索,全量回退BM25"""
print(f"[降级] {reason}")
self.fallback_to_bm25 = True
self.semantic_enabled = False

def _record_metrics(self, query: str,
mode: SearchMode,
latency_ms: float) -> None:
"""记录搜索指标"""
metrics = SearchMetrics(
query_id=hashlib.md5(
query.encode()
).hexdigest()[:8],
query_text=query,
mode=mode,
num_results=0,
latency_ms=latency_ms,
p99_latency_ms=latency_ms,
click_through_rate=0.0,
query_understanding_time_ms=0.0,
)
self.metrics.append(metrics)

def get_migration_progress(self) -> dict:
"""获取迁移进度报告"""
total = len(self.metrics)
if total == 0:
return {"status": "no_data"}

mode_counts = {}
mode_latencies = {}

for m in self.metrics:
mode_counts[m.mode] = (
mode_counts.get(m.mode, 0) + 1
)
if m.mode not in mode_latencies:
mode_latencies[m.mode] = []
mode_latencies[m.mode].append(
m.latency_ms
)

return {
"total_queries": total,
"distribution": {
mode.value: round(count / total, 3)
for mode, count in mode_counts.items()
},
"avg_latency_ms": {
mode.value: round(
sum(lats) / len(lats), 1
)
for mode, lats in (
mode_latencies.items()
)
},
"fallback_status": (
"active" if self.fallback_to_bm25
else "normal"
),
"semantic_enabled": self.semantic_enabled,
"current_distribution": {
k.value: v
for k, v in (
self.traffic_distribution.items()
)
},
}

# 使用示例
if __name__ == "__main__":
# 准备文档
documents = [
{"id": "doc1", "title": "苹果手机最新款",
"content": "iPhone 15 Pro Max搭载A17 Pro芯片"},
{"id": "doc2", "title": "苹果营养价值分析",
"content": "苹果含有丰富的维生素和膳食纤维"},
{"id": "doc3", "title": "机器学习入门指南",
"content": "机器学习是人工智能的核心分支"},
{"id": "doc4", "title": "深度学习框架对比",
"content": "PyTorch和TensorFlow的选型分析"},
{"id": "doc5", "title": "手机拍照技巧",
"content": "如何用手机拍摄专业级照片"},
]

# 构建索引
bm25 = BM25Retriever()
bm25.index_documents(documents)

semantic = SemanticRetriever()
semantic.index_documents(documents)

hybrid = HybridSearchEngine(bm25, semantic)

# 灰度迁移
canary = CanaryMigrationEngine(hybrid)

# 查询测试
queries = ["苹果手机", "机器学习教程", "拍照技巧"]

for q in queries:
results = canary.search(q, user_id="test")
print(f"\\n查询: {q}")
for r in results:
print(
f" [{r.final_rank}] {r.title} "
f"(mode={r.search_mode.value}, "
f"score={r.final_score:.3f})"
)

# 迁移进度报告
print("\\n=== 迁移进度 ===")
progress = canary.get_migration_progress()
for key, value in progress.items():
print(f" {key}: {value}")

四、工程落地中的关键决策:语义模型的选型与在线推理的延迟控制

语义搜索的核心是embedding模型的选择。中文场景的主流方案有三个候选:BGE-M3(BAAI开源,多语言支持,推理延迟约30ms@GPU);M3E(moka-ai开源,中文专精,延迟约15ms@CPU);Cohere Embed(商业API,多语言,延迟80ms+网络)。选型的权衡是精度vs延迟vs成本的三角:BGE-M3精度最高但需要GPU(成本高),M3E精度略低但CPU可推理(成本低),Cohere精度高但延迟不可控(依赖API响应时间)。

推荐的架构是三级缓存:L1(本地缓存)命中率约60%,延迟<1ms——热门查询的embedding结果缓存到Redis;L2(本地推理服务)命中率约35%,延迟15-30ms——自建BGE-M3推理服务(Triton Inference Server部署);L3(备选方案)命中率约5%,延迟80ms——Cohere API作为冷启动和fallback。通过三级缓存,平均延迟从30ms降至约8ms。

灰度迁移的推荐节奏是30天渐进式放量:Day1-7:2%语义+8%混合+90%BM25(验证模型精度和延迟基线);Day8-14:5%语义+20%混合+75%BM25(收集用户行为对比数据);Day15-21:10%语义+40%混合+50%BM25(A/B测试:对比点击率/转化率/停留时间);Day22-28:20%语义+60%混合+20%BM25(准备全面切换);Day29-30:5%语义+5%BM25+90%混合(以混合搜索为主要模式)。

五、总结

搜索系统从BM25到语义搜索的平滑迁移采用双路检索+RRF融合的架构。BM25负责精确匹配(倒排索引,确定性排序),语义搜索负责语义理解(向量检索,捕捉同义词和意图)。RRF算法通过1/(k+rank)融合两路排名,无需分数归一化,k取值60对排名融合效果最稳定。迁移策略采用渐进式灰度放量:30天从2%→90%的语义流量占比,每一步放量后都对比CTR、转化率、P99延迟三个核心指标。延迟控制采用三级缓存架构:L1 Redis本地缓存(60%命中<1ms)、L2自建Triton推理服务(35%命中15-30ms)、L3 Cohere API(5%命中80ms),平均延迟约8ms。降级策略是自动化的:语义搜索P99延迟>200ms自动触发全量回退BM25,恢复后手动逐步放量。embedding模型选型在CPU场景选择M3E(15ms@CPU),GPU场景选择BGE-M3(30ms@GPU),商业场景选择Cohere。迁移成功的关键指标不是语义搜索的比例,而是搜索质量的提升幅度——通过离线NDCG和在线CTR的A/B对比验证语义升级的实际价值。

赞(0)
未经允许不得转载:171主机测评 » 搜索系统的AI升级:从BM25到语义搜索的平滑迁移方案与工程实践
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址