🏠 总目录📚 本教程 11 · 长任务与队列
📑 本页目录(点开跳转)

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 没权限、内容被安全策略拒绝,重试一万次也是同样结果,直接进死信;只有 429500、超时才值得退避重试。


🔁 换个栈怎么对应

概念 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 ⚠️ 同名不同物:那边的「批处理」是推理引擎内部把多个请求拼进一次前向,毫秒级;本章的「批量任务」是应用层的分钟级排队

✅ 检查点

  1. 给出三条「必须异步」的判断信号。
  2. BackgroundTasks 的三个限制是什么?
  3. 为什么「先 SELECT 一个 queued 的、再 UPDATE 成 running」是错的?为什么「改写成一条语句」还不够?Postgres 上正确的写法长什么样?
  4. jobs 表里 run_after 这一列干什么用?
  5. 幂等键为什么不能用 uuid4() 生成?那段示例投递三次,账上有几笔、多少钱?
  6. 什么是毒药消息?防住它要一起做的三件事是哪三件?
  7. 退避为什么必须加抖动?
  8. 哪几类错误不该重试?
👀 答案
  1. 任意三条:可能超过 30 秒(网关硬超时,用户看到 504 但后端还在烧钱)、批量处理入库 + embedding定时任务需要重试的活
  2. ① 跑在 web 进程里,和请求抢 worker;② 进程重启全丢,而部署 / 扩缩容 / OOM 每天都在重启;③ 状态存进程内存,多进程部署下另一个进程查不到这个 id
  3. 两条语句之间有一道缝,两个 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 示例之所以安全,靠的是整库一把写锁,理由完全不同。
  4. 延迟执行的时间点。认领语句里有 run_after <= now,重试时把退避算出的秒数写进去,这个 job 自然被跳过到那个时刻之后。
  5. 幂等键必须由输入决定uuid4() 每次重试算出新键,唯一约束永远撞不上。更稳的是让客户端传 Idempotency-Key 头。示例结果是 ['charged','skipped','skipped'],账上 1 笔、30 分
  6. 永远会失败的任务,没有上限就无限循环、堵死队列。三件事:重试上限(示例 MAX_ATTEMPTS = 5)、指数退避加抖动死信(要有人看、能重新入队)。
  7. 一次上游抖动会让一批任务同时失败,没抖动它们会在同一秒一起回来,把刚恢复的上游再打挂。代码里是 base * (0.5 + random()),跑出 2.69 / 5.03 / 7.36 / 12.14 秒。
  8. 400401、内容被安全策略拒绝 —— 重试一万次结果一样,直接进死信。只有 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

打卡记录保存在你的浏览器里,首页能看到总进度