1. [Embedding 服务独立化的核心动因](#1-embedding-服务独立化的核心动因)
  2. [服务接口设计:单条、批量与流式](#2-服务接口设计单条批量与流式)
  3. [多模型抽象与动态切换](#3-多模型抽象与动态切换)
  4. [批量向量化与性能优化](#4-批量向量化与性能优化)
  5. [向量入库全流程设计](#5-向量入库全流程设计)
  6. [与向量检索服务的协作](#6-与向量检索服务的协作)
  7. [缓存与去重:避免重复计算](#7-缓存与去重避免重复计算)
  8. [独立部署与弹性扩容](#8-独立部署与弹性扩容)
  9. [监控指标与运维要点](#9-监控指标与运维要点)
  10. [系列总结与全链路回顾](#10-系列总结与全链路回顾)

1.1 Embedding 在大模型应用中的角色

Embedding(向量化)是把文本转换为高维向量的过程,它是 RAG(检索增强生成)、语义搜索、推荐召回等场景的基础能力。在完整的 AI 应用链路中,Embedding 出现在两个环节:入库环节(把文档片段向量化后存入向量库)和检索环节(把用户查询向量化后在向量库中搜索相似项)。这两个环节对 Embedding 服务的需求特征不同——入库是批量、离线、吞吐优先;检索是单条、在线、延迟优先。

1.2 为什么要独立成服务

在单体架构中,Embedding 逻辑通常以库的形式被直接调用。这种方式的问题随着业务规模增长逐渐暴露。

**资源争抢**是最直接的动因。Embedding 模型(如 bge-large-zh、text-embedding-3-large)在推理时需要较多计算资源,尤其是用 GPU 加速时。如果 Embedding 和业务逻辑在同一进程,Embedding 的计算会占用大量 CPU/GPU,拖慢业务请求。独立部署后,Embedding 服务可以选择带 GPU 的机型,业务服务用普通 CPU 机型,各取所需。

**模型切换成本**是第二个动因。Embedding 模型迭代很快,从早期的 text-embedding-ada-002 到 bge系列再到当前的各类最优模型,几乎每几个月就有更好的选择。如果 Embedding 逻辑散落在各处,换模型意味着全面改造。独立成服务后,只需修改服务内部的模型配置,调用方无感知。

**向量一致性**是第三个动因。入库和检索必须用同一个 Embedding 模型,否则向量维度不匹配、语义空间不一致,检索效果会严重退化。独立服务天然保证模型一致性——所有向量化请求都走同一个服务,不会出现"入库用 A 模型、检索用 B 模型"的错误。

1.3 与推理服务的区别

Embedding 服务和推理服务虽然都调用模型,但特征截然不同。推理服务是"低并发、长耗时、流式输出",Embedding 服务是"高并发、短耗时、批量处理"。两者的扩容策略、容错策略、接口设计都不同,因此必须分别独立部署,不能合并。

2.1 三种接口模式

Embedding 服务提供三种接口模式,覆盖不同场景。

**单条接口**用于在线检索场景:用户输入一个查询,实时向量化后在向量库中搜索。延迟敏感,要求在 50-200ms 内返回。

**批量接口**用于离线入库场景:知识库文档切分成片段后,批量向量化。吞吐敏感,一次请求可能包含数百到数千条文本,要求整体处理时间可控。

**异步任务接口**用于超大批量场景:全量重建向量库时,可能涉及百万级文档,同步处理不现实。提交一个异步任务,服务后台处理,完成后回调通知。

2.2 单条接口实现

@RestController
@RequestMapping("/v1/embeddings")
public class EmbeddingController {
    @Autowired
    private EmbeddingService embeddingService;
    @Autowired
    private EmbeddingMetrics metrics;
    @PostMapping("/single")
    public ApiResult<EmbeddingResponse> embed(
            @RequestBody @Valid EmbeddingRequest request) {
        long start = System.currentTimeMillis();
        try {
            EmbeddingResponse response = embeddingService.embed(request);
            metrics.recordLatency("single", System.currentTimeMillis() - start, "success");
            return ApiResult.success(response);
        } catch (Exception e) {
            metrics.recordLatency("single", System.currentTimeMillis() - start, "error");
            return ApiResult.error(500, "向量化失败: " + e.getMessage());
        }
    }
}
@Data
public class EmbeddingRequest {
    @NotBlank
    private String text;        // 待向量化的文本
    private String model;       // 模型标识,可选,默认用配置的模型
    private boolean normalize;  // 是否归一化,默认true
}
@Data
@Builder
public class EmbeddingResponse {
    private List<Float> vector;     // 向量
    private int dimension;          // 维度
    private String model;           // 实际使用的模型
    private int tokenCount;         // token数
}

2.3 批量接口实现

批量接口是 Embedding 服务最重要的接口,入库场景几乎都走批量。关键设计点是:接收一个文本列表,返回对应的向量列表,保持顺序一致。

@PostMapping("/batch")
public ApiResult<BatchEmbeddingResponse> embedBatch(
        @RequestBody @Valid BatchEmbeddingRequest request) {
    if (request.getTexts().size() > 500) {
        return ApiResult.error(400, "单次批量不超过500条");
    }
    long start = System.currentTimeMillis();
    try {
        BatchEmbeddingResponse response = embeddingService.embedBatch(request);
        metrics.recordLatency("batch", System.currentTimeMillis() - start, "success");
        metrics.recordBatchSize(request.getTexts().size());
        return ApiResult.success(response);
    } catch (Exception e) {
        metrics.recordLatency("batch", System.currentTimeMillis() - start, "error");
        return ApiResult.error(500, "批量向量化失败: " + e.getMessage());
    }
}
@Data
public class BatchEmbeddingRequest {
    @NotEmpty
    @Size(max = 500)
    private List<String> texts;
    private String model;
    private boolean normalize = true;
}
@Data
@Builder
public class BatchEmbeddingResponse {
    private List<List<Float>> vectors;  // 与输入texts顺序一致
    private int dimension;
    private String model;
    private int totalTokens;
}

2.4 异步任务接口

对于百万级的全量重建,提供异步任务接口。

@PostMapping("/async-job")
public ApiResult<JobResponse> submitAsyncJob(
        @RequestBody AsyncEmbeddingRequest request) {
    String jobId = embeddingService.submitAsyncJob(request);
    return ApiResult.success(JobResponse.builder()
        .jobId(jobId)
        .status("QUEUED")
        .totalItems(request.getTexts().size())
        .build());
}
@GetMapping("/async-job/{jobId}/status")
public ApiResult<JobStatus> getJobStatus(@PathVariable String jobId) {
    return ApiResult.success(embeddingService.getJobStatus(jobId));
}

异步任务的状态流转为:QUEUED -> PROCESSING -> COMPLETED/FAILED。任务结果写入对象存储,通过回调 URL 通知调用方。

3.1 模型抽象层

和推理服务类似,Embedding 服务也需要多模型抽象层。不同 Embedding 模型的 API 格式差异比 LLM 小(大多数都接受文本列表、返回向量列表),但维度、归一化方式、token 限制各不相同。抽象层统一这些差异。

public interface EmbeddingProvider {
    String modelId();
    int dimension();
    List<float[]> embed(List<String> texts, boolean normalize);
    boolean supportsBatch();
    int maxBatchSize();
}

3.2 本地 ONNX 模型实现

自建 Embedding 服务通常用 ONNX Runtime 加载模型,避免依赖外部 API。下面是本地 ONNX provider 的实现。

@Component
@Slf4j
public class OnnxEmbeddingProvider implements EmbeddingProvider {
    private final String modelId;
    private final int dimension;
    private final OrtSession session;
    private final Tokenizer tokenizer;
    private final int maxBatchSize;
    public OnnxEmbeddingProvider(OnnxModelConfig config) {
        this.modelId = config.getModelId();
        this.dimension = config.getDimension();
        this.maxBatchSize = config.getMaxBatchSize();
        try {
            OrtEnvironment env = OrtEnvironment.getEnvironment();
            this.session = env.createSession(config.getModelPath(),
                new OrtSession.SessionOptions()
                    .setOptimizationLevel(OrtSession.SessionOptions.OptLevel.ALL_OPT));
            this.tokenizer = TokenizerFactory.create(config.getTokenizerPath());
            log.info("ONNX Embedding模型已加载: {}, 维度={}", modelId, dimension);
        } catch (Exception e) {
            throw new RuntimeException("模型加载失败: " + modelId, e);
        }
    }
    @Override
    public List<float[]> embed(List<String> texts, boolean normalize) {
        try {
            // 1. tokenize
            long[][] inputIds = new long[texts.size()][];
            long[][] attentionMask = new long[texts.size()][];
            for (int i = 0; i < texts.size(); i++) {
                Encoding enc = tokenizer.encode(texts.get(i));
                inputIds[i] = enc.getIds();
                attentionMask[i] = enc.getAttentionMask();
            }
            // 2. 推理
            OnnxTensor inputTensor = OnnxTensor.createTensor(session.getInputNames().get(0), inputIds);
            OnnxTensor maskTensor = OnnxTensor.createTensor(session.getInputNames().get(1), attentionMask);
            OrtSession.Result result = session.run(Map.of(
                session.getInputNames().get(0), inputTensor,
                session.getInputNames().get(1), maskTensor));
            // 3. 提取向量 (mean pooling + normalize)
            float[][] raw = (float[][]) result.get(0).getValue();
            List<float[]> vectors = new ArrayList<>(texts.size());
            for (int i = 0; i < texts.size(); i++) {
                float[] pooled = meanPool(raw[i], attentionMask[i]);
                if (normalize) {
                    pooled = normalize(pooled);
                }
                vectors.add(pooled);
            }
            return vectors;
        } catch (Exception e) {
            throw new EmbeddingException("ONNX推理失败", e);
        }
    }
    @Override
    public boolean supportsBatch() { return true; }
    @Override
    public int maxBatchSize() { return maxBatchSize; }
    private float[] meanPool(float[] tokenVectors, long[] mask) {
        int dim = dimension;
        float[] pooled = new float[dim];
        int count = 0;
        for (int t = 0; t < mask.length; t++) {
            if (mask[t] == 1) {
                for (int d = 0; d < dim; d++) {
                    pooled[d] += tokenVectors[t * dim + d];
                }
                count++;
            }
        }
        for (int d = 0; d < dim; d++) {
            pooled[d] /= count;
        }
        return pooled;
    }
    private float[] normalize(float[] vector) {
        float norm = 0;
        for (float v : vector) norm += v * v;
        norm = (float) Math.sqrt(norm);
        if (norm == 0) return vector;
        for (int i = 0; i < vector.length; i++) vector[i] /= norm;
        return vector;
    }
}

3.3 云 API 实现

对于使用 OpenAI 或其他云厂商 Embedding API 的场景,实现一个云 API provider。

@Component
public class CloudEmbeddingProvider implements EmbeddingProvider {
    private final String modelId;
    private final int dimension;
    private final WebClient webClient;
    private final int maxBatchSize;
    public CloudEmbeddingProvider(CloudModelConfig config) {
        this.modelId = config.getModelId();
        this.dimension = config.getDimension();
        this.maxBatchSize = config.getMaxBatchSize();
        this.webClient = WebClient.builder()
            .baseUrl(config.getApiUrl())
            .defaultHeader("Authorization", "Bearer " + config.getApiKey())
            .codecs(c -> c.defaultCodecs().maxInMemorySize(32 * 1024 * 1024))
            .build();
    }
    @Override
    public List<float[]> embed(List<String> texts, boolean normalize) {
        Map<String, Object> body = Map.of(
            "model", modelId,
            "input", texts
        );
        Map<String, Object> resp = webClient.post()
            .uri("/v1/embeddings")
            .bodyValue(body)
            .retrieve()
            .bodyToMono(Map.class)
            .timeout(Duration.ofSeconds(30))
            .block();
        List<?> data = (List<?>) resp.get("data");
        List<float[]> vectors = new ArrayList<>(data.size());
        // 按index排序确保顺序一致
        data.stream()
            .map(d -> (Map<String, Object>) d)
            .sorted(Comparator.comparingInt(d -> (Integer) d.get("index")))
            .forEach(d -> {
                List<Number> embedding = (List<Number>) d.get("embedding");
                float[] vec = new float[embedding.size()];
                for (int i = 0; i < embedding.size(); i++) {
                    vec[i] = embedding.get(i).floatValue();
                }
                if (normalize) vec = normalize(vec);
                vectors.add(vec);
            });
        return vectors;
    }
}

3.4 动态切换

模型配置存在 Nacos,支持热切换。切换模型时需要注意:已有向量的维度可能与新模型不同,不能直接混用。因此模型切换通常配合全量重建。

@Component
public class EmbeddingModelManager {
    private volatile EmbeddingProvider currentProvider;
    @NacosConfigListener(dataId = "embedding-model-config.json")
    public void onConfigUpdate(String config) {
        EmbeddingModelConfig mc = JsonUtils.parse(config, EmbeddingModelConfig.class);
        EmbeddingProvider newProvider = createProvider(mc);
        // 预热新模型
        newProvider.embed(List.of("预热文本"), true);
        this.currentProvider = newProvider;
        log.info("Embedding模型已切换: {}, 维度={}",
            mc.getModelId(), mc.getDimension());
    }
    public EmbeddingProvider getProvider() {
        return currentProvider;
    }
}

4.1 批处理的重要性

Embedding 模型的推理效率在批处理时远高于单条处理。GPU 的并行计算能力在批量输入时才能充分利用——单条向量化可能耗时 20ms,批量 32 条可能只耗时 50ms,单条均摊仅 1.5ms。因此 Embedding 服务的核心性能优化就是批处理。

4.2 微批处理队列

在线场景下,请求是逐个到达的,如果每个请求都立即推理,无法利用批处理优势。解决方案是微批处理(Micro-batching):把短时间内到达的多个请求攒成一个小批次,一次推理后分发结果。

@Component
public class MicroBatchProcessor {
    private final EmbeddingProvider provider;
    private final int maxBatchSize;
    private final Duration maxWaitTime;
    private final BlockingQueue<EmbeddingTask> queue;
    @PostConstruct
    public void start() {
        Thread.ofVirtual().start(this::batchLoop);
    }
    private void batchLoop() {
        while (!Thread.currentThread().isInterrupted()) {
            List<EmbeddingTask> batch = collectBatch();
            if (batch.isEmpty()) continue;
            try {
                List<String> texts = batch.stream()
                    .map(EmbeddingTask::getText).toList();
                List<float[]> vectors = provider.embed(texts, true);
                // 分发结果
                for (int i = 0; i < batch.size(); i++) {
                    batch.get(i).complete(vectors.get(i));
                }
            } catch (Exception e) {
                batch.forEach(t -> t.completeExceptionally(e));
            }
        }
    }
    private List<EmbeddingTask> collectBatch() {
        List<EmbeddingTask> batch = new ArrayList<>(maxBatchSize);
        long deadline = System.nanoTime() + maxWaitTime.toNanos();
        try {
            // 阻塞等待第一条
            EmbeddingTask first = queue.poll(maxWaitTime.toMillis(), TimeUnit.MILLISECONDS);
            if (first != null) batch.add(first);
            // 非阻塞收集更多,直到满批或超时
            while (batch.size() < maxBatchSize
                   && System.nanoTime() < deadline) {
                EmbeddingTask t = queue.poll();
                if (t == null) break;
                batch.add(t);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        return batch;
    }
    public CompletableFuture<float[]> submit(String text) {
        EmbeddingTask task = new EmbeddingTask(text);
        queue.offer(task);
        return task.future;
    }
}

关键参数是 `maxBatchSize`(通常 32-64)和 `maxWaitTime`(通常 10-50ms)。批大小越大吞吐越高但延迟越高;等待时间越长攒批越满但单条延迟越高。需要根据场景权衡——在线检索场景偏向小批短等待,离线入库场景偏向大批长等待。

4.3 长文本截断

Embedding 模型有最大 token 限制(通常 512),超长文本需要截断。但粗暴截断会丢失信息,更好的做法是分段向量化后取平均。

@Component
public class TextTruncator {
    @Autowired
    private TokenEstimator tokenEstimator;
    private static final int MAX_TOKENS = 512;
    public List<String> splitForEmbedding(String text) {
        int tokenCount = tokenEstimator.estimate(text);
        if (tokenCount <= MAX_TOKENS) {
            return List.of(text);
        }
        // 按段落切分,每段不超过MAX_TOKENS
        List<String> segments = new ArrayList<>();
        String[] paragraphs = text.split("\n\n");
        StringBuilder current = new StringBuilder();
        int currentTokens = 0;
        for (String para : paragraphs) {
            int paraTokens = tokenEstimator.estimate(para);
            if (currentTokens + paraTokens > MAX_TOKENS) {
                if (current.length() > 0) {
                    segments.add(current.toString());
                    current = new StringBuilder();
                    currentTokens = 0;
                }
                // 单段就超限,硬截断
                if (paraTokens > MAX_TOKENS) {
                    segments.add(hardTruncate(para, MAX_TOKENS));
                } else {
                    current.append(para);
                    currentTokens = paraTokens;
                }
            } else {
                current.append(para).append("\n\n");
                currentTokens += paraTokens;
            }
        }
        if (current.length() > 0) segments.add(current.toString());
        return segments;
    }
    private String hardTruncate(String text, int maxTokens) {
        // 近似截断:按字符比例
        int ratio = maxTokens * 4; // 粗略1 token ≈ 4字符(中文偏少)
        return text.substring(0, Math.min(text.length(), ratio));
    }
}

5.1 入库流程概览

知识库文档入库是 Embedding 服务最典型的使用场景。完整流程包括:文档上传 -> 文本提取 -> 文本切分 -> 向量化 -> 写入向量库 -> 更新元信息。Embedding 服务在这个流程中承担"向量化"环节,但它通常也编排整个流程(因为切分和向量化紧密耦合)。

5.2 文本切分策略

切分质量直接影响检索效果。切得太碎,语义不完整;切得太长,向量稀释了关键信息。常用策略是"按段落切分 + 重叠窗口"。

@Component
public class TextSplitter {
    private static final int CHUNK_SIZE = 500;     // 每段约500字符
    private static final int OVERLAP = 100;        // 重叠100字符
    public List<TextChunk> split(String text, String documentId) {
        List<TextChunk> chunks = new ArrayList<>();
        // 先按段落预切
        String[] paragraphs = text.split("(?<=\n)");
        StringBuilder buffer = new StringBuilder();
        int chunkIndex = 0;
        for (String para : paragraphs) {
            buffer.append(para);
            while (buffer.length() >= CHUNK_SIZE) {
                String chunkText = buffer.substring(0, CHUNK_SIZE);
                chunks.add(TextChunk.builder()
                    .documentId(documentId)
                    .chunkIndex(chunkIndex++)
                    .content(chunkText)
                    .build());
                // 保留重叠部分
                buffer = new StringBuilder(
                    buffer.substring(CHUNK_SIZE - OVERLAP));
            }
        }
        // 处理剩余
        if (buffer.length() > 50) {
            chunks.add(TextChunk.builder()
                .documentId(documentId)
                .chunkIndex(chunkIndex)
                .content(buffer.toString())
                .build());
        }
        return chunks;
    }
}

5.3 入库编排

@Service
public class IngestionService {
    @Autowired
    private TextSplitter splitter;
    @Autowired
    private EmbeddingService embeddingService;
    @Autowired
    private VectorSearchClient vectorClient;
    @Autowired
    private DocumentMetadataRepository metadataRepo;
    @Async("ingestionExecutor")
    public CompletableFuture<IngestionResult> ingest(Document doc) {
        try {
            // 1. 切分
            List<TextChunk> chunks = splitter.split(doc.getContent(), doc.getId());
            log.info("文档{}切分为{}个片段", doc.getId(), chunks.size());
            // 2. 批量向量化
            List<String> texts = chunks.stream().map(TextChunk::getContent).toList();
            BatchEmbeddingResponse embedResp = embeddingService.embedBatch(
                new BatchEmbeddingRequest(texts, null, true));
            // 3. 写入向量库
            List<VectorRecord> records = new ArrayList<>();
            for (int i = 0; i < chunks.size(); i++) {
                records.add(VectorRecord.builder()
                    .id(doc.getId() + "_" + i)
                    .vector(embedResp.getVectors().get(i))
                    .metadata(Map.of(
                        "documentId", doc.getId(),
                        "chunkIndex", i,
                        "content", chunks.get(i).getContent()
                    ))
                    .build());
            }
            vectorClient.batchInsert(records);
            // 4. 更新文档元信息
            doc.setStatus("INDEXED");
            doc.setChunkCount(chunks.size());
            metadataRepo.save(doc);
            return CompletableFuture.completedFuture(
                IngestionResult.success(doc.getId(), chunks.size()));
        } catch (Exception e) {
            doc.setStatus("FAILED");
            metadataRepo.save(doc);
            log.error("文档入库失败: " + doc.getId(), e);
            return CompletableFuture.completedFuture(
                IngestionResult.failed(doc.getId(), e.getMessage()));
        }
    }
}

5.4 增量更新与删除

文档更新时,需要先删除旧向量再写入新向量。删除按 documentId 过滤即可。这个流程虽然简单,但必须保证原子性——不能出现"旧向量删了、新向量没写入"的中间状态。用本地消息表保证最终一致性(见第01篇)。

6.1 职责划分

Embedding 服务和向量检索服务虽然都涉及向量,但职责不同。Embedding 服务负责"文本到向量"的转换,向量检索服务负责"向量到相似向量"的查询。两者通过 API 协作:检索服务收到查询请求后,先调 Embedding 服务把查询文本转向量,再在自己的向量库中搜索。

6.2 查询链路

一次完整的语义检索链路如下:客户端 -> 检索服务 -> Embedding 服务(查询向量化)-> 检索服务(向量搜索)-> 客户端。这条链路有两个同步调用,延迟叠加。优化方式是把查询向量化缓存——相同查询的向量可以复用,避免重复计算。

@Service
public class SemanticSearchService {
    @Autowired
    private EmbeddingClient embeddingClient;
    @Autowired
    private VectorRepository vectorRepo;
    @Autowired
    private Cache<String, float[]> queryVectorCache;
    public List<SearchResult> search(String query, int topK) {
        // 1. 查询向量化(带缓存)
        float[] queryVector = queryVectorCache.get(query, k -> {
            EmbeddingResponse resp = embeddingClient.embed(
                new EmbeddingRequest(query, null, true));
            return toPrimitiveArray(resp.getVector());
        });
        // 2. 向量检索
        return vectorRepo.search(queryVector, topK);
    }
}

6.3 维度一致性保障

入库和检索必须用同一个 Embedding 模型,否则维度不匹配。保障措施是在 Embedding 服务的响应中带上 modelId 和 dimension,检索服务在写入和查询时都校验维度一致性。如果检测到维度不匹配(通常是模型切换后未重建向量库),直接拒绝并告警。

@Component
public class DimensionValidator {
    private final int expectedDimension;
    public void validate(float[] vector) {
        if (vector.length != expectedDimension) {
            throw new DimensionMismatchException(
                String.format("向量维度%d与预期%d不匹配,可能模型已切换,请重建向量库",
                    vector.length, expectedDimension));
        }
    }
}

7.1 内容哈希缓存

Embedding 计算是有成本的(尤其用云 API 时按量计费)。相同文本的向量结果是确定性的(相同模型 + 相同输入 = 相同输出),因此可以缓存。缓存 key 用文本内容的哈希值,缓存 value 是向量。

@Component
public class EmbeddingCache {
    @Autowired
    private RedisTemplate<String, String> redis;
    @Autowired
    private Cache<String, float[]> localCache;
    private static final Duration TTL = Duration.ofDays(7);
    public float[] getOrCompute(String text, String modelId,
                                 Supplier<float[]> computer) {
        String key = "emb:" + modelId + ":" + hash(text);
        // 1. 本地缓存
        float[] cached = localCache.getIfPresent(key);
        if (cached != null) return cached;
        // 2. Redis
        String json = redis.opsForValue().get(key);
        if (json != null) {
            float[] vec = JsonUtils.parse(json, float[].class);
            localCache.put(key, vec);
            return vec;
        }
        // 3. 计算
        float[] vector = computer.get();
        // 4. 回填
        localCache.put(key, vector);
        redis.opsForValue().set(key, JsonUtils.toJson(vector), TTL);
        return vector;
    }
    public Map<String, float[]> batchGetOrCompute(List<String> texts,
            String modelId, Function<List<String>, List<float[]>> batchComputer) {
        Map<String, float[]> result = new HashMap<>();
        List<String> missed = new ArrayList<>();
        // 批量查缓存
        for (String text : texts) {
            String key = "emb:" + modelId + ":" + hash(text);
            float[] cached = localCache.getIfPresent(key);
            if (cached != null) {
                result.put(text, cached);
            } else {
                missed.add(text);
            }
        }
        // 批量计算未命中的
        if (!missed.isEmpty()) {
            List<float[]> computed = batchComputer.apply(missed);
            for (int i = 0; i < missed.size(); i++) {
                String text = missed.get(i);
                float[] vec = computed.get(i);
                result.put(text, vec);
                String key = "emb:" + modelId + ":" + hash(text);
                localCache.put(key, vec);
                redis.opsForValue().set(key, JsonUtils.toJson(vec), TTL);
            }
        }
        return result;
    }
    private String hash(String text) {
        return DigestUtils.md5DigestAsHex(text.getBytes(StandardCharsets.UTF_8));
    }
}

7.2 缓存命中率监控

缓存命中率是 Embedding 服务的重要成本指标。命中率越高,实际计算的请求越少,成本越低。通过 Micrometer 暴露命中率指标。

@Scheduled(fixedRate = 5000)
public void reportCacheMetrics() {
    long hits = cacheHits.get();
    long misses = cacheMisses.get();
    double rate = (hits + misses) > 0 ? (double) hits / (hits + misses) : 0;
    meterRegistry.gauge("embedding.cache.hit_rate", rate);
}

8.1 部署架构

Embedding 服务独立部署在带 GPU 的节点上(如果用本地模型)或普通 CPU 节点上(如果用云 API)。部署架构包括:多个服务实例(无状态)、一个负载均衡器、Redis 缓存集群、Nacos 配置中心。服务实例无状态,会话级状态(如微批处理队列)在每个实例内部维护,实例故障不影响全局。

8.2 GPU 节点扩容

使用本地 ONNX 模型时,扩容依赖 GPU 资源。Kubernetes 中用 GPU 资源请求和节点亲和性确保实例调度到带 GPU 的节点。

apiVersion: apps/v1
kind: Deployment
metadata:
  name: embedding-service
spec:
  replicas: 3
  selector:
    matchLabels:
      app: embedding-service
  template:
    metadata:
      labels:
        app: embedding-service
    spec:
      nodeSelector:
        gpu-type: "nvidia-a10"     # 调度到A10 GPU节点
      containers:
      - name: embedding
        image: embedding-service:1.0.0
        resources:
          limits:
            nvidia.com/gpu: 1      # 每实例1块GPU
            memory: "8Gi"
          requests:
            nvidia.com/gpu: 1
            memory: "4Gi"
        readinessProbe:
          httpGet:
            path: /actuator/readiness
            port: 8080
          initialDelaySeconds: 30   # 模型加载需要时间
          periodSeconds: 10

8.3 混合部署策略

对于成本敏感的场景,可以采用混合部署:在线检索请求用本地 GPU 模型(低延迟),离线入库请求用云 API(按量计费、无需维护 GPU)。通过路由器按请求类型分流。

@Component
public class HybridEmbeddingRouter {
    @Autowired
    @Qualifier("localOnnxProvider")
    private EmbeddingProvider localProvider;
    @Autowired
    @Qualifier("cloudApiProvider")
    private EmbeddingProvider cloudProvider;
    public EmbeddingProvider route(EmbeddingRequest request) {
        // 在线单条请求:用本地模型(低延迟)
        if (request.getSource() == RequestSource.ONLINE_RETRIEVAL) {
            return localProvider;
        }
        // 离线批量请求:用云API(弹性扩展)
        if (request.getSource() == RequestSource.BATCH_INGESTION) {
            return cloudProvider;
        }
        // 默认本地
        return localProvider;
    }
}

9.1 核心监控指标

Embedding 服务的监控指标分为四类:性能指标(延迟、吞吐)、资源指标(GPU 利用率、内存)、业务指标(向量化的文本数、token 数)、成本指标(缓存命中率、云 API 调用量)。

指标

说明

告警阈值

------

------

---------

embedding_latency

单条向量化延迟

P99 > 200ms

batch_latency

批量向量化延迟

P99 > 5s

gpu_utilization

GPU利用率

> 85% 持续5分钟

cache_hit_rate

缓存命中率

< 30%(异常低)

batch_efficiency

批处理效率(实际批次/最大批次)

< 0.3(攒批不足)

error_rate

错误率

> 5%

queue_depth

微批处理队列深度

持续增长

9.2 指标暴露

@Component
public class EmbeddingMetrics {
    @Autowired
    private MeterRegistry meterRegistry;
    public void recordLatency(String type, long latencyMs, String outcome) {
        Timer.builder("embedding.latency")
            .tag("type", type)          // single / batch
            .tag("outcome", outcome)    // success / error
            .register(meterRegistry)
            .record(latencyMs, TimeUnit.MILLISECONDS);
    }
    public void recordBatchSize(int size) {
        meterRegistry.summary("embedding.batch.size")
            .record(size);
    }
    public void recordTokens(int count) {
        meterRegistry.counter("embedding.tokens", "type", "processed")
            .increment(count);
    }
    public void recordCache(String result) {
        meterRegistry.counter("embedding.cache", "result", result)
            .increment();
    }
}

9.3 常见运维问题

**模型加载失败**:检查模型文件路径、ONNX Runtime 版本兼容性、GPU 驱动是否正常。

**延迟突增**:通常是 GPU 利用率打满或微批处理队列积压。检查是否有突发批量请求挤占了在线请求资源,必要时分离在线和离线的实例组。

**维度不匹配**:模型切换后未重建向量库。必须执行全量重建,期间用旧模型的服务实例继续服务。

**缓存命中率下降**:可能是查询模式变化(新查询增多)。命中率下降不一定异常,但需要关注成本变化。

10.1 Embedding 服务核心要点

Embedding 向量服务独立化的核心收益是资源隔离、模型统一、弹性扩容。通过单条/批量/异步三种接口模式覆盖在线和离线场景,通过多模型抽象层支持动态切换,通过微批处理最大化吞吐,通过多级缓存避免重复计算,通过维度校验保障入库检索一致性。

10.2 五篇系列全链路回顾

本系列五篇文章构成了一套完整的大模型微服务拆分方法论:

**第 01 篇**确立了总体设计思路:四层架构(接入层/编排层/能力层/数据层)、五大拆分原则、绞杀者模式四阶段演进路线。核心是回答"为什么要拆"和"怎么拆"的方向问题。

**第 02 篇**解决了最难的边界识别问题:用 DDD 事件风暴梳理领域、用数据流分析验证边界、用三维度评估模型(业务内聚度/运行时独立性/团队所有权)量化粒度、用变更频率指导发布节奏。核心是回答"在哪里下刀"的判断问题。

**第 03 篇**实战了 AI 推理服务的独立化:多模型路由抽象、SSE 全链路透传、分层超时+熔断降级容错、token 计量、请求排队背压、模型预热。核心是"低并发长耗时流式服务"的工程实现。

**第 04 篇**实战了 Prompt 管理服务的设计:模板存储+Mustache 渲染引擎、版本管理+灰度发布、AB 测试框架、多级缓存、回归测试、可视化后台。核心是"高频变更的软资产如何资产化管理"。

**第 05 篇**(本篇)实战了 Embedding 向量服务的设计:单条/批量/异步三模式接口、多模型抽象、微批处理性能优化、向量入库全流程、缓存去重、独立部署扩容。核心是"高并发短耗时批处理服务"的工程实现。

10.3 拆分后的完整架构

完成五篇的全部拆分后,系统从单体演进为以下微服务集群:

服务

职责

扩容依据

核心技术

------

------

---------

---------

business-gateway

接入/鉴权/限流

QPS

Gateway

orchestration-service

业务编排/会话管理

QPS

Spring Boot

ai-inference-service

大模型推理

队列深度

Spring Boot + vLLM

prompt-service

Prompt 模板管理

QPS

Spring Boot + MySQL

embedding-service

文本向量化

QPS+GPU

Spring Boot + ONNX

vector-search-service

向量检索

QPS

Spring Boot + Milvus

knowledge-service

知识库管理

业务量

Spring Boot + MySQL

每个服务独立部署、独立扩容、独立发布、独立容错,通过标准 API 协作。这正是"独立服务解耦便于扩容维护"的最终落地。

10.4 最佳实践总结

回顾整个系列,大模型微服务拆分的十条最佳实践:

  1. **先有方法论再动手**:不要拍脑袋拆,用 DDD + 三维度评估 + 数据流分析做决策。
  2. **绞杀者模式渐进拆分**:灰度切换、保留兜底、随时回退。
  3. **能力层服务不互调**:通过编排层串联,避免网状依赖。
  4. **推理服务快速失败优于漫长等待**:分层超时 + 降级兜底。
  5. **SSE 必须全链路透传**:网关不缓冲、编排用 Flux、推理正确解析。
  6. **Prompt 资产化管理**:脱离代码发布周期、版本可追溯、AB 可度量。
  7. **Embedding 微批处理**:攒批是吞吐优化的关键、在线离线分流。
  8. **不共享数据库**:这是微服务松耦合的底线。
  9. **接口先行契约保护**:OpenAPI 定义 + ArchUnit 约束 + CDC 测试。
  10. **可观测性先行**:链路追踪 + 指标监控 + 结构化日志,没有可观测性就没有微服务。

10.5 后续学习方向

本系列聚焦拆分方法论和核心服务设计,以下方向值得进一步深入:

  • **向量检索服务**的深度优化:混合检索(向量+关键词)、重排序、分片策略。
  • **知识库管理服务**的完整实现:文档解析(PDF/Word/HTML)、切分策略、增量更新。
  • **大模型网关**的高级能力:多租户隔离、配额管理、内容安全审计。
  • **成本优化**:token 预算管理、模型分级路由(简单问题用小模型、复杂问题用大模型)。
  • **可观测性体系**:全链路追踪在大模型场景的适配(流式调用的 trace 特殊处理)。

希望这五篇文章能为 Java 开发者在大模型微服务架构实践中提供有价值的参考。微服务拆分不是目的,目的是让系统更灵活、更稳定、更易维护。带着这个目标去拆,才能拆得对、拆得好。

> 本文是"Java 程序员第 44 阶段"系列的第 05 篇(完结篇),聚焦 Embedding 向量服务的独立部署与接口设计。建议结合前四篇系统性阅读。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