AI 应用后端工程化:从原型到可交付系统 · 第 11 篇 · 第三章 · 数据管道
把耗时文档处理移出请求线程,并用任务状态、缓存和配置边界承接失败与重试。
创建知识库的接口如果同步完成文档解析、切片、嵌入和索引,用户可能等几十秒甚至几分钟。把函数加上“异步任务”装饰器只是第一步;请求离开后,系统还要描述任务状态、处理重复消息、保存错误,并让用户知道什么时候可以使用结果。
请求与工作分开
sequenceDiagram
participant C as 客户端
participant A as API
participant D as 数据库
participant Q as 队列
participant W as Worker
C->>A: 创建数据集
A->>D: 写入 pending 记录
A->>Q: 发布任务 ID
A-->>C: 202 Accepted
Q->>W: 投递任务
W->>D: 标记 processing
W->>W: 解析、嵌入、索引
W->>D: 标记 ready / failed
C->>A: 查询状态
A-->>C: 当前进度
API 先保存意图,再返回资源 ID。Worker 不依赖原始 HTTP 请求,而是通过任务 ID重新加载数据。这样消息体较小,敏感配置也不会长时间滞留在队列中。
队列提供投递,不提供业务正确性
消息系统通常承诺“至少投递一次”,Worker 可能在完成处理后、确认消息前崩溃,同一任务便会再次执行。每一步必须可重试:写片段使用稳定 ID,索引写入采用 upsert,状态转换检查任务版本。
def process_document(job_id: str, repository, index):
job = repository.lock(job_id)
if job.status == "ready":
return
if job.status not in {"pending", "retrying"}:
raise RuntimeError("job is not runnable")
repository.mark_processing(job_id, job.attempt + 1)
chunks = parse_and_chunk(job.file_id)
index.upsert(job.document_id, chunks)
repository.mark_ready(job_id, len(chunks))
这段代码仍需要异常时标记失败,但核心意思是:先读持久状态,再决定是否执行。不能假设“队列发来一次,函数就只运行一次”。
Redis 的三个不同角色
Redis 可能同时被用作队列 Broker、结果后端和应用缓存。三个角色的可靠性要求不同。队列消息关系到任务是否丢失,缓存丢失通常只影响性能,结果后端决定用户能否查询状态。把它们混用同一 key 前缀和过期策略,会让一次清缓存误删任务。
资源允许时应使用独立数据库或至少独立命名空间、权限和监控。缓存键要包含版本;任务状态不能只存在 Redis,若需要长期审计,应落数据库。
进度不是随便填一个百分比
文档处理各阶段耗时差异很大。解析 10%、嵌入 80%、索引 10% 只是估算;如果页面精确显示 73%,用户会误以为这是可靠承诺。更诚实的是展示阶段和已完成数量,例如“正在生成向量,320/500 个片段”。
任务状态应包含当前阶段、尝试次数、开始时间、心跳、公开错误和可否重试。Worker 定期更新心跳,调度器才能识别进程崩溃留下的永久 processing 状态。
重试需要分类
网络超时可以指数退避重试,文件损坏重试十次也不会变好,权限不足必须由用户修正。若所有异常都自动重试,会制造队列风暴和额外模型费用。
可以把错误分成 transient、permanent 和 user_action_required。暂时错误限制次数并增加抖动,永久错误立即失败,需要用户动作的错误返回明确说明。每次重试复用幂等键,避免重复创建向量和片段。
配置与密钥离开代码库
异步 Worker 是独立进程,它也需要数据库、Redis 和模型凭证。把 .env 提交到仓库或打进镜像会扩大泄露范围。仓库只保存变量名称与示例,真实值由部署环境或密钥服务注入。
配置启动时校验,但日志只输出变量是否存在和非敏感选项。密钥轮换后,新任务使用新值,正在执行任务应有明确行为,而不是把凭证复制进每条消息永久保存。
事务与发消息之间的空隙
数据库提交成功后发布队列失败,资源会永远停在 pending;先发消息后事务回滚,Worker 又找不到记录。经典解决方案是事务 outbox:业务记录与待发布事件在同一数据库事务写入,由独立发布器可靠投递并标记完成。
规模较小时也可以定时扫描超时 pending 记录补发,但要承认这是恢复机制并监控积压。不要假设两次网络操作天然原子。
向普通数据任务迁移
视频转码、月度报表、批量邮件和 CSV 导入都有相同模型。练习可设计一个任务状态表,模拟 Worker 在索引完成后崩溃,再次投递时验证不会出现重复片段;随后让 Redis 暂停一分钟,检查 pending 扫描能否恢复任务。
五次提交组成的异步底座
b7cad4c:建立数据集模块,长期任务有了领域对象。9693ef2:引入 Redis 扩展,提供队列与缓存基础。5012373:加入 Celery,耗时工作离开请求线程。dd1b324:实现创建、更新与分页,管理面开始围绕状态运转。86f4d42:移除被提交的环境文件,密钥治理进入工程边界。
异步系统的复杂度不在队列 API,而在请求消失之后仍能保持事实一致。状态、幂等、重试和配置共同决定它是否可靠。