Java 程序员第 44 阶段05:大模型微服务拆分,独立服务解耦便于扩容维护,Embedding向量服务独立部署与接口设计
[Embedding 服务独立化的核心动因](#1-embedding-服务独立化的核心动因)- [服务接口设计:单条、批量与流式](#2-服务接口设计单条批量与流式)
- [多模型抽象与动态切换](#3-多模型抽象与动态切换)
- [批量向量化与性能优化](#4-批量向量化与性能优化)
- [向量入库全流程设计](#5-向量入库全流程设计)
- [与向量检索服务的协作](#6-与向量检索服务的协作)
- [缓存与去重:避免重复计算](#7-缓存与去重避免重复计算)
- [独立部署与弹性扩容](#8-独立部署与弹性扩容)
- [监控指标与运维要点](#9-监控指标与运维要点)
- [系列总结与全链路回顾](#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 最佳实践总结
回顾整个系列,大模型微服务拆分的十条最佳实践:
- **先有方法论再动手**:不要拍脑袋拆,用 DDD + 三维度评估 + 数据流分析做决策。
- **绞杀者模式渐进拆分**:灰度切换、保留兜底、随时回退。
- **能力层服务不互调**:通过编排层串联,避免网状依赖。
- **推理服务快速失败优于漫长等待**:分层超时 + 降级兜底。
- **SSE 必须全链路透传**:网关不缓冲、编排用 Flux、推理正确解析。
- **Prompt 资产化管理**:脱离代码发布周期、版本可追溯、AB 可度量。
- **Embedding 微批处理**:攒批是吞吐优化的关键、在线离线分流。
- **不共享数据库**:这是微服务松耦合的底线。
- **接口先行契约保护**:OpenAPI 定义 + ArchUnit 约束 + CDC 测试。
- **可观测性先行**:链路追踪 + 指标监控 + 结构化日志,没有可观测性就没有微服务。
10.5 后续学习方向
本系列聚焦拆分方法论和核心服务设计,以下方向值得进一步深入:
- **向量检索服务**的深度优化:混合检索(向量+关键词)、重排序、分片策略。
- **知识库管理服务**的完整实现:文档解析(PDF/Word/HTML)、切分策略、增量更新。
- **大模型网关**的高级能力:多租户隔离、配额管理、内容安全审计。
- **成本优化**:token 预算管理、模型分级路由(简单问题用小模型、复杂问题用大模型)。
- **可观测性体系**:全链路追踪在大模型场景的适配(流式调用的 trace 特殊处理)。
希望这五篇文章能为 Java 开发者在大模型微服务架构实践中提供有价值的参考。微服务拆分不是目的,目的是让系统更灵活、更稳定、更易维护。带着这个目标去拆,才能拆得对、拆得好。
> 本文是"Java 程序员第 44 阶段"系列的第 05 篇(完结篇),聚焦 Embedding 向量服务的独立部署与接口设计。建议结合前四篇系统性阅读。
更多推荐

所有评论(0)