知识库几百上千个片段要批量向量化,任务怎么设计?失败了怎么处理?
简化版
批量向量化要逐条调用向量模型接口,耗时随片段数线性增长,不能放在一个同步请求里跑:接口只负责校验、启动后台任务并立刻返回总数,前端轮询进度。任务本身有三条关键规则:同一时间只允许一个全量任务(用原子操作抢运行标记);逐条提交,一条失败不影响其他条,进度里如实记录成功数和失败数;把错误分成两类:限流、网络抖动这类偶发错误跳过继续,欠费、Key 无效、模型名写错这类每条都会同样失败的错误,第一条就停下并报告原因。
详细版
点击「生成全部向量」
→ 接口:抢运行标记(失败说明已有任务在跑)→ 查待处理片段 → 初始化进度 → 启动后台线程 → 返回总数
→ 前端:每 1~2 秒查询一次进度,直到任务结束
后台线程逐条处理:
删这条的旧向量 → 调向量模型 → 写新向量 → 成功数 +1
出错:致命错误 → 立即结束任务,写明原因和已成功条数
偶发错误 → 失败数 +1,继续下一条
整个循环外层再包一层:循环本身崩了也要把任务标为结束、释放运行标记
| 设计点 | 做法 | 不这样做的后果 |
|---|---|---|
| 异步执行 | 接口立即返回,后台线程处理 | 请求超时,前端报失败而后端还在写库 |
| 单任务 | 原子操作抢运行标记 | 两个任务同时跑,互相覆盖向量 |
| 进度可见 | 进度对象跨线程可见,前端轮询 | 进度条不动或看不到结果 |
| 逐条提交 | 每条成功立即写库 | 中途失败时不知道到底成功了多少 |
| 错误分类 | 致命错误提前停,偶发错误跳过 | 对几百条发注定失败的请求,最后只得到一句「N 条失败」 |
| 兜底结束 | 外层异常也要结束任务 | 运行标记永远不释放,前端永远在轮询 |
完整版教学
一、为什么批量向量化不能同步跑完
粗算一下耗时(按单次向量接口调用约 0.3 秒估算,实际取决于模型服务和网络):
100 个片段 × 0.3 秒 ≈ 30 秒
556 个片段 × 0.3 秒 ≈ 167 秒
5000 个片段 × 0.3 秒 ≈ 25 分钟
片段数是不确定的,所以任何固定的请求超时都可能不够。把超时设长也不是解法:设到 10 分钟,数据再涨一点还是会断;而且前端超时并不会停掉后端,前端显示「失败了」,后端循环还在跑、还在写库,用户看到的状态和真实状态完全不一致。
正确做法是把长任务和 HTTP 请求解耦:接口只启动任务,执行放到后台,进度靠查询获取。任务从此不依赖那条 HTTP 连接,页面关掉也不影响它继续执行。同样的思路用在多步分析任务上,见「多步分析任务耗时很长,接口怎么设计?步骤成功怎么判定?」。
二、只允许一个任务在跑
用户可能连点两次,或者两个管理员同时点。两个全量任务同时跑,会各自先删再写同一批向量,结果互相覆盖。所以启动前要抢一个运行标记:
// 原子操作:只有当前是 false 才能改成 true,并返回是否成功
if (!running.compareAndSet(false, true)) {
throw new CustomException("已有全量向量化任务在执行,请等它结束后再试");
}
为什么不用 if (running) return; running = true;?两个请求可能同时通过检查,都把标记设成 true,还是启动了两个任务。compareAndSet 把「检查」和「设置」合成一个不可分割的动作。
还有一个配套要求:启动阶段出错也要把标记复位。比如向量模型没配置、查询片段失败,任务根本没启动,如果标记还是 true,以后所有全量任务都会被挡在门外。
三、进度对象与跨线程可见性
写进度的是后台线程,读进度的是处理轮询请求的 Web 线程。进度对象的引用要保证跨线程可见:
private volatile Progress progress = new Progress();
没有 volatile,轮询线程可能一直读到自己缓存的旧对象,页面进度条就不动了。进度对象至少包含:是否运行中、总数、已处理数、成功数、失败数、结束说明。先把进度对象初始化好再启动线程,保证前端第一次查询拿到的就是有效进度。
前端轮询也有三个细节:
| 细节 | 原因 |
|---|---|
用 setTimeout 递归,不用 setInterval | 上一次请求回来了才发下一次,后端慢时不会堆积请求 |
| 单次轮询失败不停止 | 任务在后端照常跑,不能因为一次网络抖动就丢掉进度 |
| 离开页面清掉定时器 | 组件销毁后不再发请求、不再写数据 |
四、逐条提交,失败只影响那一条
批量任务里最容易写错的是「一条失败就整体抛异常」:
第 30 条调用失败 → 抛异常中断
→ 前 29 条其实已经写进数据库了
→ 页面只看到一句笼统的失败,不知道成功了多少、哪些没做
正确做法是每条独立处理、独立提交:成功就写库并计数,失败就记失败数继续下一条;任务结束时进度里带着准确的成功数和失败数,没生成的片段在列表里显示为未向量化或失败状态,下次重跑补上即可。
记忆钩子:批量任务要能回答「做了多少、还差多少、为什么停」,这三个数必须在进度里。
五、错误要分两类处理
调用外部模型接口,失败是常态,但失败的性质完全不同:
| 类型 | 例子 | 特点 | 处理 |
|---|---|---|---|
| 偶发错误 | 限流、网络抖动、单次超时 | 下一条可能就成功 | 计失败数,继续 |
| 致命错误 | 欠费、Key 无效、模型名写错 | 每一条都会以同样原因失败 | 第一条就停,报告原因和已成功数 |
如果不区分,Key 填错时会对几百个片段逐个发请求,每个都返回 401,使用者盯着进度条跑完全程,最后只得到「556 条失败」。区分之后,第一条失败就能看到「API Key 无效,已成功 0 条」,立刻去改配置。判断方式通常是看翻译后的错误信息或错误码,把「重试也没用」的几类单独识别出来。
同理,批量开始前先解析一次模型配置:配置不可用就直接报错,不启动一个注定失败的任务;解析一次复用到所有片段,也避免每条都查一遍配置表。
六、兜底:任务一定要能结束
后台循环外层必须再包一层异常处理:
try {
for (Chunk chunk : chunks) {
try { embedAndSave(chunk); success++; }
catch (Exception e) { if (isFatal(e)) { finish("中断:" + reason(e)); return; } failed++; }
processed++;
}
finish("完成:成功 " + success + " 条,失败 " + failed + " 条");
} catch (Exception e) {
finish("任务异常结束:" + e.getMessage()); // 否则运行标记永远是 true
}
内层管单条失败,外层管整个循环崩溃。外层这一层最容易被省略:一旦后台线程抛出未捕获异常,运行标记停在 true,前端永远在轮询,以后也再启动不了新任务。finish 里统一做三件事:写结束说明、记录结束时间、释放运行标记。
七、粒度的另一种选择:一条知识要么全成要么全撤
上面是「片段级」提交:每个片段独立成功或失败。另一种粒度是「文档级」:一篇文档切出的所有片段要么全部向量化成功,要么全部撤掉、文档标记为失败。
| 粒度 | 优点 | 代价 |
|---|---|---|
| 片段级 | 失败影响最小,重跑只补失败的片段 | 一篇文档可能部分可检索,结果不完整 |
| 文档级 | 一篇文档要么完整可检索,要么完全不可检索 | 一个片段失败,整篇都要重做 |
文档内容前后关联紧密、缺一段就可能误导回答时,适合文档级;片段相互独立、量大时,适合片段级。
八、常见误区与追问
- 误区:把请求超时调长就能解决批量向量化。 数据量一涨还是会断,而且前端超时不会停掉后端,状态会不一致,应该改成后台任务加进度查询。
- 误区:用 if 判断运行标记就能防止重复启动。 两个请求可能同时通过检查,要用
compareAndSet这类原子操作。 - 误区:一条失败就整体抛异常最安全。 前面成功的已经写库了,整体失败反而让人不知道做了多少,应该逐条提交、分别计数。
- 误区:所有失败都跳过继续。 Key 无效、欠费这类错误每条都会失败,应该第一条就停并报告原因。
- 误区:后台循环只需要处理单条异常。 循环本身崩溃时运行标记不会释放,外层必须再兜一层。
- 追问:为什么轮询用 setTimeout 递归而不是 setInterval? 递归保证上一次请求回来才发下一次,后端变慢时不会堆积请求。
- 追问:进度对象为什么要加 volatile? 写进度和读进度在不同线程,没有它读线程可能一直看到旧值。
九、加强记忆
批量向量化耗时随片段数增长,固定超时怎么设都不对,所以接口只启动后台任务并返回总数,前端用 setTimeout 递归轮询进度、失败不停、离开页面清定时器。启动时用 compareAndSet 抢运行标记,启动失败要复位;进度对象加 volatile,先初始化再开线程。循环里逐条提交、分别计成功数和失败数;偶发错误跳过,Key 无效、欠费、模型名错这类致命错误第一条就停并报已成功数;开始前先确认模型配置可用。外层再包一层异常,保证任务一定结束、标记一定释放。片段级提交影响小,文档级全成全撤保证完整,按内容关联程度选。
项目实战落地
项目里怎么做的
《AI Agent电商平台导购与运营增长系统》的「生成全部向量」拆成启动和查进度两个接口:
- 启动:
startGenerateAll用compareAndSet(false, true)抢运行标记,查出状态为 READY 的切片,初始化进度对象,整批只解析一次向量模型配置,然后交给CompletableFuture后台执行,接口立刻返回切片总数; - 进度:
EmbeddingGenerateProgress记录是否运行中、总数、已处理数、成功数、失败数和结果说明,Service 里用volatile字段持有;前端每秒用setTimeout递归查询一次,单次失败继续重试,onBeforeUnmount清定时器; - 后台循环:每条先删旧向量再生成写入;内层
try-catch处理单条失败,欠费、Key 无效这类错误由isFatalModelError识别后立即结束;外层try-catch保证循环崩溃时也把任务置为结束,finishProgress负责释放运行标记。
《AI Agent岗位匹配与求职规划系统》用的是文档级粒度:一条求职知识的任一切片向量化失败,就删除本次已写入的切片,切片数归零,状态记为「向量化失败」并保存原因;「全部向量化」逐条复用这个逻辑,返回成功的知识条数。
为什么这样取舍
- 商城项目用片段级:商品知识切片相互独立,一条失败不影响其他切片被检索,重跑时只补失败的。
- 求职项目用文档级:一条求职知识是一篇完整的指导内容,只剩一半切片可检索时,Agent 拿到的建议可能不完整。
面试官还会追问
- 抢运行标记为什么用
compareAndSet(false, true),而不是先判断运行中再置为 true? - 进度对象为什么要用
volatile修饰?不加会出现什么现象?
学完《AI Agent电商平台导购与运营增长系统》,上面这些追问你都会迎刃而解。