📑 本页目录(点开跳转)
11 · 长任务与队列
⏱ 52 分钟 | ⭐ 别一上来就上 Celery —— 中间那一档「数据库当队列」被严重低估
🎯 一句话
只要一个请求可能超过几十秒,它就不该在请求里做完;但「扔到后台」有三档做法,绝大多数人直接跳到了最重的那一档。 真正难的也不是把活扔出去,是扔出去之后怎么保证它只被干一次、失败了不会把队列堵死。
🧩 一、什么时候必须异步
别凭感觉。满足任意一条,就不要在请求里做完:
| 信号 | 为什么 |
|---|---|
| ⭐ 可能超过 30 秒 | 很多网关的默认超时就在这个量级(🗓️ 查你那层网关的配置)。超时后用户看到 504,后端还在跑,钱照花,结果没人收 |
| 批量处理 | 「把这 800 篇文档都总结一遍」,用户不可能挂着页面等 |
| ⭐ 入库 + embedding | 一个 PDF 切成几百块、每块调一次 embedding,天然是分钟级 |
| 定时任务 | 凌晨重建索引、跑重训练流水线 —— 压根没有「请求」这个东西 |
| 要重试的活 | 请求活不过重试。第三方挂了要 10 分钟后再试,只有后台能做 |
⚠️ 单次对话生成不属于这一类。它慢,但用户就坐在那儿等,流式已经解决了体感问题 —— 别拿队列去套。
🧩 二、三档做法
| 档 | 是什么 | 适用边界 | ⭐ 该升级的信号 |
|---|---|---|---|
① BackgroundTasks |
响应发出去后,在同一进程里接着干 | 几秒到几十秒、丢了也无所谓的活:发通知、写日志、清临时文件 | 出现「丢了用户会投诉」的活,或需要重试 / 看进度 |
| ⭐ ② 数据库当队列 | 一张 jobs 表 + 一条原子认领语句 |
中小规模几乎全覆盖:分钟级任务、每分钟几十到几百个 job、少数几个 worker | 积压到轮询压力明显 / 要扇出到几十个 worker / 要复杂的优先级调度 |
| ③ 专用队列 | Redis / RabbitMQ + worker 进程 | 高吞吐、多任务类型、要现成的路由和监控面板 | —— |
⭐ 第 ② 档被严重低估:你已经有一个数据库了,把队列也放进去 —— 少一个中间件、少一套部署、少一处会不一致的状态。而且任务状态能和业务数据在同一个事务里改,这一点 Redis 队列给不了。
① BackgroundTasks
形状是:先把 job 记下来、立刻返回一个 id,真正的活挂在响应之后跑。FastAPI 里这东西叫 BackgroundTasks。
import uuid, time
from fastapi import FastAPI, BackgroundTasks
from fastapi.testclient import TestClient
app = FastAPI()
JOBS = {} # ⚠️ 状态只在进程内存里:进程一重启,任务和进度一起蒸发
def slow_work(jid, text):
for i in range(1, 5):
time.sleep(0.02)
JOBS[jid].update(status="running", progress=i * 25) # ⭐ 进度写在前端查得到的地方
JOBS[jid].update(status="done", result=text.upper())
@app.post("/jobs", status_code=202) # ⭐ 202 = 「收到了,还没做完」
def create(text: str, bg: BackgroundTasks):
jid = str(uuid.uuid4())
JOBS[jid] = {"status": "queued", "progress": 0}
bg.add_task(slow_work, jid, text) # 响应写完之后才执行
return {"job_id": jid}
@app.get("/jobs/{jid}") # 前端拿这个轮询
def status(jid: str):
return JOBS.get(jid, {"status": "unknown"})
c = TestClient(app)
jid = c.post("/jobs", params={"text": "hi"}).json()["job_id"]
print(c.get("/jobs/" + jid).json())
⚠️ 三个限制:① 跑在你的 web 进程里,CPU 密集的活会和请求抢 worker;② ⭐ 进程重启就全丢,而部署、扩缩容、OOM 每天都在重启;③ JOBS 这个字典在多进程部署下,另一个进程查不到这个 id。
⭐ ② 数据库当队列
把「谁在干这个 job」落到磁盘,上面三个毛病一起没了。
核心只有一件事:认领必须原子 ——「先 SELECT 一个 queued 的、再 UPDATE 成 running」中间那道缝,就是两个 worker 抢到同一行的原因。
⚠️⚠️ 但「写成一条语句」只是必要条件,不是充分条件。 下面这条在 Postgres 上就不安全:
-- 🗓️ 未实跑 —— 需要 Postgres。⚠️ 这是【错误示范】,别抄
UPDATE jobs SET status='running', worker=$1
WHERE id = (SELECT id FROM jobs WHERE status='queued' ORDER BY run_after LIMIT 1)
RETURNING id, kind, payload;
两个并发事务的子查询会选中同一行:第二个先卡在行锁上排队等(吞吐直接塌掉),等前一个提交之后,外层 WHERE 只比对 id、不再检查 status='queued' —— 于是这次认领要么把同一个 job 认领第二次,要么空转一次什么也没拿到。⚠️ 具体落到哪一种取决于隔离级别和 WHERE 被重新求值的细节,但两种都是 bug。
⭐ 正确写法是把锁的语义写进子查询:
-- 🗓️ 未实跑 —— 需要 Postgres
UPDATE jobs SET status='running', attempts=attempts+1, worker=$1
WHERE id IN (SELECT id FROM jobs
WHERE status='queued' AND run_after <= now()
ORDER BY run_after
LIMIT 1
FOR UPDATE SKIP LOCKED) -- ⭐ 当场锁住选中的行;被别人锁住的直接跳过,不排队等
RETURNING id, kind, payload;
FOR UPDATE 让子查询当场把选中的行锁住,SKIP LOCKED 让别的 worker 跳过这些行去拿下一条 —— 既不重复也不排队。⚠️ 事务边界也是这条语句的一部分:它要在一个事务里跑,提交之后这个 job 才算真归你;没走到提交就崩(进程被 kill),事务回滚、行锁释放,这一行自动回到可认领状态。
下面用 SQLite 演示这套形状里能当场跑的那一半。⚠️ SQLite 没有 FOR UPDATE SKIP LOCKED,它靠的是整库一把写锁(同一时刻只有一个写事务),所以下面这条语句在 SQLite 上确实是安全的 —— 但它安全的理由和 Postgres 完全不是一回事,别把这段 SQL 原样搬过去:
import sqlite3, json, time, uuid
db = sqlite3.connect(":memory:")
db.execute("""CREATE TABLE jobs(
id TEXT PRIMARY KEY, kind TEXT, payload TEXT,
status TEXT DEFAULT 'queued', -- queued / running / done / dead
attempts INTEGER DEFAULT 0,
run_after REAL DEFAULT 0, -- ⭐ 延迟重试全靠这一列
worker TEXT)""")
db.execute("CREATE INDEX idx_pick ON jobs(status, run_after)")
def enqueue(kind, payload, delay=0.0):
jid = str(uuid.uuid4())
db.execute("INSERT INTO jobs(id,kind,payload,run_after) VALUES(?,?,?,?)",
(jid, kind, json.dumps(payload), time.time() + delay))
db.commit()
return jid
def claim(worker):
# ⭐ 认领必须原子:查到和改状态之间留了缝,两个 worker 就会拿到同一行
# ⚠️ 但「一条语句」不等于安全 —— 这里安全是靠 SQLite 的整库写锁;
# Postgres 上同样的写法必须配 FOR UPDATE SKIP LOCKED(见上)
cur = db.execute("""UPDATE jobs SET status='running', attempts=attempts+1, worker=?
WHERE id = (SELECT id FROM jobs WHERE status='queued' AND run_after<=?
ORDER BY run_after LIMIT 1)
RETURNING id, kind, payload""", (worker, time.time()))
row = cur.fetchone()
db.commit()
return row
enqueue("embed", {"doc": "A"})
enqueue("embed", {"doc": "B"})
print([None if (j := claim(w)) is None else (j[0][:8], w) for w in ["w1", "w2", "w1"]])
输出 [('…','w1'), ('…','w2'), None] —— 两个 worker 各拿到不同的一行,第三次没活了。
⚠️ 要配三件事:(status, run_after) 上要有索引(否则每次认领全表扫);worker 空转别 100ms 一轮(1–2 秒足够,或用 Postgres 的 LISTEN/NOTIFY 让入队方叫醒它);必须有清扫语句把 status='running' 却超过 N 分钟没动静的行改回 queued —— worker 被 kill 时来不及改状态。
③ 专用队列
到了每秒几百个 job、十几种任务类型、需要优先级路由和现成监控面板时才值得引入。⚠️ 代价是真实的:多一个要运维的中间件、多一处可能和数据库不一致的状态、本地开发要多起一个服务。
🧩 三、状态怎么让前端知道
| 轮询 | SSE 推送 | |
|---|---|---|
| 后端 | ⭐ 无状态,哪台机器都能答 | ⚠️ 连接钉死在某台机器上,扩容和滚动更新都要处理 |
| 延迟 | 一个轮询间隔 | 近实时 |
| 选它 | ⭐ 默认。分钟级任务,2–5 秒轮一次完全够 | 秒级进度条,或任务本身就有流式输出 |
⚠️ 最常见的错:轮询间隔写成 200 毫秒 —— 一个跑 3 分钟的任务就发了 900 个「好了没」。用退避式轮询:前 10 秒 1 秒一次,之后拉到 5 秒。
⭐ 进度报阶段,别报百分比。LLM 任务时长不可预测,百分比一定卡在 90%,而卡住的进度条比没有进度条更让人焦虑。给可数的:「已处理 3 / 12 个分片」。⚠️ 进度要写在 jobs 表的一列里;写在 worker 内存里,worker 一崩前端就永远停在 40%。
🧩 四、幂等:任务一定会被执行两次
⚠️ 不是「可能」,是「一定」:客户端超时后重试、worker 干完活但改状态前崩了、清扫脚本把一个其实还活着的 job 放回队列。
⭐ 代价很具体:用户被扣两次钱、同一段文本在向量库里存两份(检索时占掉两个 top-k 名额)、同一封邮件发两遍。
做法只有一个形状:给这件事算一个由输入决定的键,写进带唯一约束的表,写不进去就说明做过了。
import sqlite3, hashlib
db = sqlite3.connect(":memory:")
db.execute("CREATE TABLE billing(user TEXT, cents INTEGER)")
db.execute("CREATE TABLE idem(key TEXT PRIMARY KEY)") # ⭐ 主键就是那把锁
def charge(user, kind, payload, cents):
# ⭐ 键必须由【输入】算出来。用随机数或时间戳,重试会算出新键 = 等于没做幂等
key = hashlib.sha256(f"{user}|{kind}|{payload}".encode()).hexdigest()[:16]
try:
with db: # ⭐ 写键和扣款在同一个事务,中间崩了不会一边有一边没有
db.execute("INSERT INTO idem VALUES(?)", (key,))
db.execute("INSERT INTO billing VALUES(?,?)", (user, cents))
return "charged"
except sqlite3.IntegrityError: # 键撞了 = 这件事已经做过
return "skipped"
print([charge("u1", "summarize", "doc-42", 30) for _ in range(3)])
print("扣了", db.execute("SELECT COUNT(*), SUM(cents) FROM billing").fetchone())
投三次,输出 ['charged', 'skipped', 'skipped'],账上只有 1 笔 30 分。
⚠️ 两个坑:① 键不能用 uuid4() —— 重试时算出的是新键,唯一约束永远撞不上,等于没做;更稳的是让客户端传 Idempotency-Key 头,重试时原样带回。② 写键和干活必须同一个事务:先干活后写键,中间崩了会重复干;先写键后干活,中间崩了这件事永远不会被做,而且没人知道。
🧩 五、重试、死信、毒药消息
⚠️ 毒药消息:一个永远会失败的任务(payload 少了字段、文档编码坏了)。没有上限,它会被无限取出、无限失败、无限放回 —— 堵死队列,顺便刷爆日志。三件事一起做:上限、退避、死信。
import random
MAX_ATTEMPTS = 5
random.seed(0)
def backoff(attempt):
base = min(60.0, 2.0 ** attempt) # 指数退避 1,2,4,8,16…,封顶 60 秒
return round(base * (0.5 + random.random()), 2) # ⭐ 抖动:不加,失败任务会同一秒一起回来
def on_fail(job, err):
if job["attempts"] >= MAX_ATTEMPTS: # ⚠️ 没有上限,毒药消息会永远循环
return "dead", err
return "retry", backoff(job["attempts"]) # 这个数写回 jobs.run_after
job = {"id": "j1", "attempts": 0}
while True:
job["attempts"] += 1
state, info = on_fail(job, "ValueError: payload 缺 doc_id")
print(job["attempts"], state, info)
if state == "dead":
break
seed(0) 固定,退避秒数依次是 2.69 / 5.03 / 7.36 / 12.14,第 5 次转成 dead。
⭐ 抖动不是可选项:一次上游抖动会让 200 个任务同时失败,没有抖动它们会在同一秒一起回来,把刚恢复的上游再打挂一次。
⭐ 死信是待办箱,不是垃圾桶 —— 要看得见(后台页面或告警),并且能改完 payload 重新入队。没人看的死信队列 = 静默丢数据。
⚠️ 还要分清能不能重试:400 参数错、401 没权限、内容被安全策略拒绝,重试一万次也是同样结果,直接进死信;只有 429、500、超时才值得退避重试。
🔁 换个栈怎么对应
| 概念 | Python / FastAPI | Node(Express / Hono) | Go |
|---|---|---|---|
| 请求内的后台活 | BackgroundTasks |
响应后继续 / setImmediate |
go func() |
| ⚠️ 这一档的通病 | 进程重启就丢 | 一样会丢 | 一样会丢 |
| 数据库当队列 | 自己写 + SKIP LOCKED |
pg-boss |
River |
| 专用队列 | Celery / RQ / arq | BullMQ | asynq |
| 定时任务 | APScheduler | node-cron | robfig/cron |
⭐ 名字全不一样,四件事完全一样:一张带状态的表、一条原子认领语句、一个由输入算出的幂等键、一套带抖动和上限的重试。这四件事不过期,上面那些包名会。
🔗 这一章连到哪里
| 去哪 | 为什么 |
|---|---|
| Agent 长时任务的骨架 | 那边讲一个 Agent 自己怎么跑很久,本章讲后端怎么托住这种活。做「跑半小时的 Agent 任务」两边都要看 |
| 05 · 流式输出 | ⭐ 另一条岔路:单次对话生成也慢,但用户就坐在那儿等 —— 那是流式解决的问题,能流式就不要上队列。本章管的是另一类:用户等不了、必须先走人的活 |
| ⭐⭐ 03c · 幂等与条件请求 | 本章第四节解决的是队列侧的幂等(同一个任务被消费两次)。⭐ 那一章解决接口侧的:本章答案里提过一句「更稳的是让客户端传 Idempotency-Key 头」——服务端拿到这个头之后该怎么办在那里(同键不同 body、并发同键、要不要重放原响应、存多久) |
| ⭐ 03b · 接口的形状 | 本章第三节那张「状态怎么让前端知道」的表只给了轮询和 SSE 两条路 —— ⚠️ 它的上游是一个建模问题:长任务该不该有自己的资源(202 Accepted + 一个 jobs 资源)。那一章讲这个 |
| 重训练 | 定时任务里最典型的一种:那边讲什么时候该重训,本章讲这条流水线怎么被调度和重试 |
| 连续批处理与 PagedAttention | ⚠️ 同名不同物:那边的「批处理」是推理引擎内部把多个请求拼进一次前向,毫秒级;本章的「批量任务」是应用层的分钟级排队 |
✅ 检查点
- 给出三条「必须异步」的判断信号。
BackgroundTasks的三个限制是什么?- 为什么「先
SELECT一个 queued 的、再UPDATE成 running」是错的?为什么「改写成一条语句」还不够?Postgres 上正确的写法长什么样? jobs表里run_after这一列干什么用?- 幂等键为什么不能用
uuid4()生成?那段示例投递三次,账上有几笔、多少钱? - 什么是毒药消息?防住它要一起做的三件事是哪三件?
- 退避为什么必须加抖动?
- 哪几类错误不该重试?
👀 答案
- 任意三条:可能超过 30 秒(网关硬超时,用户看到 504 但后端还在烧钱)、批量处理、入库 + embedding、定时任务、需要重试的活。
- ① 跑在 web 进程里,和请求抢 worker;② 进程重启全丢,而部署 / 扩缩容 / OOM 每天都在重启;③ 状态存进程内存,多进程部署下另一个进程查不到这个 id。
- 两条语句之间有一道缝,两个 worker 会认领同一行。⚠️ 但只把它并成一条语句还不够:Postgres 上
UPDATE ... WHERE id = (SELECT ... LIMIT 1)的外层WHERE只比对id、不再检查status,并发时第二个事务先卡在行锁上等,等到之后要么重复认领要么空转。正确写法是把锁写进子查询 ——UPDATE ... WHERE id IN (SELECT id ... ORDER BY run_after LIMIT 1 FOR UPDATE SKIP LOCKED) RETURNING ...:FOR UPDATE当场锁行,SKIP LOCKED让别的 worker 跳过去拿下一条。本章那段 SQLite 示例之所以安全,靠的是整库一把写锁,理由完全不同。 - 延迟执行的时间点。认领语句里有
run_after <= now,重试时把退避算出的秒数写进去,这个 job 自然被跳过到那个时刻之后。 - 幂等键必须由输入决定;
uuid4()每次重试算出新键,唯一约束永远撞不上。更稳的是让客户端传Idempotency-Key头。示例结果是['charged','skipped','skipped'],账上 1 笔、30 分。 - 永远会失败的任务,没有上限就无限循环、堵死队列。三件事:重试上限(示例
MAX_ATTEMPTS = 5)、指数退避加抖动、死信(要有人看、能重新入队)。 - 一次上游抖动会让一批任务同时失败,没抖动它们会在同一秒一起回来,把刚恢复的上游再打挂。代码里是
base * (0.5 + random()),跑出 2.69 / 5.03 / 7.36 / 12.14 秒。 400、401、内容被安全策略拒绝 —— 重试一万次结果一样,直接进死信。只有429/500/ 超时值得退避重试。
🛑 可以停在这里
⚡ 走神救援
这一章讲「一个请求做不完的活怎么扔到后台」。判断标准:可能超过 30 秒(网关默认硬超时就在这量级,超时后用户看到 504、后端还在烧钱)、批量处理、入库 + embedding、定时任务、需要重试的活 —— 任意一条就别在请求里做完。单次对话生成不算,那是流式解决的问题。 三档做法,关键是别直接跳到第三档。①
BackgroundTasks:响应发完后在同一进程里接着干,三个限制 —— 和请求抢 worker、进程重启全丢、状态放内存导致多进程部署查不到 id。② ⭐ 数据库当队列:一张jobs表 + 一条原子认领语句,中小规模几乎全覆盖,少一个中间件,而且任务状态能和业务数据在同一个事务里改。认领必须原子,「先 SELECT 再 UPDATE」中间那道缝就是两个 worker 抢到同一行的原因;⚠️⚠️ 但「写成一条语句」只是必要条件不是充分条件 —— Postgres 上UPDATE ... WHERE id = (SELECT ... LIMIT 1)的外层WHERE只比对id、不再检查status,并发时要么重复认领要么空转,必须把锁写进子查询:UPDATE ... WHERE id IN (SELECT id ... LIMIT 1 FOR UPDATE SKIP LOCKED) RETURNING ...(FOR UPDATE当场锁行、SKIP LOCKED让别人跳过去拿下一条);本章那段 SQLite 示例安全靠的是整库一把写锁,理由完全不同。配套:(status, run_after)上要有索引、空转别 100ms 一轮、要有清扫语句把卡在running的行放回去。③ 专用队列等到每秒几百个 job 再上。 状态怎么让前端知道:默认轮询(无状态、哪台机器都能答),SSE 只在要秒级进度条时用。别把间隔写成 200 毫秒 —— 三分钟的任务会发出 900 个「好了没」。进度报阶段不报百分比(「3 / 12 个分片」),LLM 时长不可预测,百分比一定卡在 90%;进度写在jobs表里,写在 worker 内存里它一崩前端就永远停在 40%。 幂等不是可选项:任务一定会执行两次,代价是用户被扣两次钱、同一段文本在向量库存两份。做法是由输入算一个键、写进带唯一约束的表。两个坑:不能用uuid4()(重试算出新键等于没做)、写键和干活必须同一事务。示例投三次,账上只有 1 笔 30 分。 失败处理三件套:上限(示例MAX_ATTEMPTS = 5)、指数退避加抖动(min(60, 2**n) * (0.5+random),跑出 2.69 / 5.03 / 7.36 / 12.14 秒后第 5 次转dead)、死信。抖动不是可选项 —— 200 个同时失败的任务没抖动会在同一秒一起回来。死信是待办箱不是垃圾桶。还要分类:400/401/ 内容被拒直接进死信,只有429/500/ 超时值得重试。换栈时不变的就是这四件事:带状态的表、原子认领、幂等键、带抖动和上限的重试。
下一节 👉 12-限流配额与成本护栏.md