RAG 知识库真正困难的部分通常不是第一次导入,而是后续文件修改、删除和同步失败后的状态恢复。本文用文件指纹、版本号和删除标记构建一条可重复执行的增量同步流水线。示例只使用 Node.js 内置模块,但保留了替换为向量数据库适配器所需的边界。

一次性导入为什么会失控

很多知识库的第一版流程是:扫描目录、读取文件、切分文本、生成向量,然后批量写入索引。这条流程适合验证原型,却不适合作为长期运行的同步任务。

文件没有变化时,重复执行会再次切分和写入,造成重复向量;文件被修改时,如果只追加新内容,旧分块仍然会被召回;文件被删除时,如果没有显式删除索引记录,知识库仍可能回答已经不存在的内容。同步任务在写入一半时进程退出,也会让源文件状态、元数据状态和索引状态互相不一致。

因此,同步系统需要把“源文件当前是什么状态”和“索引曾经做过哪些变更”分开保存。前者用于计算下一次计划,后者用于审计、重试和回滚。

本文示例采用以下模型:

对象作用关键字段
文件指纹判断内容是否变化pathsha256sizemtimeMs
文档版本标识一次有效内容documentIdversionstatus
分块记录保存当前可检索内容chunkIddocumentIdversiontext
变更事件记录同步操作runIdactionbeforeVersionafterVersion

这里的 documentId 应该稳定地来源于相对路径或业务主键,而不是来源于文件内容。内容变化只增加版本号,不改变文档身份。

先设计可重复执行的同步协议

一次同步可以分为五个阶段:扫描、比对、计划、应用、提交。扫描阶段只读取源目录并计算指纹;比对阶段找出新增、修改和删除;计划阶段生成固定的操作列表;应用阶段更新索引;提交阶段记录本次运行结果。

关键点是:不要在扫描到一个文件后立即写索引。先生成完整计划,能够让日志、人工检查和失败重试都更容易。对于同一个输入状态,计划应当稳定地产生相同的操作。

可以把状态转换写成如下规则:

当前状态文件状态操作
不存在新增upsert,版本从 1 开始
存在指纹相同skip,不重新切分
存在指纹不同先删除旧版本,再写入新版本
存在文件消失写入删除标记,并删除当前分块
已删除文件重新出现重新写入,版本继续递增

删除标记不能只依赖物理删除。物理删除负责让检索结果立即消失,删除标记负责保留“这个文档曾经存在但已被删除”的事实。这样可以避免下次扫描时把删除文档误判为从未同步,也方便审计和恢复。

版本号也不应该通过重新扫描时临时推导。版本应该存入元数据,并在每次内容变更时递增。回滚时可以把某个历史版本重新写成当前版本,或者恢复其分块,但不建议直接把审计记录改掉。

一个可运行的 Node.js 实现

下面的程序使用 Node.js 18 或更高版本运行,不依赖第三方包。它把源文件放在 ./documents,把知识库状态写入 ./rag-state.json。示例中的 Index 类模拟向量数据库适配器:真实项目中,只需要把 upsertdeleteByDocument 替换为目标数据库的批量接口,状态机和审计逻辑仍然适用。

const fs = require('node:fs/promises');
const path = require('node:path');
const crypto = require('node:crypto');

const ROOT = path.resolve('./documents');
const STATE_FILE = path.resolve('./rag-state.json');

async function exists(file) {
  try {
    await fs.access(file);
    return true;
  } catch {
    return false;
  }
}

async function loadState() {
  if (!(await exists(STATE_FILE))) {
    return { documents: {}, chunks: {}, events: [], runs: [] };
  }
  return JSON.parse(await fs.readFile(STATE_FILE, 'utf8'));
}

async function saveState(state) {
  const temporary = `${STATE_FILE}.${process.pid}.tmp`;
  await fs.writeFile(temporary, JSON.stringify(state, null, 2));
  await fs.rename(temporary, STATE_FILE);
}

async function listFiles(dir) {
  const result = [];
  for (const entry of await fs.readdir(dir, { withFileTypes: true })) {
    const full = path.join(dir, entry.name);
    if (entry.isDirectory()) result.push(...await listFiles(full));
    else if (entry.isFile() && entry.name.endsWith('.txt')) result.push(full);
  }
  return result;
}

