Skip to content

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.jsonchunks.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 各自处理。一个文档失败不会阻塞其他文档。批量接口的返回值包含每个文档的 idstatus,前端可以逐个追踪。

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 等各种格式转成纯文本,这是处理管道的第一步。

基于 MIT 协议开源