代码语言

知识点思维导图

30 个知识节点

生产工程(05) - 异步任务与队列基础

读完后,你应能完成以下任务:

  • 绘制“生产工程(05) - 异步任务与队列基础 / 同步 vs 异步:差在哪”的关键对象与数据流,解释“关键转变:接口的职责从「干完活再回」变成「登记任务并秒回」。”,并用源码位置、日志或 Trace 标注证据。
  • 为“生产工程(05) - 异步任务与队列基础 / 任务要有状态机”设计正常与异常输入,验证“除了状态,还要暴露 progress(百分比)和 current_step(正在干啥),前端才能显示「正在向量化(66%)」这种有体感的进度,而不是干巴巴一个转圈。”,输出首个偏差位置与回归测试结果。
  • 实现“生产工程(05) - 异步任务与队列基础 / 失败要能重试,但要有上限”的最小代码或配置,检验“后台任务跑得久,中间某一步偶发失败很常见(向量化接口超时、网络抖动)。”,输出命令、结果与 Diff,并说明不适用边界。

一、与进阶篇的分工

本篇保留为异步任务基础:重点讲长任务、进度和重试。进阶 Agent 定时任务请读 71 和 72,它们会把自然语言任务抽取、确认、持久化调度、幂等执行和失败审计讲完整。

二、异步任务与队列基础的真实应用场景

上一篇做完文档解析,用户在你的知识库后台一次性传了 200 个文档。每个文档要解析、切块、调向量化接口——一个文档 3 秒,200 个就是 10 分钟。

如果你在 HTTP 接口里同步处理:用户点「上传」,浏览器转圈转 10 分钟,大概率中途超时报错(网关一般 30-60 秒就掐断),用户以为失败了又点一次,雪上加霜。

正确做法:接口收到请求,立刻返回一个 task_id(秒回),真正的处理丢到后台慢慢跑。用户拿着 task_id 轮询进度,看着进度条从 0% 涨到 100%。这就是异步任务。判断标准很简单:一个操作可能超过几秒,就别让它占着 HTTP 请求

三、同步 vs 异步:差在哪

同步(错误做法):
  用户 --上传--> 接口【解析+切块+向量化,10分钟】--> 返回结果
                     用户全程干等,大概率超时

异步(正确做法):
  用户 --上传--> 接口【创建任务,秒回 task_id】
  用户 --轮询 task_id--> 接口【返回 progress: 45%】
  用户 --轮询 task_id--> 接口【返回 status: succeeded】
                     后台慢慢跑,用户随时能看进度

关键转变:接口的职责从「干完活再回」变成「登记任务并秒回」。活在后台异步执行。

四、任务要有状态机

异步任务跑在后台,用户看不见,全靠状态沟通。一个任务至少有这几种状态:

状态 含义 前端表现
pending 已提交,排队中 「已提交,等待处理」
running 执行中 进度条 + 当前步骤
succeeded 成功 「完成」,展示结果
failed 重试耗尽仍失败 错误原因 + 重试入口

除了状态,还要暴露 progress(百分比)和 current_step(正在干啥),前端才能显示「正在向量化(66%)」这种有体感的进度,而不是干巴巴一个转圈。

五、失败要能重试,但要有上限

后台任务跑得久,中间某一步偶发失败很常见(向量化接口超时、网络抖动)。不能一失败就整个任务作废,要能自动重试——但必须有上限,否则一个永远会失败的任务会无限重试,把资源吃光:

真实项目里重试之间还会加退避(第一次等 1 秒、第二次等 2 秒、第四次等 4 秒),避免下游服务正忙时被重试流量压垮。

六、工程上真正会踩的坑

  • 任务状态存在内存里。服务一重启,所有进行中的任务状态全丢,用户的 task_id 查不到了。状态要存 Redis / 数据库这种独立存储。
  • 重试没有上限或没有退避。永远失败的任务无限重试烧资源;下游正挂着时密集重试把它彻底压死。要设上限 + 指数退避。
  • 轮询太频繁。前端每 100ms 查一次,几千个任务一起轮询能把接口打爆。轮询间隔要合理(1-2 秒),或者改用 SSE/WebSocket 推送。
  • 失败了不告诉用户原因。任务 failed 但前端只显示「失败」,用户不知道是文件格式不对还是服务故障。error 字段要带可读原因,并给重试入口。
  • 用户取消了任务还在后台跑。用户关页面/点取消,后台任务要能感知并停止,否则白白消耗算力和钱。

七、一句话面试答法