function fingerprint(buffer, stat) {
  return {
    sha256: crypto.createHash('sha256').update(buffer).digest('hex'),
    size: stat.size,
    mtimeMs: stat.mtimeMs
  };
}

function splitText(text, size = 120) {
  const normalized = text.replace(/\\r\\n/g, '\\n').trim();
  if (!normalized) return [];
  const chunks = [];
  for (let start = 0; start < normalized.length; start += size) {
    chunks.push(normalized.slice(start, start + size));
  }
  return chunks;
}

class Index {
  constructor(state) {
    this.state = state;
  }

  async deleteByDocument(documentId) {
    for (const [chunkId, chunk] of Object.entries(this.state.chunks)) {
      if (chunk.documentId === documentId) delete this.state.chunks[chunkId];
    }
  }

  async upsert(documentId, version, texts) {
    await this.deleteByDocument(documentId);
    texts.forEach((text, position) => {
      const chunkId = `${documentId}:v${version}:c${position}`;
      this.state.chunks[chunkId] = { chunkId, documentId, version, position, text };
    });
  }
}

async function sync() {
  const state = await loadState();
  const index = new Index(state);
  const runId = crypto.randomUUID();
  const plan = [];
  const seen = new Set();

  await fs.mkdir(ROOT, { recursive: true });
  for (const file of await listFiles(ROOT)) {
    const relative = path.relative(ROOT, file).split(path.sep).join('/');
    const buffer = await fs.readFile(file);
    const stat = await fs.stat(file);
    const current = fingerprint(buffer, stat);
    seen.add(relative);
    const old = state.documents[relative];
    if (!old || old.status === 'deleted') plan.push({ action: 'upsert', path: relative, current, content: buffer.toString('utf8'), old });
    else if (old.sha256 !== current.sha256) plan.push({ action: 'upsert', path: relative, current, content: buffer.toString('utf8'), old });
    else plan.push({ action: 'skip', path: relative });
  }

  for (const [documentId, old] of Object.entries(state.documents)) {
    if (old.status !== 'deleted' && !seen.has(documentId)) plan.push({ action: 'delete', path: documentId, old });
  }

  const eventIds = [];
  try {
    for (const operation of plan) {
      if (operation.action === 'skip') continue;
      const old = operation.old;
      const beforeVersion = old ? old.version : null;
      const afterVersion = operation.action === 'upsert' ? (old ? old.version + 1 : 1) : old.version;

      if (operation.action === 'upsert') {
        const texts = splitText(operation.content);
        await index.upsert(operation.path, afterVersion, texts);
        state.documents[operation.path] = {
          path: operation.path,
          sha256: operation.current.sha256,
          size: operation.current.size,
          mtimeMs: operation.current.mtimeMs,
          version: afterVersion,
          status: 'active',
          updatedAt: new Date().toISOString()
        };
      } else {
        await index.deleteByDocument(operation.path);
        state.documents[operation.path] = { ...old, status: 'deleted', deletedAt: new Date().toISOString() };
      }

      const event = {
        eventId: crypto.randomUUID(),
        runId,
        action: operation.action,
        documentId: operation.path,
        beforeVersion,
        afterVersion,
        at: new Date().toISOString()
      };
      state.events.push(event);
      eventIds.push(event.eventId);
    }
    state.runs.push({ runId, status: 'succeeded', eventIds, at: new Date().toISOString() });
    await saveState(state);
    console.log(JSON.stringify({ runId, plan, status: 'succeeded' }, null, 2));
  } catch (error) {
    state.runs.push({ runId, status: 'failed', eventIds, error: error.message, at: new Date().toISOString() });
    await saveState(state);
    throw error;
  }
}

sync().catch(error => {
  console.error(error);
  process.exitCode = 1;
});

运行方式如下:

mkdir -p documents
printf '第一版内容。\\n第二段内容。' > documents/guide.txt
node sync.js
printf '修改后的内容。\\n新增一段内容。' > documents/guide.txt
node sync.js
rm documents/guide.txt
node sync.js

