RAG 知识库真正困难的部分通常不是第一次导入,而是后续文件修改、删除和同步失败后的状态恢复。本文用文件指纹、版本号和删除标记构建一条可重复执行的增量同步流水线。示例只使用 Node.js 内置模块,但保留了替换为向量数据库适配器所需的边界。
一次性导入为什么会失控
很多知识库的第一版流程是:扫描目录、读取文件、切分文本、生成向量,然后批量写入索引。这条流程适合验证原型,却不适合作为长期运行的同步任务。
文件没有变化时,重复执行会再次切分和写入,造成重复向量;文件被修改时,如果只追加新内容,旧分块仍然会被召回;文件被删除时,如果没有显式删除索引记录,知识库仍可能回答已经不存在的内容。同步任务在写入一半时进程退出,也会让源文件状态、元数据状态和索引状态互相不一致。
因此,同步系统需要把“源文件当前是什么状态”和“索引曾经做过哪些变更”分开保存。前者用于计算下一次计划,后者用于审计、重试和回滚。
本文示例采用以下模型:
| 对象 | 作用 | 关键字段 |
|---|---|---|
| 文件指纹 | 判断内容是否变化 | path、sha256、size、mtimeMs |
| 文档版本 | 标识一次有效内容 | documentId、version、status |
| 分块记录 | 保存当前可检索内容 | chunkId、documentId、version、text |
| 变更事件 | 记录同步操作 | runId、action、beforeVersion、afterVersion |
这里的 documentId 应该稳定地来源于相对路径或业务主键,而不是来源于文件内容。内容变化只增加版本号,不改变文档身份。
先设计可重复执行的同步协议
一次同步可以分为五个阶段:扫描、比对、计划、应用、提交。扫描阶段只读取源目录并计算指纹;比对阶段找出新增、修改和删除;计划阶段生成固定的操作列表;应用阶段更新索引;提交阶段记录本次运行结果。
关键点是:不要在扫描到一个文件后立即写索引。先生成完整计划,能够让日志、人工检查和失败重试都更容易。对于同一个输入状态,计划应当稳定地产生相同的操作。
可以把状态转换写成如下规则:
| 当前状态 | 文件状态 | 操作 |
|---|---|---|
| 不存在 | 新增 | upsert,版本从 1 开始 |
| 存在 | 指纹相同 | skip,不重新切分 |
| 存在 | 指纹不同 | 先删除旧版本,再写入新版本 |
| 存在 | 文件消失 | 写入删除标记,并删除当前分块 |
| 已删除 | 文件重新出现 | 重新写入,版本继续递增 |
删除标记不能只依赖物理删除。物理删除负责让检索结果立即消失,删除标记负责保留“这个文档曾经存在但已被删除”的事实。这样可以避免下次扫描时把删除文档误判为从未同步,也方便审计和恢复。
版本号也不应该通过重新扫描时临时推导。版本应该存入元数据,并在每次内容变更时递增。回滚时可以把某个历史版本重新写成当前版本,或者恢复其分块,但不建议直接把审计记录改掉。
一个可运行的 Node.js 实现
下面的程序使用 Node.js 18 或更高版本运行,不依赖第三方包。它把源文件放在 ./documents,把知识库状态写入 ./rag-state.json。示例中的 Index 类模拟向量数据库适配器:真实项目中,只需要把 upsert 和 deleteByDocument 替换为目标数据库的批量接口,状态机和审计逻辑仍然适用。
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。
还可以在事件中增加 operationId 和 checksum,由数据库适配器拒绝相同 operationId 的重复提交。对于跨系统事务无法原子提交的问题,不要假设存在真正的全局事务,而应通过对账任务修复状态:扫描元数据中的 active 版本,检查索引中是否存在对应分块,缺失时重新投递任务。
回滚应当是“生成反向变更”,而不是删除审计记录。假设文档当前是版本 3,事件中保存了版本 2 的内容快照或内容对象存储地址,那么回滚可以按以下流程执行:
- 读取目标历史版本,确认校验和与分块数量。
- 为当前文档生成新的版本 4,其内容等于版本 2。
- 删除文档当前分块,写入版本 4 的分块。
- 把文档元数据更新为 active,并记录
rollback事件,注明fromVersion: 3和sourceVersion: 2。
这里不直接把版本号改回 2,是因为版本号应描述事件顺序,而不是只描述内容编号。这样审计日志仍然是追加的,问题排查也能回答“谁在什么时间把内容恢复成了什么状态”。如果数据量较大,历史内容可以放在对象存储中,状态文件或关系数据库只保存地址、校验和与版本元数据。
工程化边界与验证清单
示例用 JSON 文件保存状态,适合本地演示,不适合作为多实例生产存储。生产环境至少应将文档元数据、同步运行记录和变更事件放到支持并发控制的数据库中,并为 documentId + version、runId 建立约束。同步任务也应使用单实例锁,或者采用队列把每个文档的操作串行化。
切分策略改变时,不能只比较文件指纹。应把 chunkerVersion 纳入文档版本计算条件;嵌入模型更换时,也应把 embeddingModel 写入分块元数据,并重新建立对应版本。否则文件内容没有变化,索引却已经不兼容,系统仍会错误地跳过同步。
上线前可以围绕以下情况写集成测试:空目录首次同步;完全相同的文件重复同步;文件修改后旧分块是否消失;文件删除后是否产生删除标记;删除文件重新出现时版本是否递增;写入中途失败后重试是否产生重复分块;历史版本回滚后查询结果是否只包含恢复后的内容。
另一个容易忽略的问题是查询过滤。检索时应过滤 status = active,并校验分块版本与当前文档版本一致。物理删除失败时,删除标记可以阻止新查询继续使用旧内容;后台清理任务随后再删除残留向量。这种“双重保护”比只依赖一次删除请求更稳妥。
总结
本文的核心不是某个向量数据库 API,而是一套可重复执行的状态管理方法:
- 用稳定的文档标识和 SHA-256 文件指纹判断内容是否变化。
- 用递增版本号区分每次有效内容,避免修改后新旧分块共存。
- 用删除标记表达文件已消失的事实,同时删除当前可检索分块。
- 先生成同步计划,再应用变更,让执行、审计和重试有清晰边界。
- 使用确定性
chunkId和幂等操作,降低同步中断后的重复写入风险。 - 通过追加式事件记录支持审计,并用生成新版本的方式执行回滚。
- 将切分器版本和嵌入模型纳入元数据,避免基础设施变化造成错误跳过。
当知识库进入长期运行阶段,增量同步应被视为数据管道的一部分,而不是导入脚本的附属功能。文件指纹负责发现变化,版本负责描述变化,删除标记负责表达消失,审计事件则让失败和回滚变成可以验证的工程流程。
评论