知识点思维导图
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 |
| 失败整体重试 + 重试上限 | execute 的 while task.attempts <= max_retries |
| 前端轮询拿进度 | poll + Task.snapshot |
| 执行与轮询并发 | main 里 asyncio.gather |
8.4 动手改
- 把
max_retries改成 0,看任务在「切分 chunk」失败后直接变failed,不再重试。 - 给
run_step加个永远失败的步骤,观察重试耗尽后status变failed、error记录原因。 - 真实项目里
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>
学完自测
选择所有正确答案;提交后逐项核对判断依据。