什么任务要异步化,怎么设计? 判断标准是可能超过几秒的操作就别占 HTTP 请求,比如文档批量入库、长文分析。设计上接口收到请求立刻创建任务、秒回 task_id,真正的活丢后台异步跑。任务有状态机 pending/running/succeeded/failed,还要暴露 progress 和 current_step 让前端显示进度。失败自动重试但要有上限加指数退避,避免无限重试烧资源。任务状态存 Redis/数据库而不是内存,防止重启丢失。前端轮询或用 SSE 拿进度,失败要给可读原因和重试入口。

八、动手实践:18 异步任务与队列基础

用 asyncio 模拟「文档入库」这种慢任务,演示异步任务的完整生命周期:提交秒回 task_id、后台跑、状态流转、前端轮询进度、失败自动重试。

8.1 在线运行

零依赖,纯标准库(asyncio)。

8.2 预期输出

=== 提交文档入库任务(异步,不阻塞)===
接口立即返回 task_id:task_001,状态:pending

  [task_001] 第1次尝试 - 执行:解析文档
  >> 前端轮询看到:{'id': 'task_001', 'status': 'running', 'progress': 0, 'current_step': '解析文档', 'attempts': 1}
  [task_001] 解析文档 完成,进度 33%
  [task_001] 第1次尝试 - 执行:切分 chunk
  >> 前端轮询看到:{'id': 'task_001', 'status': 'running', 'progress': 33, 'current_step': '切分 chunk', 'attempts': 1}
  [task_001] 失败:步骤「切分 chunk」临时失败(如向量化接口超时) -> 准备第2次重试

  [task_001] 第2次尝试 - 执行:解析文档
  >> 前端轮询看到:{'id': 'task_001', 'status': 'running', 'progress': 0, 'current_step': '解析文档', 'attempts': 2}
  [task_001] 解析文档 完成,进度 33%
  [task_001] 第2次尝试 - 执行:切分 chunk
  >> 前端轮询看到:{'id': 'task_001', 'status': 'running', 'progress': 33, 'current_step': '切分 chunk', 'attempts': 2}
  >> 前端轮询看到:{'id': 'task_001', 'status': 'running', 'progress': 33, 'current_step': '切分 chunk', 'attempts': 2}
  [task_001] 切分 chunk 完成,进度 66%
  [task_001] 第2次尝试 - 执行:向量化
  >> 前端轮询看到:{'id': 'task_001', 'status': 'running', 'progress': 66, 'current_step': '向量化', 'attempts': 2}
  [task_001] 向量化 完成,进度 100%
  [task_001] 任务成功 ✓

  >> 前端轮询看到:{'id': 'task_001', 'status': 'succeeded', 'progress': 100, 'current_step': 'done', 'attempts': 2}
=== 最终结果:{'id': 'task_001', 'status': 'succeeded', 'progress': 100, 'current_step': 'done', 'attempts': 2} ===
要点:提交秒回 task_id;中间步骤失败自动重试;前端靠轮询拿进度。

「切分 chunk」步骤第一次故意失败,触发整个任务重试(attempts: 1 → 2),第二次跑通。前端轮询全程能看到 status、progress、current_step 的变化。

注:步骤耗时和轮询是并发的,每次运行轮询打印的条数可能略有不同(多一两行少一两行都正常),但状态流转 pending→running→succeeded、attempts 从 1 到 2 是稳定的。

8.3 代码↔概念对应

概念 在 main.py 哪里
任务状态机(pending/running/succeeded/failed) 顶部常量 + Task.status
提交秒回 task_id main 里创建 Task 后立即打印 id
任务执行 + 进度更新 execute 里的 task.progress
失败整体重试 + 重试上限 executewhile task.attempts <= max_retries
前端轮询拿进度 poll + Task.snapshot
执行与轮询并发 mainasyncio.gather

8.4 动手改

  • max_retries 改成 0,看任务在「切分 chunk」失败后直接变 failed,不再重试。
  • run_step 加个永远失败的步骤,观察重试耗尽后 statusfailederror 记录原因。
  • 真实项目里 Task 存进 Redis / 数据库而非内存,poll 换成前端定时 GET /api/tasks/{id},多个任务用真正的队列(Celery / RQ)调度。

8.5 可运行源码:异步任务与队列基础

main.py

九、总结

  • 与进阶篇的分工:进阶 Agent 定时任务请读 71 和 72,它们会把自然语言任务抽取、确认、持久化调度、幂等执行和失败审计讲完整。
  • 任务要有状态机:除了状态,还要暴露 progress(百分比)和 current_step(正在干啥),前端才能显示「正在向量化(66%)」这种有体感的进度,而不是干巴巴一个转圈。
  • 失败要能重试,但要有上限:后台任务跑得久,中间某一步偶发失败很常见(向量化接口超时、网络抖动)。
  • 工程上真正会踩的坑:永远失败的任务无限重试烧资源;
  • 一句话面试答法:判断标准是可能超过几秒的操作就别占 HTTP 请求,比如文档批量入库、长文分析。

