L刘志敏

下一代 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-rankingES 原生 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 请求链路

步骤模块动作
1App用户输入 "海边拍照穿的裙子"
2Query Rewriter同义词表匹配:「拍照」→「摄影/旅拍」,场景标签匹配:「海边」→「度假/沙滩」
3Embedding Client调用本地 BGE 服务,将改写后 query 编码为 1024 维向量
4Hybrid Searcher同时发起 BM25 文本检索 + KNN 向量检索
5ES 8.x (product_static)text 字段走 IK 分词 + BM25,dense_vector 字段走 HNSW 近似搜索,RRF 融合双路排名
6ES 8.x (product_dynamic)用商品 ID 批量查询价格、库存、销量等动态字段
7App应用层合并静态检索结果与动态字段,展示最终商品列表

2. 核心模块

2.1 商品向量化 (Embedding Pipeline)

模型选择

模型维度推荐理由部署方式
BGE-large-zh-v1.51024中文 Embedding 事实标准,MTEB 中文榜领先,社区成熟1 x T4 GPU 私有化部署,用 ONNX 或 vLLM 推理
M3E-large1024更轻量,无 GPU 也可用 CPU 推理CPU 部署,适合预算极紧的场景
结论:首选 BGE-large-zh-v1.5。T4 卡租用成本约 1500 元/月,日均可处理 50 万+ 次 Embedding 请求。

数据预处理

向量质量取决于输入文本,不是只把标题丢进去就行:
输入文本 = [类目路径] + [品牌] + [标题] + [核心属性] + [场景标签]

示例:
女装/连衣裙 > 碎花连衣裙 | 花语坊 | 法式碎花吊带裙海边度假沙滩裙
属性: 材质=雪纺, 裙长=中长裙, 风格=波西米亚
场景: 海边, 度假, 拍照
预处理流程
  1. 字段清洗 → 去除 HTML 标签、特殊符号
  2. 属性标准化 → "S/M/L" 映射为 "小码/中码/大码"
  3. 场景标签 → 从标题/详情中用规则或轻量模型提取(海边、通勤、约会)
  4. 文本拼接 → 按优先级拼接,总长度不超过 512 tokens
  5. 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_constant20值越小,排名靠后的文档影响力越大;20 是多数场景的最佳起点
rank_window_size50~100从每个子检索器取多少条参与融合,必须 >= from + size
k(kNN)50向量检索返回候选数
num_candidates100~500kNN 内部候选数,建议 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 元
MySQL2C 8G~ 300 元
Redis(价格/库存/向量缓存)2C 8G~ 200 元
应用服务器2 节点 x 4C 8G~ 800 元
合计~ 5800 元/月

5. 风险提示

  1. 动静分离后,动态索引的 merge 压力仍需关注。虽然静态索引几乎不变,但 product_dynamic 高频更新仍会产生大量 segment。建议监控段数,超过 20 个时手动触发 forcemerge(低峰期执行)。
  2. RRF 不支持高亮、rescore、sort、scroll。如果搜索列表需要关键词高亮,需在 RRF 检索后对返回的文档单独发高亮请求,或在前端用前端分词高亮兜底。
  3. 同义词表需要持续维护,否则会有冷门 query 漏掉。初始阶段人工标注 Top 500 高频 query 就够了。长尾 query 会持续出现,建议每周从零结果日志中捞一批新词做补充。如果团队抽不出人手维护,那就别做规则同义词,直接依赖 Dense Vector 兜底也行。
  4. BGE 模型对电商专有名词可能犯傻A字裙JK制服Lolita 这些词在通用中文 Embedding 上的表现不确定。上线前拿商品标题跑一遍 Batch Embedding,人工抽查 Top 10 高频类目的向量相似度排序是否合理。如果有明显偏差,用几百条商品标题做 LoRA 微调,成本不高但收益明显。
  5. 熔断降级会损失语义召回能力。BGE 熔断后降级到 BM25 纯文本检索,长尾语义 query(如"海边拍照穿的裙子")的召回效果会下降。建议监控降级率,持续 > 5% 时需扩容 BGE 服务或优化推理延迟。
  6. 分页约束。RRF 要求 from + size <= rank_window_size,深度分页必须用 search_after。前端分页设计时需特别注意,避免翻到后面出现空结果。