《Python编程实战》8.2 任务队列与重试

本节先对比 Celery 与 arq 的模型差异,讲清 broker、worker、结果后端与可见性超时的机制;再用 tenacity 9.2.1 实测指数退避重试,并用纯 asyncio 加 fakeredis 2.39.0 手写一个最小可靠队列,演示 RPOPLPUSH 取任务、失败重试与死信投递的完整链路。

本节目标:判断哪些工作该异步化,理解任务队列的 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 是异步的、极简的「轻骑兵」。

维度Celeryarq
编程模型同步函数(用进程/线程池执行)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 缓存层次与键设计 · 下一节:定时任务、幂等与死信处理 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「python」更多文章

  1. 《Python高级编程》目录
  2. 《Python高级编程》11.3 PEP 流程与版本迁移策略
  3. 《Python高级编程》11.2 嵌入式与自由线程运行时