9.1 可运行实验:异步建库队列与死信处理

<!doctype html>
<html lang="zh-CN">
<head>
  <meta charset="utf-8">
  <meta name="viewport" content="width=device-width,initial-scale=1">
  <title>AA-13 在线实验</title>
  <style>
    :root{color-scheme:dark;font-family:Inter,system-ui,sans-serif}*{box-sizing:border-box}body{margin:0;background:#0f1211;color:#e7ece9;font-size:13px}.shell{padding:16px}.top{display:flex;justify-content:space-between;gap:16px;margin-bottom:14px}h1{margin:3px 0;font-size:18px}.id,.value{color:#68e0b5;font-family:ui-monospace,monospace}.summary{margin:4px 0;color:#a5afa9}.run{border:0;border-radius:6px;background:#68e0b5;color:#07110d;padding:8px 14px;font-weight:700}.grid{display:grid;grid-template-columns:minmax(220px,.8fr) minmax(0,1.8fr);gap:12px}.panel{border:1px solid #29322e;background:#141817;padding:12px}.control{display:grid;gap:5px;margin-bottom:11px}.head{display:flex;justify-content:space-between;gap:8px}select,input{width:100%;accent-color:#68e0b5;background:#0d100f;color:#e7ece9}.toggle{display:flex;justify-content:space-between;border-top:1px solid #29322e;padding-top:9px}.toggle input{width:18px}.metrics{display:grid;grid-template-columns:repeat(4,minmax(0,1fr));gap:7px}.metric{border:1px solid #29322e;padding:8px}.metric b{display:block;color:#68e0b5;font-size:16px}.stages{display:flex;gap:6px;overflow:auto;margin:10px 0}.stage{border:1px solid #8a6230;padding:7px;min-width:90px}.stage.ok{border-color:#367a61}.stage.fail{border-color:#8b4545}table{width:100%;border-collapse:collapse}td{border-top:1px solid #29322e;padding:7px}.diagnosis{margin-top:9px;border-left:3px solid #68e0b5;background:#101412;padding:9px;line-height:1.5}.danger{border-color:#ef7f7f}@media(max-width:680px){.top,.grid{display:grid;grid-template-columns:1fr}.metrics{grid-template-columns:repeat(2,1fr)}}
  </style>
</head>
<body>
  <main class="shell">
    <header class="top"><div><div class="id">AA-13 · DETERMINISTIC LAB</div><h1 id="title"></h1><p class="summary" id="summary"></p></div><button class="run" id="run">运行实验</button></header>
    <section class="grid"><div class="panel"><div id="controls"></div><label class="toggle"><span>注入典型故障</span><input id="failure" type="checkbox"></label></div><div class="panel"><div class="metrics" id="metrics"></div><div class="stages" id="stages"></div><table><tbody id="rows"></tbody></table><div class="diagnosis" id="diagnosis"></div></div></section>
  </main>
  <script>
    const scenario = { title: '异步建库队列与死信处理', summary: '观察上传任务从 Pending、Parsing、Embedding、Indexing 到完成或死信。', controls: [
    { key: 'jobs', label: '待建库任务', type: 'range', min: 10, max: 500, step: 10, value: 120, suffix: ' 个' },
    { key: 'workers', label: 'Worker 数量', type: 'range', min: 1, max: 16, value: 4, suffix: ' 个' },
    { key: 'retries', label: '最大重试', type: 'range', min: 0, max: 5, value: 3, suffix: ' 次' }
  ] };
    const controls = document.querySelector('#controls');
    const failure = document.querySelector('#failure');
    document.querySelector('#title').textContent = scenario.title;
    document.querySelector('#summary').textContent = scenario.summary;
    function renderControl(control) {
      const label = document.createElement('label'); label.className = 'control';
      const head = document.createElement('span'); head.className = 'head'; head.innerHTML = '<span>' + control.label + '</span><span class="value" data-value="' + control.key + '"></span>'; label.appendChild(head);
      const input = document.createElement(control.type === 'select' ? 'select' : 'input'); input.dataset.key = control.key;
      if (control.type === 'select') control.options.forEach(option => { const item = document.createElement('option'); item.value = option[0]; item.textContent = option[1]; item.selected = option[0] === control.value; input.appendChild(item); });
      else { input.type = 'range'; input.min = control.min; input.max = control.max; input.step = control.step || 1; input.value = control.value; }
      input.addEventListener('input', updateValues); label.appendChild(input); return label;
    }
    function updateValues() { scenario.controls.forEach(control => { const input = controls.querySelector('[data-key="' + control.key + '"]'); document.querySelector('[data-value="' + control.key + '"]').textContent = control.type === 'select' ? input.options[input.selectedIndex].text : input.value + (control.suffix || ''); }); }
    function readValues() { const values = {}; scenario.controls.forEach(control => { const input = controls.querySelector('[data-key="' + control.key + '"]'); values[control.key] = control.type === 'range' ? Number(input.value) : input.value; }); values.failure = failure.checked; return values; }
    function stage(name, state, detail) { return { name, state, detail }; }
    const aiStage = stage;
    function clamp(value, minimum, maximum) { return Math.min(maximum, Math.max(minimum, value)); }
    function simulate(values) { const fail = values.failure;
      /** 每分钟单个 Worker 的确定性处理能力。 */
      const ratePerWorker = 6;
      /** 当前队列预计清空所需分钟。 */
      const minutes = Math.ceil(values.jobs / (values.workers * ratePerWorker));
      /** 故障任务经过重试后进入死信的数量。 */
      const deadLetters = fail ? Math.max(1, Math.floor(values.jobs * 0.03)) : 0;
      /** 重试产生的额外任务执行次数。 */
      const attempts = values.jobs + deadLetters * values.retries;
      return { metrics: [[minutes + 'm', '预计清空'], [attempts, '总执行次数'], [deadLetters, '死信任务'], [values.workers * ratePerWorker + '/m', '消费速率']], stages: [aiStage('Pending', 'ok', values.jobs), aiStage('Parsing', fail ? 'warn' : 'ok', values.workers), aiStage('Embedding', 'ok', 'batch'), aiStage('Indexing', 'ok', 'idempotent'), aiStage('Retry', deadLetters ? 'warn' : 'ok', values.retries), aiStage('Dead Letter', deadLetters ? 'fail' : 'ok', deadLetters), aiStage('Done', deadLetters ? 'warn' : 'ok', values.jobs - deadLetters)], rows: [['幂等键', 'tenant_id + document_id + checksum + pipeline_version'], ['重试策略', values.retries > 3 ? '重试过多会阻塞正常任务,应指数退避并进入死信' : '有限重试后转死信'], ['可观测性', deadLetters ? '记录阶段、异常类型、源文件和最后一次 Trace ID' : '全部任务完成']], diagnosis: deadLetters ? '部分任务进入死信,在线索引不能把它们标记为可查询。' : '队列吞吐、幂等和状态迁移正常。', danger: deadLetters > 0 };
     }
    function render() { const result = simulate(readValues()); document.querySelector('#metrics').innerHTML = result.metrics.map(item => '<div class="metric"><b>' + item[0] + '</b><span>' + item[1] + '</span></div>').join(''); document.querySelector('#stages').innerHTML = result.stages.map(item => '<div class="stage ' + item.state + '"><b>' + item.name + '</b><div>' + item.detail + '</div></div>').join(''); document.querySelector('#rows').innerHTML = result.rows.map(item => '<tr><td>' + item[0] + '</td><td>' + item[1] + '</td></tr>').join(''); const diagnosis = document.querySelector('#diagnosis'); diagnosis.textContent = result.diagnosis; diagnosis.className = 'diagnosis' + (result.danger ? ' danger' : ''); }
    scenario.controls.forEach(control => controls.appendChild(renderControl(control))); updateValues(); document.querySelector('#run').addEventListener('click', render); render();
  </script>
</body>
</html>

学完自测

选择所有正确答案;提交后逐项核对判断依据。

1在“异步任务与队列基础”中,需要同时满足“与进阶篇的分工”与“异步任务与队列基础的真实应用场景”。给定正文约束“进阶 Agent 定时任务请读 71 和 72,它们会把自然语言任务抽取、确认、持久化调度、幂等执行和失败审计讲完整。”,哪些判断保持了原有处理机制?多选
2“异步任务与队列基础”出现偏差:“在“异步任务与队列基础 / 同步 vs 异步:差在哪”中,即使不满足“接口的职责从「干完活再回」变成「登记任务并秒回」”,结果与副作用仍会保持不变。”已成为实际行为。围绕“同步 vs 异步:差在哪”与“任务要有状态机”,哪些判断能定位被改变的职责或边界?多选
3评审“异步任务与队列基础”方案时,验收条件包含“不能一失败就整个任务作废,要能自动重试——但必须有上限,否则一个永远会失败的任务会无限重试,把资源吃光。”。关于“失败要能重试,但要有上限”与“一句话面试答法”的哪些决策符合正文机制?多选