知识库 5 分钟同步生效:文档增量更新、向量热替换与版本回滚的工程实现

运营同事改了一条退款政策,5 分钟后智能客服就能按新政策回答用户。这个"5 分钟"背后不是简单的全量重建,而是一套增量更新 + 向量热替换的工程方案。

一、为什么"5 分钟"很重要

1.1 业务驱动

趣玩搭作为一个社交活动平台,业务规则变更非常频繁。运营团队每周都可能修改活动报名条件、调整退款政策时间窗口、上线新的积分活动规则。如果每次改规则后知识库要等几小时甚至一天才能同步,期间智能客服就会给出过时的错误回答,直接引发用户投诉。

最初的方案是全量重建——每次文档有变更就把整个知识库的所有文档重新切分、重新 Embedding、重新写入 Milvus。2000 条 chunk 全量重建大约需要 15-20 分钟(主要耗时在调用 Embedding API),对于频繁更新的运营场景来说太慢了。

目标是把这个时间压缩到 5 分钟以内,而且要做到对线上服务零影响——更新过程中不能出现"知识库暂时不可用"或者"一半新一半旧"的中间状态。

1.2 挑战

增量更新看似简单(只重新处理变更的文档),但有几个工程细节需要处理:一篇文档修改后,它对应的 chunk 数量可能会变化(原来切成 5 段,改完后变成 3 段或 7 段),旧的 chunk 怎么清理?新旧 chunk 的 ID 怎么映射?更新过程中如果有用户查询命中了正在被替换的 chunk 怎么办?如果新文档的 Embedding 出错了,怎么回滚到旧版本?

二、整体架构

运营管理后台
    │ 编辑/新增/删除文档
    ▼
知识库管理服务
    │ 检测变更类型
    ├── 新增文档 → 切分 + Embedding + 写入 Milvus
    ├── 修改文档 → 增量更新流程(核心)
    └── 删除文档 → 标记删除 + 异步清理
    │
    ▼
变更事件 → RocketMQ → 异步处理
    │
    ▼
向量热替换(双版本切换)
    │
    ▼
知识库版本记录(支持回滚)

三、增量更新的核心流程

3.1 变更检测

运营在管理后台编辑文档后,系统通过文档版本号和内容哈希检测变更:

@Service
public class KnowledgeDocumentService {
​
    public void saveDocument(KnowledgeDocument doc) {
        KnowledgeDocument existing = docMapper.selectById(doc.getId());
​
        if (existing == null) {
            // 新增文档
            doc.setVersion(1);
            doc.setContentHash(DigestUtils.md5Hex(doc.getContent()));
            docMapper.insert(doc);
            publishEvent(new DocChangeEvent(doc.getId(), ChangeType.CREATE));
        } else {
            String newHash = DigestUtils.md5Hex(doc.getContent());
            if (!newHash.equals(existing.getContentHash())) {
                // 内容确实发生了变化
                doc.setVersion(existing.getVersion() + 1);
                doc.setContentHash(newHash);
                docMapper.updateById(doc);
                publishEvent(new DocChangeEvent(doc.getId(), ChangeType.UPDATE));
            }
            // hash 相同说明内容没变(可能只改了标题等元数据),只更新元数据
        }
    }
​
    private void publishEvent(DocChangeEvent event) {
        rocketMQTemplate.syncSend("knowledge-doc-change", event);
    }
}

通过内容哈希避免了"假更新"——运营点了保存但实际没改内容的情况,不会触发无意义的重新 Embedding。

3.2 增量处理消费者

@RocketMQMessageListener(
    topic = "knowledge-doc-change",
    consumerGroup = "knowledge-update-consumer"
)
@Service
public class KnowledgeUpdateConsumer implements RocketMQListener<DocChangeEvent> {
​
    @Override
    public void onMessage(DocChangeEvent event) {
        log.info("收到知识库变更事件: docId={}, type={}", event.getDocId(), event.getChangeType());
​
        switch (event.getChangeType()) {
            case CREATE -> handleCreate(event.getDocId());
            case UPDATE -> handleUpdate(event.getDocId());
            case DELETE -> handleDelete(event.getDocId());
        }
    }
​
    /**
     * 文档修改:增量更新核心逻辑
     */
    private void handleUpdate(String docId) {
        long startTime = System.currentTimeMillis();
​
        // 1. 加载最新文档
        KnowledgeDocument doc = docMapper.selectById(docId);
​
        // 2. 切分为 chunks
        List<DocumentChunk> newChunks = documentSplitter.split(doc);
​
        // 3. 批量 Embedding
        List<float[]> embeddings = embeddingService.batchEmbed(
            newChunks.stream().map(DocumentChunk::getContent).toList()
        );
        for (int i = 0; i < newChunks.size(); i++) {
            newChunks.get(i).setEmbedding(embeddings.get(i));
            newChunks.get(i).setChunkId(generateChunkId(docId, i));
        }
​
        // 4. 向量热替换(原子操作)
        vectorStoreService.atomicReplace(docId, newChunks);
​
        // 5. 记录版本(支持回滚)
        versionService.recordVersion(docId, doc.getVersion(), newChunks);
​
        long elapsed = System.currentTimeMillis() - startTime;
        log.info("文档 {} 增量更新完成,新chunk数: {}, 耗时: {}ms", docId, newChunks.size(), elapsed);
    }
}

