本节目标:判断哪些工作该异步化,理解任务队列的 broker / worker / 结果后端模型,掌握指数退避重试的写法,并能自己实现一个「不丢任务」的最小可靠队列。
适用版本:Python 3.12+(实测 3.14.6);tenacity 9.2.1、fakeredis 2.39.0
8.2 任务队列与重试
8.1 节我们用缓存挡下了读压力,但有些工作根本不该出现在请求线程里:发邮件、生成报表、调用一个可能超时的第三方接口、批量重算推荐结果。这些活的特点是慢、会失败、不需要立刻返回结果。把它们塞进请求里,接口响应时间就被它们拖死;丢进一个后台线程又管不好重启和重试。任务队列就是为这类工作准备的。
本节有两个目标:一是搞清楚成熟框架(Celery / arq)的模型,这样选型和排查问题时有坐标系;二是亲手实现一个最小可靠队列,因为只有自己写过「取任务、失败、重试、进死信」这条链路,才会真正理解框架替你做了什么。
8.2.1 判断该不该异步
不是所有耗时操作都值得异步化。判断标准有三条:能否容忍延迟、失败后能否重试、结果是否被同步需要。三条都满足,才适合丢进队列。
| 场景 | 是否适合异步 | 原因 |
|---|---|---|
| 注册后发欢迎邮件 | 适合 | 晚几秒无感,失败可重试,不需要同步结果 |
| 下单扣库存 | 不适合 | 必须同步返回结果,失败要立刻反馈用户 |
| 生成月度报表 | 适合 | 耗时分钟级,用户可稍后下载 |
| 实时搜索建议 | 不适合 | 要求毫秒级返回,延迟不可接受 |
| 调用外部支付回调 | 适合 | 网络抖动常见,需要重试与幂等 |
反面例子尤其重要:「会失败」不等于「该异步」。扣库存也会失败,但它必须在事务里同步完成。把强一致操作丢进队列,只会让问题从「报错」变成「状态不一致」。
8.2.2 Celery 与 arq 的模型对比
两个框架代表了两种设计取向。Celery 是同步的、功能全的「重装步兵」;arq 是异步的、极简的「轻骑兵」。
| 维度 | Celery | arq |
|---|---|---|
| 编程模型 | 同步函数(用进程/线程池执行) | async def 原生协程 |
| 依赖 | 需要 broker(Redis/RabbitMQ) | 只要 Redis |
| 定时任务 | celery beat 独立进程 | 内置 cron 表达式 |
| 任务编排 | chain / group / chord 等 canvas | 无,靠代码组合 |
| 结果后端 | 支持多种后端,可查任务状态 | 通过 job.result() 等待 |
| 适用 | 大团队、复杂编排、成熟生态 | asyncio 技术栈、轻量够用 |
选型上没有绝对优劣:技术栈是 asyncio、任务形态简单,选 arq;需要复杂编排或已有 Celery 生态,选 Celery。
Celery 的一个最小配置长这样(本机未安装 Celery,以下为示意,未实测):
from celery import Celery
app = Celery("tasks", broker="redis://localhost:6379/0",
backend="redis://localhost:6379/1")
app.conf.update(
task_acks_late=True, # 任务执行完才 ack,崩溃可重新投递
task_reject_on_worker_lost=True, # worker 被杀时拒绝任务,而非静默丢失
task_default_retry_delay=10,
broker_transport_options={"visibility_timeout": 3600},
)
@app.task(bind=True, max_retries=3,
autoretry_for=(ConnectionError,),
retry_backoff=True, retry_jitter=True)
def send_email(self, to):
...
task_acks_late 是 Celery 里最关键的一个开关:默认情况下任务一取出来就 ack(确认),worker 中途崩溃任务就永久丢了;开了它,ack 推迟到任务真正执行完成,崩溃的任务会被 broker 重新投递。代价是任务可能被执行两次——所以业务必须幂等(8.3 节展开)。
arq 的写法则是纯协程(同样未实测,仅示意):
from arq import cron
from arq.connections import RedisSettings
async def send_email(ctx, to):
...
async def startup(ctx):
ctx["pool"] = ... # 建连接池
class WorkerSettings:
functions = [send_email]
cron_jobs = [cron(send_email, hour=9, minute=0)]
redis_settings = RedisSettings(host="localhost", port=6379)
max_tries = 3
arq 把「任务」「定时」「重试次数」都收在一个 WorkerSettings 类里,没有 beat 进程,定时逻辑内嵌在 worker 内。
8.2.3 broker 与可见性超时
无论哪个框架,队列的可靠性都建立在一个机制上:可见性超时(visibility timeout)。任务被 worker 取走时,broker 不会立刻删除它,而是把它挪到一个「处理中」的临时位置并开始计时:
- worker 在超时前确认完成 → broker 彻底删除该任务。
- worker 崩溃或超时未确认 → broker 认为任务失败,把它重新投回队列。
这正是「不丢任务」的保证。代价是至少一次(at-least-once)投递:任务可能被执行两次。想清楚这一点,就理解了为什么幂等是异步任务的必修课,而不是可选项。
8.2.4 重试策略与指数退避
任务失败后立刻重试往往没用——下游还在抖,立刻重试只会再失败一次,还可能把下游打得更惨。正确做法是指数退避:每次重试的间隔翻倍,并加一点随机抖动避免「重试风暴」。
tenacity 9.2.1 把这件事做成声明式装饰器:
import logging
from tenacity import (
retry, stop_after_attempt, wait_exponential,
retry_if_exception_type, before_sleep_log,
)
logging.basicConfig(level=logging.INFO, format="%(message)s")
log = logging.getLogger("demo")
class Transient(Exception):
pass
attempts = 0
@retry(
stop=stop_after_attempt(4),
wait=wait_exponential(multiplier=0.05, min=0.05, max=0.4),
retry=retry_if_exception_type(Transient),
before_sleep=before_sleep_log(log, logging.WARNING),
reraise=True,
)
def flaky():
global attempts
attempts += 1
if attempts < 3:
raise Transient(f"第 {attempts} 次失败")
return f"第 {attempts} 次成功"
print("结果:", flaky())
print("总尝试次数:", attempts)
Retrying __main__.flaky in 0.05 seconds as it raised Transient: 第 1 次失败.
Retrying __main__.flaky in 0.1 seconds as it raised Transient: 第 2 次失败.
结果: 第 3 次成功
总尝试次数: 3
三次尝试、两次退避,间隔 0.05s → 0.1s 逐次翻倍(日志里 tenacity 会把每次等待的秒数打出来),累计等待约 0.15 秒。几个关键参数值得记住:
| 参数 | 作用 |
|---|---|
stop_after_attempt(n) | 最多尝试 n 次(含首次) |
wait_exponential(multiplier, min, max) | 指数退避,间隔 = multiplier × 2ⁿ,受 min/max 夹逼 |
retry_if_exception_type(...) | 只对可重试的异常重试,参数错误不该重试 |
before_sleep | 每次重试前打日志,方便观测 |
reraise=True | 次数耗尽后抛原始异常,而非 RetryError |
retry_if_exception_type 这一条最容易被忽略:只有瞬时故障才该重试。ValueError(参数错)、PermissionError(权限错)重试一万次也不会成功,只会浪费时间并堆积死信。
8.2.5 最小可靠队列:asyncio + fakeredis
现在把上面的机制亲手实现一遍。我们用纯 asyncio + fakeredis 2.39.0 写一个最小队列,它包含三个 Redis 结构:
- 主队列(List):待执行的任务。
- 处理中列表(List):已取出但未确认的任务,用于可见性超时。
- 死信队列(List):重试耗尽的任务,等待人工处理。
取任务用 RPOPLPUSH——它原子地从主队列取出并放进处理中列表,这一步是「不丢任务」的核心:
import asyncio
import json
import uuid
from dataclasses import dataclass, asdict
import fakeredis
QUEUE = "app:q:default"
PROCESSING = "app:q:default:processing"
DLQ = "app:q:default:dead"
MAX_ATTEMPTS = 3
@dataclass
class Task:
id: str
name: str
payload: dict
attempts: int = 0
def to_json(self) -> str:
return json.dumps(asdict(self), ensure_ascii=False)
@classmethod
def from_json(cls, s: str) -> "Task":
return cls(**json.loads(s))
class Broker:
def __init__(self, redis):
self.redis = redis
def submit(self, name, payload):
task = Task(id=uuid.uuid4().hex[:8], name=name, payload=payload)
self.redis.lpush(QUEUE, task.to_json())
return task.id
def fetch(self):
"""RPOPLPUSH:从待处理队列取出并放进 processing,保证不丢"""
raw = self.redis.rpoplpush(QUEUE, PROCESSING)
return Task.from_json(raw) if raw else None
def ack(self, task):
self.redis.lrem(PROCESSING, 1, task.to_json())
def retry(self, task):
self.redis.lrem(PROCESSING, 1, task.to_json())
task.attempts += 1
self.redis.lpush(QUEUE, task.to_json())
def dead_letter(self, task, reason):
self.redis.lrem(PROCESSING, 1, task.to_json())
self.redis.lpush(DLQ, json.dumps(
{"task": json.loads(task.to_json()), "reason": reason}, ensure_ascii=False))
worker 循环取任务、调用处理器,成功就 ack,失败就重试或进死信:
HANDLERS = {}
def handler(name):
def deco(fn):
HANDLERS[name] = fn
return fn
return deco
@handler("send_email")
def send_email(payload):
if payload.get("to") == "bad@example.com":
raise ValueError("SMTP 拒收")
return f"sent to {payload['to']}"
async def worker(broker, stop: asyncio.Event):
while not stop.is_set():
task = broker.fetch()
if task is None:
await asyncio.sleep(0.01)
continue
try:
result = HANDLERS[task.name](task.payload)
print(f" [{task.name}] 成功 -> {result}")
broker.ack(task)
except Exception as e:
if task.attempts + 1 >= MAX_ATTEMPTS:
print(f" [{task.name}] 连续失败 {MAX_ATTEMPTS} 次 -> 进死信")
broker.dead_letter(task, f"{type(e).__name__}: {e}")
else:
print(f" [{task.name}] 第 {task.attempts + 1} 次失败,重试")
broker.retry(task)
async def main():
redis = fakeredis.FakeRedis(decode_responses=True)
broker = Broker(redis)
broker.submit("send_email", {"to": "ok@example.com"})
broker.submit("send_email", {"to": "bad@example.com"})
stop = asyncio.Event()
task = asyncio.create_task(worker(broker, stop))
await asyncio.sleep(0.3) # 留出 worker 处理的时间
stop.set()
await task
print("待处理:", redis.llen(QUEUE), " 处理中:", redis.llen(PROCESSING),
" 死信:", redis.llen(DLQ))
dead = json.loads(redis.lindex(DLQ, 0))
print("死信内容:", dead["task"]["name"], "attempts=", dead["task"]["attempts"],
"reason=", dead["reason"])
asyncio.run(main())
跑起来的输出:
[send_email] 成功 -> sent to ok@example.com
[send_email] 第 1 次失败,重试
[send_email] 第 2 次失败,重试
[send_email] 连续失败 3 次 -> 进死信
待处理: 0 处理中: 0 死信: 1
死信内容: send_email attempts= 2 reason= ValueError: SMTP 拒收
两个任务:一个成功、一个连续失败 3 次后被投进死信。注意队列的终态——主队列和处理中列表都清空了,死信里恰好一条,这就是一个可靠队列应有的「不丢不重不漏」。
8.2.6 结果后端与任务状态
有些任务调用方需要知道「跑完了没、结果是什么」。这时要引入结果后端:任务执行完把返回值写进一个 Redis 键(如 result:<task_id>),调用方用 task_id 轮询。Celery 的 backend、arq 的 job.result() 都是这个思路。
但结果后端要谨慎使用:它把 Redis 变成了「临时数据库」,键必须设 TTL,否则任务结果会无限堆积。更重要的是——大多数任务根本不需要结果后端。发邮件不需要返回值,报表结果直接写对象存储。只有当调用方真的要用这个结果做后续决策时,才值得引入轮询或回调。
延伸阅读:Celery 的任务编排(chain / group / chord)、结果后端与 beat 定时可参考专题 Python Celery 任务队列 。
小结
- 异步化的判断标准是「能否容忍延迟 + 失败能否重试 + 结果是否同步需要」,三条都满足才该进队列;强一致操作绝不能异步。
- Celery 功能全、适合复杂编排;arq 纯协程、轻量,适合 asyncio 技术栈。选型看团队与场景,没有绝对优劣。
- 队列的可靠性来自可见性超时,代价是「至少一次」投递——任务可能被跑两次,所以业务必须幂等。
- 重试用指数退避加抖动,且只重试瞬时故障;
tenacity用装饰器声明stop/wait/retry_if即可。 - 可靠队列的最小骨架是三个结构:主队列、处理中列表、死信队列,靠
RPOPLPUSH保证取出不丢。
我们已经有了「失败会重试、耗尽进死信」的队列。但重试带来的「至少一次」语义还留着一个问题:任务被跑两次,业务怎么不重复执行?下一节我们就来解决幂等,并补上定时任务与死信重放的完整闭环。
阅读导航:上一节:Redis 缓存层次与键设计 · 下一节:定时任务、幂等与死信处理 。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。