第一次运行会创建版本 1,第二次运行会创建版本 2,删除文件后会保留 status: "deleted"。相同文件再次运行时,计划中只有 skip,不会重新切分或写入索引。

失败重试与安全回滚

上面的实现已经具备两个基础能力:每次运行有唯一 runId,每个变更有独立事件;写入失败时会记录失败运行。生产环境还需要进一步处理“索引已写入、状态文件尚未提交”的情况。

一个实用做法是把每个操作设计成幂等命令,并让写入顺序固定为:先根据稳定的 documentId 删除旧分块,再写入指定版本的分块,最后提交文档元数据。任务重试时,即使上一次只完成了一半,重新执行同一操作也不会产生重复分块。向量数据库侧最好使用确定性的 chunkId,例如示例中的 文档标识:版本:分块序号,不要使用随机 ID。

还可以在事件中增加 operationIdchecksum,由数据库适配器拒绝相同 operationId 的重复提交。对于跨系统事务无法原子提交的问题,不要假设存在真正的全局事务,而应通过对账任务修复状态:扫描元数据中的 active 版本,检查索引中是否存在对应分块,缺失时重新投递任务。

回滚应当是“生成反向变更”,而不是删除审计记录。假设文档当前是版本 3,事件中保存了版本 2 的内容快照或内容对象存储地址,那么回滚可以按以下流程执行:

  1. 读取目标历史版本,确认校验和与分块数量。
  2. 为当前文档生成新的版本 4,其内容等于版本 2。
  3. 删除文档当前分块,写入版本 4 的分块。
  4. 把文档元数据更新为 active,并记录 rollback 事件,注明 fromVersion: 3sourceVersion: 2

这里不直接把版本号改回 2,是因为版本号应描述事件顺序,而不是只描述内容编号。这样审计日志仍然是追加的,问题排查也能回答“谁在什么时间把内容恢复成了什么状态”。如果数据量较大,历史内容可以放在对象存储中,状态文件或关系数据库只保存地址、校验和与版本元数据。

工程化边界与验证清单

示例用 JSON 文件保存状态,适合本地演示,不适合作为多实例生产存储。生产环境至少应将文档元数据、同步运行记录和变更事件放到支持并发控制的数据库中,并为 documentId + versionrunId 建立约束。同步任务也应使用单实例锁,或者采用队列把每个文档的操作串行化。

切分策略改变时,不能只比较文件指纹。应把 chunkerVersion 纳入文档版本计算条件;嵌入模型更换时,也应把 embeddingModel 写入分块元数据,并重新建立对应版本。否则文件内容没有变化,索引却已经不兼容,系统仍会错误地跳过同步。

上线前可以围绕以下情况写集成测试:空目录首次同步;完全相同的文件重复同步;文件修改后旧分块是否消失;文件删除后是否产生删除标记;删除文件重新出现时版本是否递增;写入中途失败后重试是否产生重复分块;历史版本回滚后查询结果是否只包含恢复后的内容。

另一个容易忽略的问题是查询过滤。检索时应过滤 status = active,并校验分块版本与当前文档版本一致。物理删除失败时,删除标记可以阻止新查询继续使用旧内容;后台清理任务随后再删除残留向量。这种“双重保护”比只依赖一次删除请求更稳妥。

总结

本文的核心不是某个向量数据库 API,而是一套可重复执行的状态管理方法:

  • 用稳定的文档标识和 SHA-256 文件指纹判断内容是否变化。
  • 用递增版本号区分每次有效内容,避免修改后新旧分块共存。
  • 用删除标记表达文件已消失的事实,同时删除当前可检索分块。
  • 先生成同步计划,再应用变更,让执行、审计和重试有清晰边界。
  • 使用确定性 chunkId 和幂等操作,降低同步中断后的重复写入风险。
  • 通过追加式事件记录支持审计,并用生成新版本的方式执行回滚。
  • 将切分器版本和嵌入模型纳入元数据,避免基础设施变化造成错误跳过。

当知识库进入长期运行阶段,增量同步应被视为数据管道的一部分,而不是导入脚本的附属功能。文件指纹负责发现变化,版本负责描述变化,删除标记负责表达消失,审计事件则让失败和回滚变成可以验证的工程流程。