文档解析异步化
1. 为什么要把解析放到后台
分段预览原来是同步接口:浏览器上传文件后一直等在那里,后端把整份文档解析完再一次性返回分段。整份解析要逐页做远程 OCR,一份 35 页、12 MB 的扫描 PDF 需要几分钟到十几分钟,于是同步实现有三个绕不开的问题:
- 请求线程被长时间占用。HTTP 连接、超时配置和中间的反向代理都要为最长的那份文档让路,用户看到的进度条也分不清是“正在解析”还是“已经卡死”。
- 任何一页失败都会让整份文档失败。远程 OCR 是外部服务,偶发失败、超时和限流都会发生;同步实现里一页抛异常,前面已经完成的页也一起白做。
- 结果无法复用。用户关掉页面、刷新或者重新点一次“预览”,前端拿不到上一次的结果,后端也只能重算。
现在的实现把“接收文件”和“解析文件”拆成两个动作:上传接口只做保存与建任务,立即返回任务 ID;解析在后台虚拟线程里逐页推进,把进度和最终分段写进任务表;前端按固定间隔轮询任务状态,拿到分段后继续原来的确认与入库流程。
2. 上传立即返回,解析在后台推进
上传接口不再等解析:
package nexus.io.mosskb.controller;
import java.util.ArrayList;
import java.util.List;
import com.jfinal.kit.Kv;
import nexus.io.annotation.Get;
import nexus.io.annotation.Post;
import nexus.io.annotation.RequestPath;
import nexus.io.jfinal.aop.Aop;
import nexus.io.mosskb.service.SystemFileService;
import nexus.io.mosskb.service.kb.MossKbDocumentSplitService;
import nexus.io.mosskb.service.kb.MossKbDocumentSplitTaskService;
import nexus.io.mosskb.vo.MossKbDocumentSplitTaskVo;
import nexus.io.model.result.ResultVo;
import nexus.io.model.upload.UploadFile;
import nexus.io.model.upload.UploadResult;
import nexus.io.tio.boot.http.TioRequestContext;
@RequestPath("/api/dataset/document")
public class ApiDatasetDocumentController {
/**
* 分段预览:只负责上传与建任务,解析在后台线程逐页推进。
*
* <p>大文档解析要几分钟,请求线程不再等待,前端拿 task_id 轮询 {@link #splitTask(Long)} 取结果。
*/
@Post("/split")
public ResultVo split(nexus.io.tio.http.common.HttpRequest request) {
Object[] files = request.getParams().get("file");
if (files == null || files.length == 0) {
return ResultVo.fail("请求体中未找到文件");
}
Long userId = TioRequestContext.getUserIdLong();
List<Long> taskIds = new ArrayList<>();
List<Kv> failures = new ArrayList<>();
for (Object item : files) {
if (!(item instanceof UploadFile)) {
return ResultVo.fail("文件格式不正确");
}
UploadFile file = (UploadFile) item;
UploadResult uploaded = Aop.get(SystemFileService.class).upload(file, "default", "default");
if (uploaded == null) {
failures.add(Kv.by("name", file.getName()).set("message", "文件保存失败"));
continue;
}
try {
taskIds.add(Aop.get(MossKbDocumentSplitService.class).splitAsync(file.getData(), uploaded, userId));
} catch (Exception e) {
failures.add(Kv.by("name", file.getName()).set("message", e.getMessage()));
}
}
if (taskIds.isEmpty()) {
return ResultVo.fail(failures.isEmpty() ? "没有可解析的文件" : String.valueOf(failures.get(0).get("message")));
}
return ResultVo.ok(Kv.by("task_id_list", taskIds).set("failures", failures));
}
}
上传多个文件时每个文件一个任务,响应里的 task_id_list 与 failures 分别表示“已经进入解析”和“连保存都没成功”。文件本身仍然走原来的上传服务,因此同一份文件重复上传会复用已有的文件 ID,前端提交分段时回传的仍是这个 ID。
后台任务用虚拟线程执行,一份文档一个线程:
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
public class ExecutorServiceUtils {
/**
* 文档解析专用执行器。
*
* <p>解析一份大文档要几分钟且几乎全在等远程 OCR,用每任务一个虚拟线程的线程池,避免占用固定池的
* 100 个工作线程,也避免多份大文档互相顶掉。
*/
private static final ExecutorService documentParseExecutor =
Executors.newThreadPerTaskExecutor(Thread.ofVirtual().name("doc-parse-", 0).factory());
public static ExecutorService getDocumentParseExecutor() {
return documentParseExecutor;
}
}
3. 任务状态表
任务表只保存进度与结果,不参与文档、分段的正规存储:
-- 分段预览任务:上传后立刻返回 task_id,解析在后台线程里按页推进,
-- 前端轮询这个表拿到进度或最终分段。result 直接存前端要用的分段列表。
CREATE TABLE IF NOT EXISTS "public"."moss_kb_document_split_task" (
"id" BIGINT NOT NULL PRIMARY KEY,
"user_id" BIGINT,
"file_id" BIGINT,
"file_name" VARCHAR NOT NULL,
"file_size" BIGINT,
"status" VARCHAR(16) NOT NULL DEFAULT 'running',
"progress" SMALLINT NOT NULL DEFAULT 0,
"total" INT NOT NULL DEFAULT 0,
"result" JSONB,
"error_message" TEXT,
"creator" VARCHAR(64) DEFAULT '',
"create_time" TIMESTAMP WITHOUT TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP,
"updater" VARCHAR(64) DEFAULT '',
"update_time" TIMESTAMP WITHOUT TIME ZONE NOT NULL DEFAULT CURRENT_TIMESTAMP,
"deleted" SMALLINT NOT NULL DEFAULT 0,
"tenant_id" BIGINT NOT NULL DEFAULT 0
);
CREATE INDEX IF NOT EXISTS "moss_kb_document_split_task_user_id" ON "public"."moss_kb_document_split_task" USING btree ("user_id");
| 字段 | 含义 |
|---|---|
status | running、success、failed |
progress | 已完成百分比,0—99 表示还在解析,成功时写 100 |
total | 总页数,非分页文档为 1 |
result | 成功后保存前端要用的分段列表 |
error_message | 失败原因,直接返回给界面提示 |
进度写库按 5 页或 5% 节流,逐页回调不会变成逐页写库:
/** 记录已完成页数。 */
public void markProgress(Long taskId, int completed, int total) {
if (total <= 0) {
return;
}
int percent = (int) Math.min(99, Math.round(completed * 100.0 / total));
boolean lastPage = completed >= total;
// 页数多时按页步进节流,页数少时按百分比节流,两种情况都能看到推进。
boolean reachStep = completed % PROGRESS_STEP == 0
|| percent >= currentPercent(taskId) + PROGRESS_PERCENT_STEP;
if (!lastPage && !reachStep) {
return;
}
Db.update("update moss_kb_document_split_task set progress=?, update_time=now() where id=?", percent, taskId);
}
任务按 user_id 校验归属,status 和 error_message 只在任务失败时对外返回:
@Get("/split/task/{taskId}")
public ResultVo splitTask(Long taskId) {
Long userId = TioRequestContext.getUserIdLong();
MossKbDocumentSplitTaskVo task = Aop.get(MossKbDocumentSplitTaskService.class).get(userId, taskId);
if (task == null) {
return ResultVo.fail("任务不存在或无权访问");
}
if (MossKbDocumentSplitTaskService.STATUS_FAILED.equals(task.getStatus())) {
return ResultVo.fail(task.getErrorMessage() == null ? "文档解析失败" : task.getErrorMessage());
}
if (MossKbDocumentSplitTaskService.STATUS_SUCCESS.equals(task.getStatus())) {
return ResultVo.ok(task.getResult() == null ? new ArrayList<>() : task.getResult());
}
return ResultVo.ok();
}
4. 两个接口的约定
| 接口 | 说明 |
|---|---|
POST /api/dataset/document/split | 保存文件并建任务,返回 task_id_list 与 failures |
GET /api/dataset/document/split/task/{taskId} | 查询任务状态 |
状态接口的响应按 status 分成三种,前端据此决定继续轮询还是渲染分段:
| 状态 | code | data | 前端行为 |
|---|---|---|---|
running | 200 | 空 | 继续轮询 |
success | 200 | 分段列表 | 渲染分段预览 |
failed | 400 | 空,message 是失败原因 | 提示失败并结束 |
success 时的 data 与原来的同步响应完全一致,仍然是 name、id、content、parse_strategy、page_count 组成的数组,因此提交分段的后续流程不需要改动。分段之外还会带上解析信息:
[
{
"name": "示例文件.pdf",
"id": 696358168286838784,
"content": [{ "title": "", "content": "> Page 1 ..." }],
"parse_strategy": "pdf-mixed",
"page_count": 8,
"ocr_page_count": 1,
"failed_page_count": 0
}
]
5. 前端轮询
接口层把“上传 + 轮询 + 取结果”封装成一个方法,页面代码保持不变:
import { Result } from '@/request/Result'
import { get, post } from '@/request/index'
const prefix = '/dataset'
/**
* 分段预览(上传文档)
*
* 后端只负责保存文件并建解析任务,立即返回 { task_id_list };分段结果由本方法内部轮询
* /dataset/document/split/task/{taskId} 取得,调用方拿到的仍然是分段列表。
*/
const postSplitDocument: (data: any) => Promise<Result<any>> = async (data) => {
const started: any = await post(
`${prefix}/document/split`,
data,
undefined,
undefined,
1000 * 60 * 10
)
const taskIdList: Array<string> = started?.data?.task_id_list || []
if (!taskIdList.length) {
return started
}
const lists = await Promise.all(taskIdList.map((taskId) => pollSplitTask(taskId)))
return { ...started, data: lists.flat() } as Result<any>
}
/** 解析一份大文档要几分钟,所以按 2 秒间隔轮询任务状态,最多等待 60 分钟。 */
const pollSplitTask: (taskId: string) => Promise<Array<any>> = async (taskId) => {
const deadline = Date.now() + 1000 * 60 * 60
for (;;) {
try {
const res: any = await get(`${prefix}/document/split/task/${taskId}`, {}, undefined, {
silent: true
})
if (Array.isArray(res?.data)) {
return res.data
}
} catch (error: any) {
// 后端返回的失败原因需要提示用户,但统一拦截器已经弹过一次,这里不再重复。
throw new Error(error?.message || '文档解析失败')
}
if (Date.now() > deadline) {
throw new Error('文档解析超时,请稍后在文档列表查看解析结果')
}
await new Promise((resolve) => setTimeout(resolve, 2000))
}
}
轮询请求使用 silent: true:任务失败时后端返回的就是给用户看的失败原因,请求层已经弹过一次提示,轮询循环不再重复弹出。多文件上传对每个任务并行轮询,全部完成后再一次性渲染,保持原来的界面行为。
SetRules.vue 里的 splitDocument() 完全不需要修改,它仍然只关心 res.data 是分段列表:
function splitDocument() {
loading.value = true
let fd = new FormData()
documentsFiles.value.forEach((item) => {
if (item?.raw) {
fd.append('file', item?.raw)
}
})
documentApi
.postSplitDocument(fd)
.then((res: any) => {
paragraphList.value = res.data
loading.value = false
})
.catch(() => {
loading.value = false
})
}
6. 解析本身的改动
异步化解决了“等不起”,但整份解析仍然要面对远程 OCR 的偶发失败。解析器做了三处调整。
6.1 空白页不再作废整份文档
远程 OCR 对空白页、纯图片页会正常返回空正文。这曾经直接抛 OCR未返回正文,让一份 35 页的文档只因为最后一页是空白就整份失败。现在空结果按空页处理,页码位置照常保留:
protected String recognize(byte[] data, String filename) throws Exception {
GiteeClient client = new GiteeClient();
GiteeDocumentParseRequest request = new GiteeDocumentParseRequest();
request.setModel(OCR_MODEL);
request.setPrompt(GiteePromptConst.pdf_to_markdown_prompt);
request.setInclude_image(false);
request.setInclude_image_base64(false);
GiteeTaskResponse task = client.parseDocument(data, filename, request);
long deadline = System.nanoTime() + java.util.concurrent.TimeUnit.MINUTES.toNanos(8);
while (!Arrays.asList("success", "succeeded", "completed").contains(String.valueOf(task.getStatus()).toLowerCase(Locale.ROOT))) {
if (task.getStatus() == null || task.getTask_id() == null) {
throw new IOException("OCR服务未返回有效任务状态");
}
if (Arrays.asList("failed", "failure", "cancelled", "canceled").contains(task.getStatus().toLowerCase(Locale.ROOT))) {
throw new IOException("OCR任务失败: " + task.getTask_id());
}
if (System.nanoTime() > deadline) {
throw new IOException("OCR任务超时: " + task.getTask_id());
}
Thread.sleep(2000); task = client.getTask(task.getTask_id());
}
// 扫描页可能是空白页或整页图片,服务会正常返回空正文,这里按空内容处理,由调用方决定是否记录。
String markdown = textOnly(GiteeSimpleMarkdownUtils.toMarkdown(task.getOutput(), null));
if (markdown == null || markdown.isBlank()) {
log.info("OCR returned no text for {}", filename);
return "";
}
return markdown;
}
真正的失败(无效任务状态、任务失败、超时)仍然抛异常,交给重试处理;只有“服务正常但没识别出文字”被当成空页。
6.2 失败按退避重试,失败页单独记账
单页 OCR 失败按 1 秒、2 秒退避重试,最多三次;三次都失败就记为空页并继续:
/**
* OCR 重试:失败时按 1s、2s 退避重试,最后一次仍失败就记为失败页。
*
* <p>“未识别到内容”是正常返回,不是异常,因此不在这里重试。
*
* @return 是否拿到结果
*/
private boolean withOcrRetry(OcrCall call, String[] result, boolean[] failed, String describe) {
for (int attempt = 1; attempt <= OCR_MAX_ATTEMPTS; attempt++) {
try {
result[0] = call.recognize();
return true;
} catch (Exception e) {
log.warn("OCR attempt {}/{} failed for {}: {}", attempt, OCR_MAX_ATTEMPTS, describe, e.getMessage());
if (attempt < OCR_MAX_ATTEMPTS) {
try {
Thread.sleep(1000L * attempt);
} catch (InterruptedException interrupted) {
Thread.currentThread().interrupt();
failed[0] = true;
return false;
}
}
}
}
log.warn("OCR gave up for {}", describe);
failed[0] = true;
return false;
}
解析结果把“走 OCR 的页数”和“最终失败的页数”一起返回,策略标识会带上 -partial 后缀,前端预览里也能看到这两项:
| 情况 | 策略标识 |
|---|---|
| 全是文本页 | pdf-text |
| 全是扫描页,没有失败页 | pdf-ocr |
| 混合文档,没有失败页 | pdf-mixed |
| 有页重试后仍失败 | pdf-ocr-partial 或 pdf-mixed-partial |
如果一份文档每一页 OCR 都失败,返回值就不再是“未提取到内容”,而是明确提示可以重试:
if (result.text().isBlank()) {
// 全部走 OCR 的文档如果每一页都失败,原因是远程服务不可用,不是文档本身没有内容。
if (result.ocrPages() > 0 && result.ocrPages() == result.failedPages()) {
throw new IOException("扫描页OCR识别失败,请稍后重试");
}
throw new IllegalArgumentException("未提取到内容,请检查文档或使用扫描PDF");
}
6.3 逐页并发,页序不变
一份 35 页的扫描 PDF 如果逐页串行调用 OCR,光等待就要几分钟。页与页之间没有依赖,因此改为并发提取、按页码回填,正文顺序与串行时完全一致:
/** 并发提取的页面数:OCR 是网络调用,页与页之间并发即可,不需要按 CPU 核数拉满。 */
private static final int OCR_PAGE_PARALLELISM = Math.max(2, Math.min(4, Runtime.getRuntime().availableProcessors()));
/** 逐页提取 PDF:文本层优先,扫描页走 OCR,页面之间并发处理但按原页序输出。 */
private Parsed parsePdf(byte[] data, Consumer<Progress> progress) throws Exception {
try (PDDocument pdf = PDDocument.load(data)) {
int totalPages = pdf.getNumberOfPages();
PDFTextStripper stripper = new PDFTextStripper();
stripper.setSortByPosition(true);
String sourceHash = Md5Utils.md5Hex(data);
PageResult[] results = new PageResult[totalPages];
AtomicInteger completed = new AtomicInteger();
ForkJoinPool pool = new ForkJoinPool(OCR_PAGE_PARALLELISM);
try {
List<ForkJoinTask<?>> tasks = new ArrayList<>(totalPages);
for (int page = 1; page <= totalPages; page++) {
int pageNumber = page;
tasks.add(pool.submit(() -> {
try {
results[pageNumber - 1] = extractPdfPage(pdf, stripper, sourceHash, pageNumber);
} catch (IOException e) {
throw new UncheckedIOException("解析第 " + pageNumber + " 页失败", e);
}
if (progress != null) {
progress.accept(new Progress(completed.incrementAndGet(), totalPages));
}
}));
}
for (ForkJoinTask<?> task : tasks) {
task.join();
}
} finally {
pool.shutdown();
}
PDFTextStripper 和 PDDocument 都不是线程安全的,所以提取文本层、序列化单页这两步加锁串行,耗时的网络 OCR 在锁外并发执行。
6.4 页面缓存键与原文件绑定
序列化单页时 PDFBox 会写入新的内部文档 ID,同一页重复序列化的字节并不相同,直接拿这些字节做缓存键会让缓存永远不命中。缓存键固定为“原文件内容摘要 + 页码”:
// 缓存键固定用“原文件内容 + 页码”:序列化单页时 PDFBox 会写入新的内部文档 ID,
// 同一页重复序列化的字节并不相同,用这些字节做键会让缓存永远不命中。
boolean[] failed = new boolean[1];
String[] recognized = new String[1];
boolean ok = withOcrRetry(() -> ocrPdfPage(sourceHash, pageNumber, pageData), recognized, failed,
sourceHash + ":page:" + pageNumber + ":" + OCR_CACHE_VERSION);
return new PageResult(pageNumber, ok ? recognized[0] : "", true, failed[0]);
同一份内容并发出现时(例如几页完全相同的空白页)按内容加锁,后到的页等第一次的结果落库后直接读缓存,不重复付费调用:
private String ocrWithCache(byte[] data, String filename, String hash) throws Exception {
String cached = readOcrCache(hash);
if (cached != null && !cached.isBlank()) {
return textOnly(cached);
}
synchronized (OCR_LOCK_STRIPES[Math.floorMod(hash.hashCode(), OCR_LOCKS)]) {
// 同一份内容的第一个请求可能已经写好缓存,这里再查一次就能省掉重复的付费调用。
cached = readOcrCache(hash);
if (cached != null && !cached.isBlank()) {
return textOnly(cached);
}
String markdown = recognize(data, filename);
if (markdown != null && !markdown.isBlank()) {
writeOcrCache(hash, markdown);
}
return markdown;
}
}
空白结果不写缓存,避免一次失败长期生效;下次请求会重新提交。
7. 上传体积上限必须写进本地配置
扫描版 PDF 动辄十几 MB,而 t-io 默认只接收 2 MB 的请求体,超出时在解析请求头阶段就直接断开连接,表现为 Request body exceeds the configured limit,浏览器端看到的是连接被重置而不是业务错误。
t-io 从 http.multipart.max-request-size 与 http.multipart.max-file-size 读取上限,所以这两项要和 server.port、数据库连接一起放在本地配置文件里:
server.port=10060
app.env=dev
# 上传体积上限:扫描版 PDF 动辄十几 MB,t-io 默认只收 2MB。
http.multipart.max-request-size=110100480
http.multipart.max-file-size=104857600
默认值分别对应 105 MB 的请求体和 100 MB 的单文件,比解析器自身的 100 MB 限制略大,让超限请求能在业务层得到明确的提示文本。
8. 验证
用两份真实的政法政策 PDF 走完整链路,过程中不跳过 OCR:
| 文件 | 页数 | 结果 |
|---|---|---|
| 扫描件为主,约 12 MB | 35 | 上传 1.5 秒返回任务 ID;pdf-mixed,28 页走 OCR、7 页走文本层;约 2 分钟解析出 14 个分段,失败页 0 |
| 文字为主,约 220 KB | 8 | 上传 0.1 秒返回任务 ID;pdf-mixed,1 页走 OCR;解析出 2 个分段 |
12 MB 那份文档的第 35 页是空白页,远程 OCR 对它返回空正文。改动之前,这一页会让整份 35 页的文档直接失败并返回 OCR未返回正文;现在空页按空页处理,其余 34 页的正文与分段正常产出。
任务表里能看到推进过程与最终结果:
select id, file_name, status, progress, total, (result is not null) as has_result, error_message
from moss_kb_document_split_task
order by id desc limit 5;
id | file_name | status | progress | total | has_result | error_message
--------------------+-------------+---------+----------+-------+------------+---------------
696360496210333696 | sample1.pdf | success | 100 | 35 | t |
696358168358141952 | sample2.pdf | success | 100 | 8 | t |
同一份文件第二次预览会全部命中页面缓存,不再产生 OCR 调用:缓存键是「原文件内容 + 页码」,与序列化单页时生成的临时文档 ID 无关。
界面上的分段预览与以前一致:分段列表、解析策略和页数都照常显示,区别只在于等待期间进度是持续推进的,而且关掉页面再回来仍然可以取到同一个任务的结果。
