主题
14.03-文档上传流程
要点
- 文档上传的核心挑战不在接收文件,在于接收之后的异步处理链——解析、切分、向量化、入库每一步都可能耗时数十秒到数分钟
- 上传接口和处理 Worker 必须分离:上传立即返回
document_id,处理在后台异步进行 - 状态机是管道的骨架——合法的状态转换表约束了所有操作边界,也是断点续传和前端轮询的基础
- 每个处理步骤的中间产物必须立即持久化——这是断点续传的前提,不是优化
- 队列选型按规模决策:日处理量几百份以内用数据库轮询够用,之上切 BullMQ
- 删除文档时需要同步清理对象存储、向量数据库和关系型数据库,三者一致是真正的难点
1. 同步处理为什么行不通
最直觉的文档上传:接收文件,同步走完解析、切分、向量化、入库全流程,然后返回结果。
typescript
// ❌ 同步处理——大文档会超时
app.post('/documents', async (c) => {
const file = await c.req.parseBody()
const content = await (file['file'] as File).text()
// 同步走完全部流程:解析 → 切分 → 向量化 → 入库
const parsed = parseDocument(content)
const chunks = chunkText(parsed.content)
const embeddings = await callEmbeddingAPI(chunks.map(c => c.text))
await vectorDB.upsert(/* ... */)
return c.json({ ok: true })
})这个实现有五个绕不过去的硬伤:
大文档会超时。 10MB 的 PDF 解析可能要 30 秒,向量化要几分钟。HTTP 请求等不了这么久——网关、代理、客户端任何一层都会先超时。
没有状态追踪。 用户发起上传后只能干等,不知道处理到哪一步了,也不知道还要等多久。
失败代价是全部重来。 处理到第 800 个 chunk 时 embedding API 超时,之前所有的解析和切分计算全部浪费,用户只能重新上传。
没有资源控制。 没有文件大小限制、格式校验和去重机制,也没有用户配额——恶意用户可以上传超大文件撑爆存储。
没有并发保护。 同一份文档可以被反复上传,产生重复向量。
解决方案是把上传和处理拆开。 上传接口只负责接收文件、创建记录、立即返回 document_id。后续处理交给后台 Worker 异步完成。用户通过轮询或 WebSocket 追踪状态。
这样做的关键收益:上传接口的响应时间与文件大小、处理耗时完全解耦。处理失败不影响上传。每个步骤有独立状态记录,失败后可以从断点恢复。
2. 异步管道与状态机
2.1 管道架构
把上传和处理拆开之后,整体架构变成一条异步管道:
用户上传 ─→ 接收文件 ─→ 返回 document_id ─→ 完成(HTTP 202)
│
▼
加入处理队列
│
▼
┌─────── 后台 Worker ───────┐
│ 1. 解析文档 → parsed │
│ 2. 切分文本 → chunks │
│ 3. 向量化 → embeddings │
│ 4. 写入向量库 │
│ 5. 更新状态 │
└───────────────────────────┘
│
▼
状态变更通知上传和处理之间通过队列解耦。Worker 按自己的节奏消费任务,不会因为上传量大而被压垮。
2.2 状态机
文档从上传到完成(或失败),经历一条线性的状态链:
pending → parsing → chunking → embedding → indexing → indexed
↘
任何阶段都可能 → failed → pending(重试)每个状态对应管道中的一个具体步骤。状态转换只能沿箭头方向前进或标记为失败,不能跳跃或逆序。
typescript
// src/services/rag/document-states.ts
export type DocumentStatus =
| 'pending' // 已上传,等待处理
| 'parsing' // 解析中
| 'chunking' // 切分中
| 'embedding' // 向量化中
| 'indexing' // 写入向量库中
| 'indexed' // 完成(终态)
| 'failed' // 失败(可重试)
export type Document = {
id: string
userId: string
fileName: string
fileSize: number
mimeType: string
status: DocumentStatus
error?: string
chunksCount?: number
createdAt: string
updatedAt: string
}
// 合法的状态转换表
const VALID_TRANSITIONS: Record<DocumentStatus, DocumentStatus[]> = {
pending: ['parsing', 'failed'],
parsing: ['chunking', 'failed'],
chunking: ['embedding', 'failed'],
embedding: ['indexing', 'failed'],
indexing: ['indexed', 'failed'],
indexed: [], // 终态,不再转换
failed: ['pending'], // 重试:回到 pending
}
export async function transitionStatus(
docId: string,
newStatus: DocumentStatus,
error?: string
) {
const doc = await getDocument(docId)
if (!VALID_TRANSITIONS[doc.status].includes(newStatus)) {
throw new Error(
`Invalid transition: ${doc.status} → ${newStatus}`
)
}
await db.query(
`UPDATE documents
SET status = ?, error = ?, updated_at = ?
WHERE id = ?`,
[newStatus, error ?? null, new Date().toISOString(), docId]
)
}状态机看起来简单,但它解决了文档处理中最容易出的问题——状态混乱。没有这个约束,重试可能在还在处理时就触发,失败后无法回到初始状态重新来,同一文档被并发处理导致向量库出现重复数据。
状态转换表就是这份管道的操作契约。 它定义了所有合法操作,也是前端判断「当前能做什么」的依据。
3. 上传接口:只做 I/O,不做计算
上传接口的设计原则:快速接收文件,快速返回 ID。不要在上传请求里做任何耗时操作——解析、切分、向量化全部交给后台 Worker。
typescript
// src/routes/documents.ts
import { Hono } from 'hono'
import { nanoid } from 'nanoid'
const MAX_FILE_SIZE = 50 * 1024 * 1024 // 50MB
const ALLOWED_MIME_TYPES = [
'application/pdf',
'application/vnd.openxmlformats-officedocument.wordprocessingml.document',
'text/plain',
'text/markdown',
'text/html',
]
const documents = new Hono()
documents.post('/', async (c) => {
const user = c.get('user')
// 1. 配额检查
const quota = await getUserQuota(user.id)
if (quota.used >= quota.limit) {
return c.json({ error: '配额已满' }, 403)
}
// 2. 解析文件
const body = await c.req.parseBody()
const file = body['file'] as File
if (!file) {
return c.json({ error: '缺少文件' }, 400)
}
// 3. 格式校验——只接受白名单内的 MIME 类型
if (!ALLOWED_MIME_TYPES.includes(file.type)) {
return c.json(
{ error: `不支持的文件格式: ${file.type}` }, 400
)
}
// 4. 大小校验
if (file.size > MAX_FILE_SIZE) {
return c.json(
{ error: `文件过大,最大 ${MAX_FILE_SIZE / 1024 / 1024}MB` },
400
)
}
// 5. 创建文档记录
const docId = nanoid()
await db.query(
`INSERT INTO documents
(id, user_id, file_name, file_size, mime_type, status, created_at)
VALUES (?, ?, ?, ?, ?, 'pending', ?)`,
[docId, user.id, file.name, file.size, file.type,
new Date().toISOString()]
)
// 6. 存储原始文件到对象存储
const fileBuffer = await file.arrayBuffer()
await storage.put(
`documents/${docId}/raw`,
new Uint8Array(fileBuffer)
)
// 7. 加入处理队列(不阻塞响应)
await enqueueProcessing(docId)
// 8. 立即返回 202 Accepted
return c.json({
id: docId,
status: 'pending',
message: '文档已上传,正在处理中',
}, 202)
})
// 查询文档状态
documents.get('/:id', async (c) => {
const doc = await getDocument(c.req.param('id'))
if (!doc) return c.json({ error: 'Not found' }, 404)
if (doc.userId !== c.get('user').id) {
return c.json({ error: 'Forbidden' }, 403)
}
return c.json(doc)
})这个接口里没有任何耗时操作——只有数据库写入、对象存储上传和队列入队三个 I/O 动作。响应时间通常在 1 秒以内,和文件大小关系不大。
生产环境注意: 示例中文件存在本地
storage,生产环境应换成 S3、Cloudflare R2 或 MinIO 等对象存储,路径规则不变。
常见错误: 如果 enqueueProcessing 失败(比如队列服务宕机),文档记录已经创建但不会有 Worker 处理。需要在入队失败时回滚——把文档状态设为 failed 并记录错误信息,返回 500 让前端提示重试。
4. 后台 Worker:逐步推进,每步持久化
Worker 从队列拿到任务后,按状态链逐步推进。关键设计:每步的中间产物立即持久化到对象存储,不要只存在内存里。
typescript
// src/workers/document-processor.ts
export async function processDocument(docId: string) {
const doc = await getDocument(docId)
if (!doc) return
try {
// ── 步骤 1:解析 ──
await transitionStatus(docId, 'parsing')
const rawFile = await storage.get(`documents/${docId}/raw`)
const parsed = await parseDocument(rawFile, doc.mimeType)
// 持久化解析结果——失败恢复的依赖
await storage.put(
`documents/${docId}/parsed.json`,
new TextEncoder().encode(JSON.stringify(parsed))
)
// ── 步骤 2:切分 ──
await transitionStatus(docId, 'chunking')
const chunks = chunkText(parsed.content, {
chunkSize: 1000,
overlap: 200,
})
await storage.put(
`documents/${docId}/chunks.json`,
new TextEncoder().encode(JSON.stringify(chunks))
)
// ── 步骤 3:向量化(分批,避免 API 限流)──
await transitionStatus(docId, 'embedding')
const BATCH_SIZE = 50
const allEmbeddings: number[][] = []
for (let i = 0; i < chunks.length; i += BATCH_SIZE) {
const batch = chunks.slice(i, i + BATCH_SIZE)
const embeddings = await callEmbeddingAPI(
batch.map((c) => c.text)
)
allEmbeddings.push(...embeddings)
}
// ── 步骤 4:写入向量库 ──
await transitionStatus(docId, 'indexing')
const vectors = chunks.map((chunk, i) => ({
id: `${docId}-chunk-${i}`,
vector: allEmbeddings[i],
metadata: {
text: chunk.text,
documentId: docId,
documentTitle: parsed.title,
chunkIndex: i,
source: doc.fileName,
userId: doc.userId,
},
}))
await vectorDB.upsert(vectors)
// ── 步骤 5:标记完成 ──
await db.query(
`UPDATE documents
SET status = 'indexed', chunks_count = ?, updated_at = ?
WHERE id = ?`,
[chunks.length, new Date().toISOString(), docId]
)
} catch (err) {
const message = err instanceof Error
? err.message : 'Unknown error'
await transitionStatus(docId, 'failed', message)
console.error(`Document ${docId} failed:`, err)
}
}每步完成后把中间产物写入对象存储(parsed.json、chunks.json),而不是只保存在变量里。这是断点续传的前提——第 6 节会展开。
为什么分批向量化? Embedding API 通常有批次限制(比如 OpenAI 单次最多 2048 个 input)。BATCH_SIZE = 50 是保守值,既能减少 API 调用次数,又不会触发限流。实际值需要根据你使用的模型调整。
5. 队列选型:两条路径
队列有两种实用方案,按规模选择:
| 维度 | 数据库轮询 | Redis / BullMQ |
|---|---|---|
| 复杂度 | 低,零额外依赖 | 中,需要 Redis |
| 延迟 | 秒级(取决于轮询间隔) | 毫秒级 |
| 适合规模 | 日处理量 < 几百份 | 日处理量 > 几千份 |
| 生产特性 | 无 | 优先级、重试、死信队列 |
5.1 数据库轮询(小规模可用)
typescript
// 入队:写入待处理记录
async function enqueueProcessing(docId: string) {
await db.query(
`INSERT INTO processing_queue
(document_id, status, created_at)
VALUES (?, 'pending', ?)`,
[docId, new Date().toISOString()]
)
}
// Worker 定时轮询
setInterval(async () => {
const pending = await db.query(
`SELECT document_id FROM processing_queue
WHERE status = 'pending'
ORDER BY created_at LIMIT 1`
)
if (pending.length === 0) return
const docId = pending[0].document_id
await db.query(
`UPDATE processing_queue SET status = 'processing'
WHERE document_id = ?`,
[docId]
)
await processDocument(docId)
await db.query(
`UPDATE processing_queue SET status = 'done'
WHERE document_id = ?`,
[docId]
)
}, 5000) // 每 5 秒轮询一次数据库轮询的代价:轮询间隔决定了任务等待的下限延迟,频繁轮询会增加数据库负载。小规模应用(每天几百个文档以内)完全够用。
5.2 BullMQ(生产环境推荐)
typescript
import { Queue, Worker } from 'bullmq'
const ragQueue = new Queue('rag-processing', {
connection: { host: 'localhost', port: 6379 }
})
async function enqueueProcessing(docId: string) {
await ragQueue.add('process-document', { docId }, {
attempts: 3, // 最多重试 3 次
backoff: { type: 'exponential', delay: 1000 },
})
}
const worker = new Worker('rag-processing', async (job) => {
await processDocument(job.data.docId)
}, {
connection: { host: 'localhost', port: 6379 },
concurrency: 5, // 最多同时处理 5 个文档
})BullMQ 提供生产环境需要的关键能力:优先级队列(付费用户上传优先)、自动重试(embedding API 限流时指数退避)、并发控制(concurrency: 5 避免打爆 API 配额)、死信队列(重试耗尽的进 dead letter,人工介入)。
示例可用 vs 生产可用: 数据库轮询是示例可用的起点。当你遇到轮询延迟影响用户体验、需要优先级或重试策略、数据库查询成为瓶颈这三个问题中的任何一个,就该切 BullMQ。好消息是:无论选哪种队列,Worker 代码(
processDocument)完全不变——变的只是任务来源和结果存储。
6. 断点续传:示例可用与生产可用的分界线
大文档处理可能耗时几分钟。如果 embedding 处理到一半失败(API 限流、网络中断),不应该从解析重新开始。
断点续传的核心思路: Worker 启动时先检查当前状态和已有的中间产物,从上次中断的步骤继续。
typescript
async function processDocumentWithResume(docId: string) {
const doc = await getDocument(docId)
const progress = await getProgress(docId)
const BATCH_SIZE = 50
// 从上次中断的位置继续
if (progress?.currentStep === 'embedding') {
// 解析和切分已完成,直接读缓存的 chunks
const chunks = JSON.parse(
new TextDecoder().decode(
await storage.get(`documents/${docId}/chunks.json`)
)
)
await transitionStatus(docId, 'embedding')
const embeddings: number[][] = []
// 从上次中断的 batch 位置继续
for (
let i = progress.processedChunks;
i < chunks.length;
i += BATCH_SIZE
) {
const batch = chunks.slice(i, i + BATCH_SIZE)
const batchEmbeddings = await callEmbeddingAPI(
batch.map((c) => c.text)
)
embeddings.push(...batchEmbeddings)
// 每批完成后更新进度
await updateProgress(docId, {
processedChunks: i + batch.length,
})
}
// 继续写入向量库 ...
} else {
// 没有可恢复的进度,从头开始
await processDocument(docId)
}
}为什么断点续传重要? 假设一份 5000 chunk 的文档处理到第 4000 个时 embedding API 超时。没有断点续传,用户必须重新上传、重新解析、重新切分、从第 1 个 chunk 重新向量化。有了断点续传,Worker 从第 4000 个继续,前面 4000 个的计算结果不会浪费。
生产环境注意: 示例中的进度追踪只记录了「处理到第几个 chunk」。生产环境还需要处理:
- 重试策略:embedding API 失败通常是限流,应该用指数退避(1s → 2s → 4s),最多重试 3 次。直接重试大概率继续失败
- 幂等性:每批 embedding 完成后立即持久化该批结果,而不是全部完成后才写入。否则中途失败,已完成的 batch 结果也会丢失
- 失败上限:同一个文档重试超过 N 次(比如 3 次)后,标记为
failed并通知用户,不要无限重试
7. 文档生命周期:批量上传、删除与清理
7.1 批量上传
用户通常需要一次上传多个文档。批量接口为每个文件独立创建记录和入队,单个失败不影响其他文档。
typescript
documents.post('/batch', async (c) => {
const user = c.get('user')
const body = await c.req.parseBody()
const files = Object.entries(body)
.filter(([key]) => key.startsWith('file'))
.map(([, file]) => file as File)
if (files.length === 0) {
return c.json({ error: '没有文件' }, 400)
}
if (files.length > 20) {
return c.json(
{ error: '单次最多上传 20 个文件' }, 400
)
}
const results = []
for (const file of files) {
// 校验、存储、入队——复用单个上传的逻辑
if (!ALLOWED_MIME_TYPES.includes(file.type)) continue
if (file.size > MAX_FILE_SIZE) continue
const docId = nanoid()
await db.query(/* ... */)
await storage.put(
`documents/${docId}/raw`,
new Uint8Array(await file.arrayBuffer())
)
await enqueueProcessing(docId)
results.push({
id: docId, fileName: file.name, status: 'pending'
})
}
return c.json({ documents: results }, 202)
})关键设计: 每个文档独立入队,Worker 各自处理。一个文档失败不会阻塞其他文档。批量接口的返回值包含每个文档的 id 和 status,前端可以逐个追踪。
7.2 删除文档:跨三个存储层的清理
删除文档不是删一行数据库记录。一份文档的数据分散在三个地方:对象存储(原始文件、中间产物)、向量数据库(所有 chunk 的向量)、关系型数据库(文档元数据)。三者必须同时清理。
typescript
documents.delete('/:id', async (c) => {
const doc = await getDocument(c.req.param('id'))
if (!doc) return c.json({ error: 'Not found' }, 404)
if (doc.userId !== c.get('user').id) {
return c.json({ error: 'Forbidden' }, 403)
}
// 1. 从向量库删除所有相关向量
// 大多数向量数据库支持按 metadata 过滤删除
await vectorDB.deleteByFilter({ documentId: doc.id })
// 2. 删除对象存储中的文件
await storage.delete(`documents/${doc.id}/raw`)
await storage.delete(`documents/${doc.id}/parsed.json`)
await storage.delete(`documents/${doc.id}/chunks.json`)
// 3. 删除数据库记录
await db.query('DELETE FROM documents WHERE id = ?', [doc.id])
return c.json({ ok: true })
})删除顺序很重要: 先删向量库 → 再删对象存储 → 最后删数据库记录。如果中途失败,至少数据库记录还在,可以重试清理。反过来如果先删数据库,向量库的孤儿数据就再也找不到关联了。
常见错误: 如果向量数据库不支持按 metadata 过滤删除,需要在写入向量时把 chunk ID 记录到关系型数据库,删除时先查出所有 chunk ID 再逐个删除。向量库清理失败时,建议的做法是标记文档为「删除中」,由后台定时任务重试清理,而不是直接报错回滚。
8. 验收清单
到这里,文档上传管道的完整链路已经走完。用以下清单验收你的实现:
- 上传接口在 1 秒内返回
document_id和 202 状态码,不阻塞处理 - 状态转换严格遵循
pending → parsing → chunking → embedding → indexing → indexed,任何跳跃或逆序都被拒绝 - 每步状态变更都有记录,前端可以通过
GET /documents/:id实时查询 - Worker 失败后文档状态变为
failed,错误信息写入error字段 - 重试从上次中断的步骤继续,不重复已完成的解析和切分
- 删除文档同时清理对象存储、向量库和数据库,三者一致
- 批量上传单个失败不影响其他文档
前端轮询建议: 上传后用 GET /documents/:id 轮询状态,间隔 2 秒、超时 5 分钟。更好的方式是 SSE 推送状态变更,避免无效请求。
文档上传完成后,文档进入向量库,可以被检索了。下一篇讲文档解析——怎么把 PDF、Word、HTML、Markdown 等各种格式转成纯文本,这是处理管道的第一步。