3.3 批量 Embedding 优化

单条 Embedding 调用的网络延迟约 100-200ms,如果一篇文档切成 10 个 chunk 逐条调用,光 Embedding 就要 1-2 秒。通义的 Embedding API 支持批量调用(一次最多 25 条文本),大幅减少了网络往返。

@Service
public class EmbeddingService {
​
    private static final int BATCH_SIZE = 25;
​
    /**
     * 批量 Embedding,自动分批处理
     */
    public List<float[]> batchEmbed(List<String> texts) {
        List<float[]> allEmbeddings = new ArrayList<>();
​
        // 按 BATCH_SIZE 分批
        for (int i = 0; i < texts.size(); i += BATCH_SIZE) {
            List<String> batch = texts.subList(i, Math.min(i + BATCH_SIZE, texts.size()));
​
            List<float[]> batchResult = embeddingModel.embed(batch);
            allEmbeddings.addAll(batchResult);
        }
​
        return allEmbeddings;
    }
}

一篇普通文档(切成 5-10 个 chunk)的批量 Embedding 只需要一次 API 调用,耗时约 200-400ms。

3.4 向量热替换:不是删了再插,而是原子切换

最初的"删旧插新"方案有一个致命问题:在删除旧 chunk 和插入新 chunk 之间的时间窗口内,如果有用户查询,会找不到这篇文档的任何内容。

我的解决方案是"先插新、再删旧",通过 chunk ID 的版本标识实现原子切换:

@Service
public class VectorStoreService {
​
    /**
     * 原子替换某个文档的所有 chunks
     * 策略:先插入新版本 → 更新路由 → 再删除旧版本
     */
    public void atomicReplace(String docId, List<DocumentChunk> newChunks) {
        // 1. 查出旧版本的所有 chunk IDs
        List<String> oldChunkIds = getChunkIdsByDocId(docId);
​
        // 2. 插入新版本的 chunks(此时新旧共存)
        insertChunks(newChunks);
​
        // 3. 更新文档版本映射表(原子操作,此刻起检索会命中新版本)
        docVersionMapping.put(docId, newChunks.stream()
            .map(DocumentChunk::getChunkId)
            .toList());
​
        // 4. 异步删除旧版本 chunks(等待一段时间确保没有进行中的查询引用旧数据)
        CompletableFuture.runAsync(() -> {
            try {
                Thread.sleep(3000); // 等 3 秒,确保进行中的查询已完成
                deleteChunks(oldChunkIds);
                log.info("文档 {} 旧版本 {} 个chunks已清理", docId, oldChunkIds.size());
            } catch (Exception e) {
                log.error("清理旧chunks失败,docId={}", docId, e);
                // 清理失败不影响业务,下次更新时会再尝试清理
            }
        });
    }
}

chunk ID 的命名规则包含了文档版本号:

private String generateChunkId(String docId, int chunkIndex) {
    int version = docMapper.selectById(docId).getVersion();
    return String.format("%s_v%d_c%d", docId, version, chunkIndex);
    // 例:doc_refund_policy_001_v3_c2
    // 文档 refund_policy_001 的第3个版本的第2个chunk
}

检索时通过 Milvus 的过滤条件排除旧版本:

// 检索时只查最新版本的 chunks
String versionFilter = String.format("doc_version == %d", latestVersion);

实际上更简洁的做法是:检索不做版本过滤(新旧 chunk 都可能被命中),但在结果中如果同一个文档有新旧两个版本的 chunk 同时出现,只保留新版本的。这样避免了版本过滤条件对检索性能的影响。

四、版本记录与回滚

4.1 为什么需要回滚

运营同事不小心改错了文档怎么办?比如把退款政策的"48 小时"误改成"24 小时",上线后发现了要回滚。如果没有版本管理,就得手动改回文档再等一次增量更新,费时费力。

4.2 版本管理实现

CREATE TABLE `knowledge_doc_version` (
    `id` BIGINT PRIMARY KEY AUTO_INCREMENT,
    `doc_id` VARCHAR(64) NOT NULL,
    `version` INT NOT NULL,
    `content_snapshot` MEDIUMTEXT NOT NULL COMMENT '文档内容快照',
    `chunk_count` INT NOT NULL,
    `chunk_ids` TEXT NOT NULL COMMENT '该版本对应的chunk ID列表(JSON)',
    `operator` VARCHAR(64) COMMENT '操作人',
    `created_at` DATETIME DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY `uk_doc_version` (`doc_id`, `version`)
);
@Service
public class KnowledgeVersionService {

