下一代 AI 原生电商搜索引擎架构方案(中小电商简化版)
刘志敏
2026-06-20
87 分钟阅读
Elasticsearch
AI搜索
向量检索
电商
BGE
RRF
Spring Boot
下一代 AI 原生电商搜索引擎架构方案(中小电商简化版)
1. 核心设计思路
传统搜索的瓶颈:用户输入自然语言意图("海边拍照穿的裙子"),系统在做字面关键词匹配。解决的路径只有一条:在倒排索引之外引入语义向量检索。
但中小电商的约束很现实——没有专门算法团队、运维资源有限、单次搜索的毛利撑不起 GPT-4 调用。方案原则:用 ES 8.x 一个引擎搞定所有事,不引入额外分布式系统,不上 LLM,月成本控制在 6000 以内。
1.1 四刀切下去
| 大厂方案 | 本方案 | 理由 |
|---|---|---|
| Milvus 独立部署 | 只用 ES 8.x 存标量+向量 | ES 8.x 自带 HNSW 向量索引,百万级商品(维度 1024、索引量 < 500 万)性能足够。少维护一个分布式系统,不存在数据同步的一致性问题,运维团队只需懂 ES。 |
| LLM Query 改写 | 同义词表 + 规则引擎 | 高频 query("便宜好用的"、"送女朋友")预配同义词映射和属性标签规则。低频长尾 query 靠 Dense Vector 的语义匹配能力兜底。 |
| LLM Re-ranking | ES 原生 RRF(倒数排名融合) | 不单独部署 Re-ranking 服务。BM25 负责关键词精确匹配,KNN 负责语义召回,ES 8.16+ 原生 retriever + rrf 做双路结果融合。RRF 只看相对排名不看绝对分,无需调参,效果优于线性加权。日均搜索量 < 50 万的场景下覆盖 90%+ 语义需求。 |
| RAG 回复生成 | 砍掉 | 用户搜商品是为了浏览和下单,不是跟 AI 聊天。搜索栏只返回商品列表。 |
1.2 整体架构
flowchart TB subgraph Client A[App / Web / H5] end subgraph Search_Service B[Search Service<br/>Spring Boot] C[Query Rewriter<br/>同义词表 + 规则引擎] D[Embedding Client<br/>HTTP 调用 BGE 服务] E[Hybrid Searcher<br/>BM25 + KNN + Rescore] end subgraph Storage F[(Elasticsearch 8.x<br/>product_static<br/>向量 + 文本静态索引)] F2[(Elasticsearch 8.x<br/>product_dynamic<br/>标量动态索引)] G[(MySQL<br/>SPU/SKU)] H[(Redis<br/>价格/库存/向量缓存)] end subgraph AI_Service I[Embedding Model<br/>BGE-large-zh-v1.5<br/>1 x T4 私有化部署] end A -->|搜索请求| B B --> C C -->|改写后 Query| D D -->|Query Vector| E C -->|结构化过滤条件| E E -->|RRF 混合检索| F F -->|Top-N 结果| E E -->|批量查动态字段| F2 F2 -->|价格/库存/状态| E E -->|最终结果| B B --> A F -.->|商品静态数据同步| G F2 -.->|价格/库存同步| G H -.-> D H -.-> E
1.3 请求链路
| 步骤 | 模块 | 动作 |
|---|---|---|
| 1 | App | 用户输入 "海边拍照穿的裙子" |
| 2 | Query Rewriter | 同义词表匹配:「拍照」→「摄影/旅拍」,场景标签匹配:「海边」→「度假/沙滩」 |
| 3 | Embedding Client | 调用本地 BGE 服务,将改写后 query 编码为 1024 维向量 |
| 4 | Hybrid Searcher | 同时发起 BM25 文本检索 + KNN 向量检索 |
| 5 | ES 8.x (product_static) | text 字段走 IK 分词 + BM25,dense_vector 字段走 HNSW 近似搜索,RRF 融合双路排名 |
| 6 | ES 8.x (product_dynamic) | 用商品 ID 批量查询价格、库存、销量等动态字段 |
| 7 | App | 应用层合并静态检索结果与动态字段,展示最终商品列表 |
2. 核心模块
2.1 商品向量化 (Embedding Pipeline)
模型选择
| 模型 | 维度 | 推荐理由 | 部署方式 |
|---|---|---|---|
| BGE-large-zh-v1.5 | 1024 | 中文 Embedding 事实标准,MTEB 中文榜领先,社区成熟 | 1 x T4 GPU 私有化部署,用 ONNX 或 vLLM 推理 |
| M3E-large | 1024 | 更轻量,无 GPU 也可用 CPU 推理 | CPU 部署,适合预算极紧的场景 |
结论:首选 BGE-large-zh-v1.5。T4 卡租用成本约 1500 元/月,日均可处理 50 万+ 次 Embedding 请求。
数据预处理
向量质量取决于输入文本,不是只把标题丢进去就行:
输入文本 = [类目路径] + [品牌] + [标题] + [核心属性] + [场景标签]
示例:
女装/连衣裙 > 碎花连衣裙 | 花语坊 | 法式碎花吊带裙海边度假沙滩裙
属性: 材质=雪纺, 裙长=中长裙, 风格=波西米亚
场景: 海边, 度假, 拍照
预处理流程:
- 字段清洗 → 去除 HTML 标签、特殊符号
- 属性标准化 → "S/M/L" 映射为 "小码/中码/大码"
- 场景标签 → 从标题/详情中用规则或轻量模型提取(海边、通勤、约会)
- 文本拼接 → 按优先级拼接,总长度不超过 512 tokens
- Batch Embedding → 批量编码,提升 GPU 利用率
更新机制
| 变动类型 | 同步策略 | 实现 |
|---|---|---|
| 价格 / 库存 / 销量变动 | 准实时 | Canal 监听 MySQL binlog → Kafka → 只更新 product_dynamic 索引 + Redis,不触碰静态索引 |
| 标题 / 详情 / 属性变更 | 实时 | 触发重新编码 → 更新 product_static 的 text 和 dense_vector 字段 |
| 新商品上架 | 异步 | 商品创建消息入 Kafka → Embedding Pipeline 消费 → 同时写入 product_static + product_dynamic |
| 全量重建 | 离线 | 定时任务读取 MySQL 全表 → 批量 Embedding → 重建 product_static 索引,product_dynamic 不受影响 |
关键设计:
- 动静分离双索引:
product_static(向量+文本,几乎不变)与product_dynamic(价格/库存/销量,高频更新)物理隔离 - 价格库存等频繁变动的标量字段只更新动态索引,彻底避免 HNSW 向量图因 Segment Merge 频繁重建导致的 CPU 暴涨
- 静态索引可定期
forcemerge到 1 段,彻底消除 merge 开销
2.2 查询理解 (Query Understanding)
无需 BERT 分类器,无需 LLM,用同义词表 + 规则引擎完成。
意图识别
规则匹配(毫秒级):
- 查订单: "我的订单", "订单号 *", "物流", "快递"
- 闲聊: "你好", "谢谢", "在吗"
- 商品搜索: 以上规则都不命中
对中小电商来说,订单查询和闲聊的比例通常 < 5%,用规则匹配完全够用。如果后续发现误伤率高,再考虑加 BERT 二分类。
同义词改写
模糊 query 的改写靠同义词映射表 + 规则模板:
用户输入: "便宜好用的蓝牙耳机"
同义词表命中:
"便宜" → [价格<200, 折扣>30%]
"好用" → [评分>4.5, 好评率>95%]
规则模板推导:
"蓝牙耳机" → class=蓝牙耳机, category=数码配件
结构化输出:
query: "蓝牙耳机"
filters: {price: {lte: 200}, rating: {gte: 4.5}}
同义词表维护:
- 初始阶段:从搜索日志中提取 Top 500 高频 query,人工标注属性映射
- 持续优化:每周从零结果 query 中挖掘新词,补充到同义词表
- 工具:用 Excel 或 Airtable 维护,Search Service 启动时加载到本地内存+Redis
2.3 排序策略
不依赖 CTR 模型,不用 LLM,靠 ES 原生的 RRF(Reciprocal Rank Fusion,倒数排名融合) 做双路结果合并:
为什么放弃线性加权?
原方案用
0.3 × BM25 + 0.7 × cosine 线性相加,存在两个根本问题:- BM25 分数(0
几十)与 cosine 相似度(-11)量纲不同,直接相加无意义 - 权重 0.3/0.7 需要反复调参,换个数据集就要重新调
RRF 核心原理:
score = Σ 1 / (k + rank_i)
k(rank_constant):排名常数,推荐 20(Elastic 官方建议从 20 开始,而非默认 60)rank_i:文档在第 i 个子检索器结果集中的排名(从 1 开始)
关键特性:RRF 只看相对排名,不看绝对分数,天然解决了 BM25 与 KNN 量纲不可比的问题,且无需调参。
ES 8.16+ 推荐写法(retriever API GA):
GET product_static/_search { "retriever": { "rrf": { "retrievers": [ { "knn": { "field": "title_vector", "query_vector": [0.1, 0.2, ...], "k": 50, "num_candidates": 100 } }, { "standard": { "query": { "multi_match": { "query": "搜索关键词", "fields": ["title", "description"] } } } } ], "rank_window_size": 50, "rank_constant": 20 } }, "size": 10 }
参数调优:
| 参数 | 推荐值 | 说明 |
|---|---|---|
rank_constant | 20 | 值越小,排名靠后的文档影响力越大;20 是多数场景的最佳起点 |
rank_window_size | 50~100 | 从每个子检索器取多少条参与融合,必须 >= from + size |
k(kNN) | 50 | 向量检索返回候选数 |
num_candidates | 100~500 | kNN 内部候选数,建议 k * 2~10 |
分页注意:RRF 要求
from + size <= rank_window_size,否则返回空结果。深度分页建议用 search_after。业务加权(应用层):
RRF 输出基础排序后,在应用层叠加业务因子(销量、评分、新品加权),公式:
最终分数 = RRF_score × (1 + business_boost)
例如:新品(上架 < 7 天)
business_boost = 0.1,销量 Top 10% business_boost = 0.05。3. 代码级示例
3.1 Java Spring Boot + Elasticsearch 8.x 向量写入与检索
Maven 依赖
<dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>co.elastic.clients</groupId> <artifactId>elasticsearch-java</artifactId> <version>8.16.0</version> </dependency> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> </dependency> <dependency> <groupId>com.squareup.okhttp3</groupId> <artifactId>okhttp</artifactId> <version>4.12.0</version> </dependency> <!-- Resilience4j 熔断降级 --> <dependency> <groupId>io.github.resilience4j</groupId> <artifactId>resilience4j-spring-boot3</artifactId> <version>2.1.0</version> </dependency> <dependency> <groupId>io.github.resilience4j</groupId> <artifactId>resilience4j-circuitbreaker</artifactId> <version>2.1.0</version> </dependency> <dependency> <groupId>io.github.resilience4j</groupId> <artifactId>resilience4j-timelimiter</artifactId> <version>2.1.0</version> </dependency> <dependency> <groupId>org.projectlombok</groupId> <artifactId>lombok</artifactId> <optional>true</optional> </dependency> </dependencies>
ES 配置类
@Configuration public class ElasticsearchConfig { @Value("${es.host:localhost}") private String host; @Value("${es.port:9200}") private int port; @Bean public ElasticsearchClient esClient() { RestClient restClient = RestClient.builder( new HttpHost(host, port, "http") ).build(); ElasticsearchTransport transport = new RestClientTransport(restClient, new JacksonJsonpMapper()); return new ElasticsearchClient(transport); } }
索引初始化(应用启动时执行)
@Component @Slf4j public class EsIndexInitializer implements CommandLineRunner { private static final String STATIC_INDEX = "product_static"; private static final String DYNAMIC_INDEX = "product_dynamic"; private static final int VECTOR_DIM = 1024; @Autowired private ElasticsearchClient esClient; @Override public void run(String... args) throws IOException { initStaticIndex(); initDynamicIndex(); } private void initStaticIndex() throws IOException { boolean exists = esClient.indices().exists(e -> e.index(STATIC_INDEX)).value(); if (exists) { log.info("[ES] Index {} already exists, skipping init.", STATIC_INDEX); return; } esClient.indices().create(c -> c .index(STATIC_INDEX) .settings(s -> s .refreshInterval(Time.of(t -> t.time("30s"))) .numberOfShards("6") .numberOfReplicas("1")) .mappings(m -> m .properties("product_id", p -> p.keyword(k -> k)) .properties("title", p -> p.text(t -> t .analyzer("ik_max_word") .searchAnalyzer("ik_smart"))) .properties("description", p -> p.text(t -> t .analyzer("ik_max_word"))) .properties("brand", p -> p.keyword(k -> k)) .properties("category", p -> p.keyword(k -> k)) .properties("category_path", p -> p.keyword(k -> k)) .properties("attributes", p -> p.nested(n -> n .properties("name", np -> np.keyword(k -> k)) .properties("value", np -> np.keyword(k -> k)))) .properties("title_vector", p -> p .denseVector(d -> d .dims(VECTOR_DIM) .similarity("cosine") .indexOptions(iv -> iv .type("hnsw") .m(16) .efConstruction(200)))) ) ); log.info("[ES] Static index {} created (dim={}, shards=6).", STATIC_INDEX, VECTOR_DIM); } private void initDynamicIndex() throws IOException { boolean exists = esClient.indices().exists(e -> e.index(DYNAMIC_INDEX)).value(); if (exists) { log.info("[ES] Index {} already exists, skipping init.", DYNAMIC_INDEX); return; } esClient.indices().create(c -> c .index(DYNAMIC_INDEX) .settings(s -> s .refreshInterval(Time.of(t -> t.time("5s"))) .translog(t -> t.durability("async").syncInterval(Time.of(ti -> ti.time("30s")))) .numberOfShards("3") .numberOfReplicas("1")) .mappings(m -> m .properties("product_id", p -> p.keyword(k -> k)) .properties("price", p -> p.scaledFloat(sf -> sf.scalingFactor(100.0))) .properties("original_price", p -> p.scaledFloat(sf -> sf.scalingFactor(100.0))) .properties("stock", p -> p.integer(i -> i)) .properties("sales_count", p -> p.long_(l -> l)) .properties("rating", p -> p.float_(f -> f)) .properties("status", p -> p.keyword(k -> k)) .properties("updated_at", p -> p.date(d -> d.format("strict_date_optional_time"))) ) ); log.info("[ES] Dynamic index {} created (shards=3, refresh=5s).", DYNAMIC_INDEX); } }
Resilience4j 熔断配置(application.yml)
resilience4j: circuitbreaker: instances: embedding-service: sliding-window-size: 10 failure-rate-threshold: 50 wait-duration-in-open-state: 30s permitted-number-of-calls-in-half-open-state: 3 timelimiter: instances: embedding-service: timeout-duration: 50ms cancel-running-future: true retry: instances: embedding-service: max-attempts: 2 retry-exceptions: - java.net.SocketTimeoutException
Embedding 服务调用封装(带熔断降级)
@Service @Slf4j public class EmbeddingService { private final OkHttpClient httpClient = new OkHttpClient.Builder() .connectTimeout(100, TimeUnit.MILLISECONDS) .readTimeout(50, TimeUnit.MILLISECONDS) .build(); private static final String EMBEDDING_API = "http://localhost:8001/embed"; @CircuitBreaker(name = "embedding-service", fallbackMethod = "embeddingFallback") @TimeLimiter(name = "embedding-service") @Retry(name = "embedding-service") public List<Float> getEmbedding(String text) { String json = String.format("{\"text\": \"%s\"}", text.replace("\"", "\\\"")); RequestBody body = RequestBody.create(json, okhttp3.MediaType.parse("application/json")); Request request = new Request.Builder() .url(EMBEDDING_API) .post(body) .build(); try (Response response = httpClient.newCall(request).execute()) { if (!response.isSuccessful() || response.body() == null) { throw new RuntimeException("Embedding API call failed: " + response.code()); } String respStr = response.body().string(); ObjectMapper mapper = new ObjectMapper(); JsonNode root = mapper.readTree(respStr); JsonNode vectorNode = root.get("vector"); List<Float> vector = new ArrayList<>(vectorNode.size()); for (JsonNode node : vectorNode) { vector.add((float) node.asDouble()); } return vector; } catch (IOException e) { throw new RuntimeException("Failed to call Embedding API", e); } } /** * Fallback:BGE 服务超时/异常时返回 null,触发上层降级到 BM25 纯文本检索 */ private List<Float> embeddingFallback(String text, Exception e) { log.warn("[CircuitBreaker] BGE Embedding 服务降级,query='{}', reason='{}'", text, e.getMessage()); return null; } /** * Fallback(超时专用):记录更具体的降级原因 */ private List<Float> embeddingFallback(String text, TimeoutException e) { log.warn("[CircuitBreaker] BGE Embedding 超时降级(>50ms),query='{}'", text); return null; } }
同义词改写服务
@Service public class SynonymRewriter { private final Map<String, Map<String, Object>> synonymRules; @PostConstruct public void init() { synonymRules = new HashMap<>(); // 从配置文件/数据库加载同义词规则 synonymRules.put("便宜", Map.of("price_max", 200f)); synonymRules.put("好用", Map.of("rating_min", 4.5f)); synonymRules.put("送女朋友", Map.of("tags", List.of("礼物", "女生"))); // ... } public RewriteResult rewrite(String query) { String cleanedQuery = query; Map<String, Object> filters = new HashMap<>(); for (Map.Entry<String, Map<String, Object>> rule : synonymRules.entrySet()) { if (query.contains(rule.getKey())) { cleanedQuery = cleanedQuery.replace(rule.getKey(), "").trim(); filters.putAll(rule.getValue()); } } return new RewriteResult(cleanedQuery.isEmpty() ? query : cleanedQuery, filters); } @Data @AllArgsConstructor public static class RewriteResult { private String query; private Map<String, Object> filters; } }
商品写入与混合搜索 Controller
@RestController @RequestMapping("/api/search") @Slf4j public class ProductSearchController { private static final String INDEX_NAME = "ecommerce_products"; @Autowired private ElasticsearchClient esClient; @Autowired private EmbeddingService embeddingService; @Autowired private SynonymRewriter synonymRewriter; /** * 写入商品(标量 + 向量统一存入 ES) */ @PostMapping("/products") public Map<String, String> addProduct(@RequestBody ProductDTO product) { List<Float> vector = embeddingService.getEmbedding(product.getTitle()); ProductDocument doc = new ProductDocument(); doc.setProductId(product.getProductId()); doc.setTitle(product.getTitle()); doc.setCategory(product.getCategory()); doc.setPrice(product.getPrice()); doc.setTitleVector(vector); try { esClient.index(i -> i .index(INDEX_NAME) .id(product.getProductId()) .document(doc) ); } catch (IOException e) { throw new RuntimeException("ES index failed: " + e.getMessage(), e); } log.info("[ES] Product {} indexed, vector dim={}", product.getProductId(), vector.size()); return Map.of("status", "ok", "product_id", product.getProductId()); } /** * 混合搜索:同义词改写 → RRF 融合 (BM25 + KNN) → 查动态字段 → 业务加权 */ @PostMapping("/hybrid") public Map<String, Object> hybridSearch(@RequestBody SearchRequestDTO req) { // 1. 同义词改写 SynonymRewriter.RewriteResult rewritten = synonymRewriter.rewrite(req.getQuery()); // 2. Query 向量化(带熔断降级:BGE 异常时返回 null) List<Float> queryVector = embeddingService.getEmbedding(rewritten.getQuery()); try { SearchResponse<ProductDocument> response; if (queryVector == null) { // 3a. 降级路径:BGE 熔断,只走 BM25 纯文本检索 log.warn("[Fallback] Embedding 服务不可用,降级到 BM25 纯文本检索,query='{}'", req.getQuery()); response = esClient.search(s -> s .index(STATIC_INDEX) .size(20) .query(q -> q .multi_match(m -> m .query(rewritten.getQuery()) .fields("title", "description", "brand") ) ) .minScore(0.3), ProductDocument.class ); } else { // 3b. 正常路径:RRF 融合 BM25 + KNN response = esClient.search(s -> s .index(STATIC_INDEX) .size(20) .retriever(r -> r .rrf(rrf -> rrf .retrievers( // KNN 向量检索 new Retriever.Builder<ProductDocument>() .knn(k -> k .field("title_vector") .queryVector(queryVector) .k(50) .numCandidates(100)) .build(), // BM25 文本检索 new Retriever.Builder<ProductDocument>() .standard(st -> st .query(q -> q .multi_match(m -> m .query(rewritten.getQuery()) .fields("title", "description", "brand") ) )) .build() ) .rankConstant(20) .rankWindowSize(50) ) ), ProductDocument.class ); } // 4. 提取静态索引结果的商品 ID 列表 List<String> productIds = response.hits().hits().stream() .map(h -> h.source().getProductId()) .collect(Collectors.toList()); // 5. 批量查询动态字段(价格、库存、销量) Map<String, DynamicProduct> dynamicMap = fetchDynamicFields(productIds); // 6. 合并结果 + 业务加权 List<SearchResultVO> hits = new ArrayList<>(); for (Hit<ProductDocument> hit : response.hits().hits()) { ProductDocument staticDoc = hit.source(); if (staticDoc == null) continue; DynamicProduct dynamic = dynamicMap.get(staticDoc.getProductId()); if (dynamic == null || dynamic.getStock() <= 0) continue; // 过滤无库存 double rrfScore = hit.score(); double businessBoost = calcBusinessBoost(dynamic); double finalScore = rrfScore * (1 + businessBoost); hits.add(new SearchResultVO( staticDoc.getProductId(), staticDoc.getTitle(), staticDoc.getCategory(), dynamic.getPrice(), dynamic.getStock(), dynamic.getSalesCount(), (double) Math.round(finalScore * 10000) / 10000 )); } // 按最终分数排序 hits.sort((a, b) -> Double.compare(b.getScore(), a.getScore())); log.info("[ES] Hybrid search '{}' rewritten='{}' returned {} results (vector={})", req.getQuery(), rewritten.getQuery(), hits.size(), queryVector != null); return Map.of("query", req.getQuery(), "results", hits); } catch (IOException e) { throw new RuntimeException("ES search failed: " + e.getMessage(), e); } } /** * 批量查询动态字段(product_dynamic 索引) */ private Map<String, DynamicProduct> fetchDynamicFields(List<String> productIds) throws IOException { if (productIds.isEmpty()) return Map.of(); SearchResponse<DynamicProduct> response = esClient.search(s -> s .index(DYNAMIC_INDEX) .size(productIds.size()) .query(q -> q .terms(t -> t .field("product_id") .terms(ts -> ts.value(productIds.stream() .map(FieldValue::of) .collect(Collectors.toList()))) ) ), DynamicProduct.class ); return response.hits().hits().stream() .collect(Collectors.toMap( h -> h.source().getProductId(), Hit::source, (a, b) -> a )); } /** * 业务加权:新品 +10%,热销 +5% */ private double calcBusinessBoost(DynamicProduct dynamic) { double boost = 0.0; if (dynamic.getSalesCount() > 1000) boost += 0.05; // 上架 < 7 天视为新品(需要静态索引提供 created_at) return boost; } /** * 带过滤条件的 KNN 搜索(价格区间 + 类目) * 注意:价格过滤在 product_dynamic 索引,需要先在 product_static 做 KNN,再过滤动态字段 */ @PostMapping("/knn-filter") public Map<String, Object> knnWithFilter(@RequestBody FilterSearchDTO req) { List<Float> queryVector = embeddingService.getEmbedding(req.getQuery()); // BGE 熔断降级 if (queryVector == null) { log.warn("[Fallback] Embedding 不可用,降级到 BM25 + 动态字段过滤"); return fallbackBm25WithFilter(req); } try { // Step 1: 在 product_static 做 KNN 召回(不带价格过滤,价格不在静态索引) SearchResponse<ProductDocument> staticResponse = esClient.search(s -> s .index(STATIC_INDEX) .size(100) .knn(k -> k .field("title_vector") .queryVector(queryVector) .k(50) .numCandidates(100) .filter(f -> f .bool(b -> b .must(m -> m.term(t -> t .field("category") .value(req.getCategory()))) )) ), ProductDocument.class ); List<String> productIds = staticResponse.hits().hits().stream() .map(h -> h.source().getProductId()) .collect(Collectors.toList()); // Step 2: 在 product_dynamic 过滤价格 + 库存 Map<String, DynamicProduct> dynamicMap = fetchDynamicFieldsWithFilter( productIds, req.getPriceMin(), req.getPriceMax()); List<SearchResultVO> hits = new ArrayList<>(); for (Hit<ProductDocument> hit : staticResponse.hits().hits()) { ProductDocument doc = hit.source(); if (doc == null) continue; DynamicProduct dynamic = dynamicMap.get(doc.getProductId()); if (dynamic == null) continue; hits.add(new SearchResultVO( doc.getProductId(), doc.getTitle(), doc.getCategory(), dynamic.getPrice(), dynamic.getStock(), dynamic.getSalesCount(), (double) Math.round(hit.score() * 10000) / 10000 )); } return Map.of("query", req.getQuery(), "results", hits); } catch (IOException e) { throw new RuntimeException("ES KNN filter search failed: " + e.getMessage(), e); } } /** * 降级路径:BM25 + 动态字段过滤(Embedding 不可用时) */ private Map<String, Object> fallbackBm25WithFilter(FilterSearchDTO req) { try { // 先在静态索引用 BM25 召回 SearchResponse<ProductDocument> staticResponse = esClient.search(s -> s .index(STATIC_INDEX) .size(100) .query(q -> q .bool(b -> b .must(m -> m.match(t -> t.field("title").query(req.getQuery()))) .must(m -> m.term(t -> t.field("category").value(req.getCategory()))) ) ), ProductDocument.class ); List<String> productIds = staticResponse.hits().hits().stream() .map(h -> h.source().getProductId()) .collect(Collectors.toList()); Map<String, DynamicProduct> dynamicMap = fetchDynamicFieldsWithFilter( productIds, req.getPriceMin(), req.getPriceMax()); List<SearchResultVO> hits = new ArrayList<>(); for (Hit<ProductDocument> hit : staticResponse.hits().hits()) { ProductDocument doc = hit.source(); if (doc == null) continue; DynamicProduct dynamic = dynamicMap.get(doc.getProductId()); if (dynamic == null) continue; hits.add(new SearchResultVO( doc.getProductId(), doc.getTitle(), doc.getCategory(), dynamic.getPrice(), dynamic.getStock(), dynamic.getSalesCount(), (double) Math.round(hit.score() * 10000) / 10000 )); } return Map.of("query", req.getQuery(), "results", hits, "fallback", true); } catch (IOException e) { throw new RuntimeException("Fallback search failed: " + e.getMessage(), e); } } /** * 带价格/库存过滤的动态字段查询 */ private Map<String, DynamicProduct> fetchDynamicFieldsWithFilter( List<String> productIds, float priceMin, float priceMax) throws IOException { if (productIds.isEmpty()) return Map.of(); SearchResponse<DynamicProduct> response = esClient.search(s -> s .index(DYNAMIC_INDEX) .size(productIds.size()) .query(q -> q .bool(b -> b .must(m -> m.terms(t -> t .field("product_id") .terms(ts -> ts.value(productIds.stream() .map(FieldValue::of) .collect(Collectors.toList()))))) .must(m -> m.range(r -> r .number(n -> n.field("price") .gte((double) priceMin) .lte((double) priceMax)))) .must(m -> m.range(r -> r .number(n -> n.field("stock").gt(0.0)))) ) ), DynamicProduct.class ); return response.hits().hits().stream() .collect(Collectors.toMap( h -> h.source().getProductId(), Hit::source, (a, b) -> a )); } }
DTO 定义
@Data public class ProductDTO { private String productId; private String title; private String category; private float price; } @Data public class SearchRequestDTO { private String query; private int topK = 10; } @Data public class FilterSearchDTO { private String query; private String category; private float priceMin; private float priceMax; } @Data @AllArgsConstructor public class SearchResultVO { private String productId; private String title; private String category; private float price; private double score; } @Data public class ProductDocument { private String productId; private String title; private String description; private String brand; private String category; private String categoryPath; private List<Attribute> attributes; private List<Float> titleVector; } @Data public class DynamicProduct { private String productId; private float price; private float originalPrice; private int stock; private long salesCount; private float rating; private String status; } @Data public class Attribute { private String name; private String value; }
3.2 Elasticsearch 8.x 动静分离索引 + RRF 混合检索 DSL
// ========== 1. 创建静态索引(向量 + 文本,几乎不变) ========== PUT /product_static { "settings": { "number_of_shards": 6, "number_of_replicas": 1, "refresh_interval": "30s" }, "mappings": { "properties": { "product_id": { "type": "keyword" }, "title": { "type": "text", "analyzer": "ik_max_word", "search_analyzer": "ik_smart" }, "description": { "type": "text", "analyzer": "ik_max_word" }, "brand": { "type": "keyword" }, "category": { "type": "keyword" }, "category_path": { "type": "keyword" }, "attributes": { "type": "nested", "properties": { "name": { "type": "keyword" }, "value": { "type": "keyword" } } }, "title_vector": { "type": "dense_vector", "dims": 1024, "similarity": "cosine", "index_options": { "type": "hnsw", "m": 16, "ef_construction": 200 } } } } } // ========== 2. 创建动态索引(价格/库存/销量,高频更新) ========== PUT /product_dynamic { "settings": { "number_of_shards": 3, "number_of_replicas": 1, "refresh_interval": "5s", "translog.durability": "async", "translog.sync_interval": "30s" }, "mappings": { "properties": { "product_id": { "type": "keyword" }, "price": { "type": "scaled_float", "scaling_factor": 100 }, "original_price": { "type": "scaled_float", "scaling_factor": 100 }, "stock": { "type": "integer" }, "sales_count": { "type": "long" }, "rating": { "type": "float" }, "status": { "type": "keyword" }, "updated_at": { "type": "date" } } } } // ========== 3. 写入静态数据 ========== POST /product_static/_doc/1001 { "product_id": "1001", "title": "法式碎花吊带裙海边度假沙滩裙", "description": "雪纺材质,波西米亚风格,适合海边拍照", "brand": "花语坊", "category": "连衣裙", "category_path": "女装/连衣裙/碎花连衣裙", "attributes": [ { "name": "材质", "value": "雪纺" }, { "name": "裙长", "value": "中长裙" } ], "title_vector": [0.12, -0.05, 0.33, ...] } // ========== 4. 写入动态数据 ========== POST /product_dynamic/_doc/1001 { "product_id": "1001", "price": 129.00, "original_price": 199.00, "stock": 356, "sales_count": 1280, "rating": 4.7, "status": "on_sale", "updated_at": "2026-06-20T10:30:00Z" } // ========== 5. RRF 混合检索(BM25 + KNN)========== GET /product_static/_search { "retriever": { "rrf": { "retrievers": [ { "knn": { "field": "title_vector", "query_vector": [0.15, -0.02, 0.28, ...], "k": 50, "num_candidates": 100 } }, { "standard": { "query": { "multi_match": { "query": "海边拍照裙子", "fields": ["title", "description", "brand"] } } } } ], "rank_window_size": 50, "rank_constant": 20 } }, "size": 20 } // ========== 6. 带类目过滤的 RRF 检索 ========== GET /product_static/_search { "retriever": { "rrf": { "retrievers": [ { "knn": { "field": "title_vector", "query_vector": [0.15, -0.02, 0.28, ...], "k": 50, "num_candidates": 100, "filter": { "term": { "category": "连衣裙" } } } }, { "standard": { "query": { "bool": { "must": [ { "multi_match": { "query": "海边拍照裙子", "fields": ["title", "description"] } }, { "term": { "category": "连衣裙" } } ] } } } } ], "rank_window_size": 50, "rank_constant": 20 } }, "size": 20 } // ========== 7. 批量查询动态字段(价格/库存)========== GET /product_dynamic/_search { "size": 20, "query": { "bool": { "must": [ { "terms": { "product_id": ["1001", "1002", "1003", ...] } }, { "range": { "price": { "gte": 50, "lte": 300 } } }, { "range": { "stock": { "gt": 0 } } } ] } } }
4. 冷启动与评估
4.1 冷启动方案
新商品没有用户行为数据,但向量检索引擎不依赖行为,只依赖文本内容。
| 策略 | 实现 | 效果 |
|---|---|---|
| 语义保底 | Embedding 只看标题详情,新商品只要有文本就能被召回 | 解决"搜不到" |
| 新品加权 | 在 ES script_score 中对上架时间 < 7 天的商品额外 +10% 分数 | 解决"没曝光" |
| 类目填充 | 新商品缺少属性时,用同类目下热销商品的属性均值填充 | 解决"属性缺失" |
| 冷启动流量池 | 5% 的搜索流量强制分配新商品,收集初始点击数据 | 加速数据积累 |
4.2 AB Test 设计
实验分组:
| 组别 | 排序策略 |
|---|---|
| 对照组(A) | 传统 ES BM25 + 销量排序 |
| 实验组(B) | RRF 混合检索(BM25 + KNN)+ 动静分离索引(本方案) |
分流方式:按
user_id % 100 分流,保证同一用户始终在同一组。初期流量配比 A:B = 90:10。核心指标:
- Primary:搜索转化率 CVR = 搜索后下单数 / 搜索 UV
- Secondary:点击率 CTR、人均 GMV、零结果率
实验周期:最少 2 周,覆盖工作日 + 周末。每组至少 10 万搜索 UV。
4.3 成本估算
以日均 10 万搜索 UV、50 万商品为例:
| 组件 | 规格 | 月成本 |
|---|---|---|
| Elasticsearch 8.x(静态索引) | 3 节点 x 8C 32G,1TB SSD | ~ 3000 元 |
| Elasticsearch 8.x(动态索引) | 复用静态索引集群,独立索引 | 0 元 |
| BGE Embedding 服务 | 1 节点 x 4C 16G + T4 | ~ 1500 元 |
| MySQL | 2C 8G | ~ 300 元 |
| Redis(价格/库存/向量缓存) | 2C 8G | ~ 200 元 |
| 应用服务器 | 2 节点 x 4C 8G | ~ 800 元 |
| 合计 | ~ 5800 元/月 |
5. 风险提示
-
动静分离后,动态索引的 merge 压力仍需关注。虽然静态索引几乎不变,但
product_dynamic高频更新仍会产生大量 segment。建议监控段数,超过 20 个时手动触发 forcemerge(低峰期执行)。 -
RRF 不支持高亮、rescore、sort、scroll。如果搜索列表需要关键词高亮,需在 RRF 检索后对返回的文档单独发高亮请求,或在前端用前端分词高亮兜底。
-
同义词表需要持续维护,否则会有冷门 query 漏掉。初始阶段人工标注 Top 500 高频 query 就够了。长尾 query 会持续出现,建议每周从零结果日志中捞一批新词做补充。如果团队抽不出人手维护,那就别做规则同义词,直接依赖 Dense Vector 兜底也行。
-
BGE 模型对电商专有名词可能犯傻。A字裙、JK制服、Lolita 这些词在通用中文 Embedding 上的表现不确定。上线前拿商品标题跑一遍 Batch Embedding,人工抽查 Top 10 高频类目的向量相似度排序是否合理。如果有明显偏差,用几百条商品标题做 LoRA 微调,成本不高但收益明显。
-
熔断降级会损失语义召回能力。BGE 熔断后降级到 BM25 纯文本检索,长尾语义 query(如"海边拍照穿的裙子")的召回效果会下降。建议监控降级率,持续 > 5% 时需扩容 BGE 服务或优化推理延迟。
-
分页约束。RRF 要求
from + size <= rank_window_size,深度分页必须用search_after。前端分页设计时需特别注意,避免翻到后面出现空结果。