提出用数据库表 + 状态机实现后台任务队列的最小实践,列出了需要从请求Handler迁移到队列的四个判断条件。
当一个请求周期不再够用
阈值不是一段时间,而是一组属性。当满足以下任一条件时,就该迁移到任务队列:
工作时长超过浏览器可等待的时间。负载均衡器、CDN 和平台网关通常在 30 秒、60 秒或 100 秒处切断连接,无论你的服务器在做什么。
工作必须能跨越部署存活。请求处理器中的十分钟批处理任务在每次发布时都会随进程一起终止。
结果值得保留。如果用户合理地期望关闭标签页后再回来,结果就需要有自己的地址。
工作必须在用户间做限流。队列是全局并发限制真正可以存在的地方;请求处理器无法知道其他请求在做什么。
如果答案在三十秒以内且用户正在观看,那就用流式输出——一个带流式输出的 FastAPI 端点比在任务轮询上挂个加载动画体验好得多。
状态机,展开来说
五个状态。在代码之前先把它们写下来,这能避免之后出现"是卡住了还是还在跑"这种对话。
queued ──claim──▶ running ──success──▶ succeeded (terminal)
▲ │
│ ├──retryable failure, attempts < max──┐
└─────────────────┘ │
│ │
├──permanent failure──▶ failed ───────┤ (terminal)
│ │
└──worker died, lease expired──────────┘
│
back to queued ◀──────┘
Invariants:
* exactly one worker may hold a job in RUNNING
* every non-terminal job has a lease_expires_at in the future
* a job in RUNNING past its lease is available to be reclaimed
* FAILED is never reached without a stored error message
租约(lease)是大多数手写队列遗漏的部分,缺失它正是任务在 worker 被 OOM 杀死后永远处于"running"状态的原因。Worker 不拥有任务;它持有的是一个需要不断续期的有时限声明。
这张图的另一个后果是任务可能会运行两次——丢失租约的 worker 可能在另一个 worker 认领同一行时仍在工作。因此工作必须能够安全地重复,这与任何重试操作对幂等性的要求相同,这也是一个同时要发邮件的任务需要用 job id 作为发送键的原因。AI 特性的后台任务从产品角度涵盖了相同的结构,部分失败涵盖了如何在任务进行到一半时向用户展示状态。
SQLite 已经足够支撑单写入者场景的生产环境,schema 可以不加修改地迁移到 Postgres,类型除外。
-- schema.sql
CREATE TABLE IF NOT EXISTS jobs (
id TEXT PRIMARY KEY,
kind TEXT NOT NULL,
payload TEXT NOT NULL, -- JSON
state TEXT NOT NULL DEFAULT 'queued',
attempts INTEGER NOT NULL DEFAULT 0,
max_attempts INTEGER NOT NULL DEFAULT 3,
progress_done INTEGER NOT NULL DEFAULT 0,
progress_total INTEGER NOT NULL DEFAULT 0,
result TEXT, -- JSON, set on success
error TEXT, -- set on failure
lease_expires_at REAL, -- unix seconds
created_at REAL NOT NULL,
updated_at REAL NOT NULL
);
CREATE INDEX IF NOT EXISTS jobs_claimable
ON jobs (state, lease_expires_at);
# jobs.py
import json
import sqlite3
import time
import uuid
DB = "jobs.db"
LEASE_SECONDS = 120.0
def connect() -> sqlite3.Connection:
conn = sqlite3.connect(DB, isolation_level=None) # autocommit; explicit BEGIN
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA busy_timeout=5000")
return conn
def enqueue(conn: sqlite3.Connection, kind: str, payload: dict,
total: int = 0) -> str:
job_id = uuid.uuid4().hex
now = time.time()
conn.execute(
"INSERT INTO jobs (id, kind, payload, progress_total, created_at, updated_at)"
" VALUES (?, ?, ?, ?, ?, ?)",
(job_id, kind, json.dumps(payload), total, now, now),
)
return job_id
def claim(conn: sqlite3.Connection) -> sqlite3.Row | None:
"""Atomically take one claimable job. Returns None if there is nothing to do."""
now = time.time()
conn.execute("BEGIN IMMEDIATE")
try:
row = conn.execute(
"SELECT * FROM jobs"
" WHERE (state = 'queued')"
" OR (state = 'running' AND lease_expires_at < ?)"
" ORDER BY created_at LIMIT 1",
(now,),
).fetchone()
if row is None:
conn.execute("COMMIT")
return None
conn.execute(
"UPDATE jobs SET state='running', attempts = attempts + 1,"
" lease_expires_at = ?, updated_at = ? WHERE id = ?",
(now + LEASE_SECONDS, now, row["id"]),
)
conn.execute("COMMIT")
except Exception:
conn.execute("ROLLBACK")
raise
return conn.execute("SELECT * FROM jobs WHERE id = ?", (row["id"],)).fetchone()
BEGIN IMMEDIATE 是使认领具有原子性的关键:它在 SELECT 之前获取写锁,这样两个 worker 就不会同时读到同一个 queued 行并同时认领它。普通的 BEGIN 会延迟锁获取,正好允许这种竞态。在 Postgres 上,等效的操作是 SELECT ... FOR UPDATE SKIP LOCKED。
# worker.py
import json
import time
import traceback
from jobs import connect, claim, LEASE_SECONDS
POLL_SECONDS = 1.0
class PermanentError(Exception):
"""Do not retry: bad input, a 400, an unsupported model."""
def renew(conn, job_id: str, done: int) -> None:
conn.execute(
"UPDATE jobs SET lease_expires_at = ?, progress_done = ?, updated_at = ?"
" WHERE id = ?",
(time.time() + LEASE_SECONDS, done, time.time(), job_id),
)
def run_job(conn, job) -> dict:
payload = json.loads(job["payload"])
items = payload["items"]
results = []
for index, item in enumerate(items, start=1):
results.append(process_one(item)) # your model call
if index % 5 == 0:
renew(conn, job["id"], index) # heartbeat AND progress
return {"results": results}
def main() -> None:
conn = connect()
while True:
job = claim(conn)
if job is None:
time.sleep(POLL_SECONDS)
continue
try:
result = run_job(conn, job)
except PermanentError as exc:
conn.execute(
"UPDATE jobs SET state='failed', error=?, updated_at=? WHERE id=?",
(f"permanent: {exc}", time.time(), job["id"]),
)
except Exception:
detail = traceback.format_exc(limit=5)
if job["attempts"] >= job["max_attempts"]:
conn.execute(
"UPDATE jobs SET state='failed', error=?, updated_at=? WHERE id=?",
(detail, time.time(), job["id"]),
)
else:
conn.execute(
"UPDATE jobs SET state='queued', error=?, lease_expires_at=NULL,"
" updated_at=? WHERE id=?",
(detail, time.time(), job["id"]),
)
else:
conn.execute(
"UPDATE jobs SET state='succeeded', result=?, error=NULL,"
" progress_done=progress_total, updated_at=? WHERE id=?",
(json.dumps(result), time.time(), job["id"]),
)
if __name__ == "__main__":
main()
心跳和进度更新故意合并成同一条语句。两个独立的机制会漂移——你会看到一个任务显示完成了 90% 但已经死了十分钟——而一条语句就不会有这个问题。
对于处理多个条目的任务,要让 process_one 做到幂等并记录哪些条目已完成,这样被回收的任务就能从断点恢复而不是从头开始。这让每次崩溃的代价减半,和大模型分类任务做检查点保存是同样的纪律。
客户端可轮询的进度
@app.post("/jobs")
def create_job(body: JobRequest) -> dict:
conn = connect()
job_id = enqueue(conn, "classify", {"items": body.items}, total=len(body.items))
return {"job_id": job_id, "state": "queued"}
@app.get("/jobs/{job_id}")
def get_job(job_id: str) -> dict:
conn = connect()
row = conn.execute("SELECT * FROM jobs WHERE id = ?", (job_id,)).fetchone()
if row is None:
raise HTTPException(status_code=404, detail="no such job")
return {
"job_id": row["id"],
"state": row["state"],
"done": row["progress_done"],
"total": row["progress_total"],
"result": json.loads(row["result"]) if row["result"] else None,
"error": row["error"] if row["state"] == "failed" else None,
"stale": (
row["state"] == "running"
and (row["lease_expires_at"] or 0) < time.time()
),
}
stale 字段只增加一次比较,却回答了用户真正在问的问题。一个显示 40% 的进度条什么信息都没有;stale: true 意味着持有这个任务的 worker 不再续租,任务即将被回收。
用退避策略轮询——前十秒每秒一次,之后每五秒一次——而不是固定的一秒间隔,否则一百个空闲的浏览器标签就会变成每秒一百次对数据库的请求。
一条被读取的死信路径
每个队列都有一个失败任务去往的地方。大多数情况下没有人看它,这就等于没有。有四件事能带来改变:
存储错误,而不是一个标志位。上面的 error 列存的是截断后的栈追踪。布尔类型的 failed 列意味着了解发生了什么事的唯一办法是复现它。
用一条查询搞定。SELECT kind, substr(error, 1, 80), count(*) FROM jobs WHERE state = 'failed' GROUP BY 1, 2 ORDER BY 3 DESC 一秒钟就能告诉你是一个 bug 还是四十个 bug。把这个放到仓库里的脚本中,而不是某人的 shell 历史里。
对速率告警,而不是对数量。失败的总数只会增长,所以没人去看。值得告警的数字是过去一小时的失败数,或者失败数与成功数的比值。
让回放成为一个命令。把一组 id 的 state 改回 queued、attempts 改回零,这就是整个回放机制。它必须在故障发生前就存在,因为故障期间没人会写它。
SQLite、Postgres 还是 Redis
上面的状态机不关心自己存在哪里。随存储变化的是你能运行多少个 worker 以及一次崩溃的代价是什么,这三个答案确实是不同的,而不是一个需要攀登的阶梯。
对于这个工作负载,有两个属性比吞吐量更重要。AI 任务以分钟计,所以一个每秒处理十个认领的队列已经是压倒性的能力了——瓶颈从来不是队列。结果值得保留,这是在说你需要一个六周后仍能查询的存储,而不是一个为投递优化的存储。
无论选哪个,如果模型输出很大,都不要把它放到任务行里。每个任务的结果列如果持有 40 KB 的文本,就会把 jobs 表变成文档存储并使每次状态轮询都要读它。把 payload 写到对象存储或单独的表里,只保留一个引用。
Celery 或 RQ 中的同样东西
一旦状态机清晰了,采用框架就是映射工作而不是重新设计。这两个都需要一个 broker——Redis 或 RabbitMQ——这是迁移的真正成本。
Celery 的默认配置是为短任务调优的。worker_prefetch_multiplier 大于 1 意味着一个 worker 一次预留多个任务,这对分钟级的 AI 任务来说是错的——一个 worker 坐拥四个任务,而另外三个空闲。对于这个工作负载,把它设为 1。Celery 4 和 5 之间的配置名称不同;根据你锁定的版本确认。
For further actions, you may consider blocking this person and/or reporting abuse