    /**
     * 回滚到指定版本
     */
    public void rollback(String docId, int targetVersion) {
        // 1. 查出目标版本的文档快照
        KnowledgeDocVersion versionRecord = versionMapper.selectByDocVersion(docId, targetVersion);
        if (versionRecord == null) {
            throw new BusinessException("版本不存在");
        }

        // 2. 用快照内容覆盖当前文档
        KnowledgeDocument doc = docMapper.selectById(docId);
        doc.setContent(versionRecord.getContentSnapshot());
        doc.setVersion(doc.getVersion() + 1); // 回滚也是一次新版本
        doc.setContentHash(DigestUtils.md5Hex(doc.getContent()));
        docMapper.updateById(doc);

        // 3. 触发增量更新(和正常编辑走同一条链路)
        publishEvent(new DocChangeEvent(docId, ChangeType.UPDATE));

        log.info("文档 {} 回滚到版本 {},新版本号 {}", docId, targetVersion, doc.getVersion());
    }
}

回滚操作的本质就是"用旧版本的内容生成一个新版本",然后走正常的增量更新流程。

4.3 管理后台的版本对比

运营管理后台提供了文档版本历史和内容差异对比功能。运营同事可以看到每次修改了什么内容、由谁修改、什么时候生效的。这在出现客服回答错误时非常有用——可以快速定位是哪次文档修改引入了错误。

五、耗时分析与 5 分钟的保障

一次增量更新的完整耗时拆解:

环节

耗时

说明

变更检测 + MQ 投递

~50ms

内容哈希比对 + MQ 发送

MQ 消费延迟

~200ms

RocketMQ 消费延迟

文档切分

~100ms

本地计算

批量 Embedding

200-800ms

取决于 chunk 数量,通常 1 次 API 调用

Milvus 写入

~200ms

批量 upsert

旧数据清理

3-5s

异步执行,不影响可用性

端到端总耗时

约 4-6 秒

单篇文档

单篇文档的增量更新在 10 秒以内完成。"5 分钟"的承诺留出了大量的安全余量,即使同时更新几十篇文档(排队串行处理),也能在 5 分钟内全部完成。

实际上运营的使用模式通常是"改一篇发布一篇",大多数情况下 10 秒内就生效了。"5 分钟"是对批量更新场景(如运营集中修改一批活动规则)的承诺。

六、全量重建的保底方案

增量更新覆盖了 99% 的场景,但仍然保留了全量重建的能力,用于以下情况:

Embedding 模型更换。 如果换了 Embedding 模型,所有 chunk 都需要重新生成向量,只能全量重建。

切分策略大调整。 如果修改了文档切分逻辑(比如调整了 chunk size),所有文档需要重新切分。

数据一致性修复。 如果因为 bug 或其他原因导致 Milvus 中的数据和源文档不一致,可以通过全量重建修复。

全量重建通过一个管理后台按钮触发,在后台异步执行,不影响线上服务(重建到一个新的 Milvus Collection,完成后原子切换)。

public void fullRebuild() {
    String newCollection = "knowledge_chunks_" + System.currentTimeMillis();

    // 1. 创建新 Collection
    milvusClient.createCollection(newCollection, schema);

    // 2. 遍历所有文档,切分 + Embedding + 写入新 Collection
    List<KnowledgeDocument> allDocs = docMapper.selectAll();
    for (KnowledgeDocument doc : allDocs) {
        List<DocumentChunk> chunks = documentSplitter.split(doc);
        List<float[]> embeddings = embeddingService.batchEmbed(
            chunks.stream().map(DocumentChunk::getContent).toList());
        // ... 写入 newCollection
    }

    // 3. 原子切换 Collection 别名
    milvusClient.createAlias(newCollection, "knowledge_active");
    milvusClient.dropCollection(oldCollection);

    log.info("全量重建完成,新Collection: {}", newCollection);
}

七、经验总结

增量更新的核心是"先插新再删旧"。 绝对不能"先删再插",否则中间状态会导致用户查询缺失。

内容哈希可以过滤掉大量无效更新。 运营同事经常"手滑"点保存但实际没改内容,没有哈希校验的话每次都会触发一轮完整的切分 + Embedding,浪费 API 调用费用。

版本管理是运营的安全网。 没有回滚能力的知识库系统,运营团队不敢大胆修改文档。有了版本对比和一键回滚,运营效率提升了很多。

全量重建是兜底而不是常态。 日常更新走增量(秒级),大版本变更走全量(分钟级),两套方案互补。


如果这篇文章对你有帮助,欢迎访问我的博客 robinzhu.top 获取更多实战分享。