《Python编程实战》8.3 定时任务、幂等与死信处理

本节讲定时任务的三种形态与延时队列的 ZSET 调度,重点落在幂等键的设计:用 SET NX 抢占、状态机去重与结果缓存保证重复投递只执行一次,并区分「至少一次」投递语义下如何做到业务幂等;最后给出死信队列的人工重放流程与告警口径,全部用 fakeredis 2.39.0 实测。

本节目标:掌握定时与延时任务的调度方式,理解「至少一次」投递为什么要求幂等,能设计出可靠的幂等键与状态机,并搭建死信队列与人工重放的闭环。
适用版本:Python 3.12+(实测 3.14.6);fakeredis 2.39.0

8.3 定时任务、幂等与死信处理

8.2 节我们做出了一个会重试、会进死信的队列。但它还留着两个没解决的问题:一是时间——怎么让任务「每天早上 9 点跑」或「3 分钟后跑」;二是重复——可见性超时带来的「至少一次」投递意味着同一个任务可能被执行两次,业务怎么保证不会重复扣款、重复发货。

这一节把这两块补齐,并给死信一个正式的出口:人工重放。这三件事看起来独立,其实指向同一个工程主题——让异步系统在不可靠的网络与进程里,表现得像可靠的一样。

8.3.1 定时任务的三种形态

「定时」其实有三种不同的时间语义,混在一起讲会乱:

形态语义例子实现
周期性(cron)固定时刻反复触发每天 9:00 发日报celery beat / arq cron / 系统 crontab
周期性(interval)每隔固定时长触发一次每 30 秒探活循环 sleep 或调度器
一次性延时(delayed)N 秒后执行一次下单 15 分钟后未支付则关单延时队列(ZSET)

前两种靠调度器反复触发,第三种是一次性的、带具体执行时间点的任务。延时任务在业务里极其常见(超时关单、延迟重试、定时提醒),但大多数队列框架的「延时」能力都很弱,所以实战中往往自己用 Redis 的有序集合实现。

固定间隔的任务用 asyncio 写起来很轻,sleep 加循环就是一个最朴素的 interval 调度器:

import asyncio

async def run_every(interval: float, job, ticks: int):
    """每隔 interval 秒执行一次 job,共执行 ticks 次"""
    for i in range(ticks):
        await asyncio.sleep(interval)          # 固定间隔
        await job(i)

async def heartbeat(i):
    print(f"  [t={i}] 心跳任务执行")

async def main():
    await run_every(0.05, heartbeat, 3)

asyncio.run(main())
  [t=0] 心跳任务执行
  [t=1] 心跳任务执行
  [t=2] 心跳任务执行

三个 tick 各间隔 0.05 秒,总耗时约 0.15 秒。但这个朴素写法有个累积漂移的毛病:job 本身若耗时 0.02 秒,下一个 tick 实际是 0.05 + 0.02 之后才开始,越跑越偏。要求严格对齐时钟的场景(「每天 9:00」)应该算出下一个绝对触发时刻再 sleep(下一个时刻 - now),而不是简单地 sleep(interval)。

8.3.2 延时队列:用 ZSET 按时间排序

有序集合(Sorted Set)的 score 是浮点数,把它当成执行时间戳,就得到一个按「何时该执行」排序的延时队列。入队时 ZADD,取出时用 ZRANGEBYSCORE 捞所有 score 小于等于当前时刻的任务:

import json
import time
import fakeredis

redis = fakeredis.FakeRedis(decode_responses=True)
SCHED = "app:q:delayed"      # score = 计划执行的时间戳(秒)

def enqueue_at(task_id: str, delay_s: float) -> float:
    score = time.time() + delay_s
    redis.zadd(SCHED, {json.dumps({"id": task_id}): score})
    return score

def poll_due(now: float | None = None) -> list[dict]:
    """取出所有到点的任务(score <= now)"""
    now = now if now is not None else time.time()
    due = redis.zrangebyscore(SCHED, "-inf", now)
    if due:
        redis.zremrangebyscore(SCHED, "-inf", now)
    return [json.loads(x) for x in due]

enqueue_at("t1", 0.05)
enqueue_at("t2", 0.15)
print("入队后待调度:", redis.zcard(SCHED))
time.sleep(0.08)
print("t=0.08 到点:", poll_due())
time.sleep(0.1)
print("t=0.18 到点:", poll_due())
print("剩余待调度:", redis.zcard(SCHED))
入队后待调度: 2
t=0.08 到点: [{'id': 't1'}]
t=0.18 到点: [{'id': 't2'}]
剩余待调度: 0

两个任务分别在 0.05s、0.15s 后到期,poll_due 在正确的时间点各捞出一个。这里有一个必须注意的原子性问题:ZRANGEBYSCORE 取出后、ZREMRANGEBYSCORE 删除前,如果 worker 崩溃,任务会被下次轮询重复取出。生产环境要用 Lua 脚本把「取出 + 删除」做成原子操作,或改成「取出 + 标记为处理中」。这正是「至少一次」语义的又一次现身——重复执行是常态,不是意外。

8.3.3 幂等键:让重复执行无害

既然任务会被跑两次,「保证只执行一次」这个目标本身就不现实。工程上的正确目标不是「不重复」,而是**「重复执行和一次执行的效果相同」**——这就是幂等。

幂等的第一把工具是幂等键(idempotency key):给每个业务操作分配一个全局唯一标识(订单号、请求 ID),执行前先抢占这个键,抢到才真正执行。Redis 的 SET NX 天生适合做这件事:

def run_once(idem_key: str, fn):
    if not redis.set(f"app:idem:{idem_key}", "1", nx=True, ex=3600):
        return "已处理过,跳过"
    return fn()

calls = 0

def charge():
    global calls
    calls += 1
    return f"扣款成功 #{calls}"

print("第一次:", run_once("order-1001-charge", charge))
print("重复提交:", run_once("order-1001-charge", charge))
print("实际扣款次数:", calls)
第一次: 扣款成功 #1
重复提交: 已处理过,跳过
实际扣款次数: 1

三次调用只真正扣款一次。但只用 SET NX 有一个缺点:重复请求拿到的是「跳过」,而不是原本的结果。对于「用户重复点了支付按钮」这种场景,调用方期望拿到和第一次一样的响应,而不是一句「已处理」。这时要把幂等键升级成状态机 + 结果缓存:

def process_payment(order_id: str, amount: int):
    key = f"app:idem:pay:{order_id}"
    if redis.set(key, json.dumps({"status": "processing"}), nx=True, ex=3600):
        result = {"order": order_id, "charged": amount}   # 真正扣款
        redis.set(key, json.dumps({"status": "done", "result": result}), ex=3600)
        return result, "首次执行"
    state = json.loads(redis.get(key))
    if state["status"] == "done":
        return state["result"], "重复请求,返回缓存结果"
    return None, "处理中,请稍后"

print(process_payment("O-1", 99))
print(process_payment("O-1", 99))
print(process_payment("O-1", 99))
({'order': 'O-1', 'charged': 99}, '首次执行')
({'order': 'O-1', 'charged': 99}, '重复请求,返回缓存结果')
({'order': 'O-1', 'charged': 99}, '重复请求,返回缓存结果')

这里藏着三个关键设计。第一,抢占和写入状态是两步,中间有极小的时间窗,所以 processing 状态是对外可见的——第二个请求看到它就知道「有人正在处理」,而不是误以为「没人在做」。第二,结果被缓存,重复请求返回完全相同的数据,对调用方来说这就是「只执行了一次」。第三,processing 状态需要超时兜底:如果第一次执行的进程在写 done 之前崩了,键会永远停在 processing。解决办法是给 processing 状态设一个较短的 TTL,或记录一个「处理开始时间」让后续请求能接管超时的任务。

8.3.4 幂等的边界:哪些操作天然幂等

不是所有操作都需要幂等键。区分「天然幂等」和「需要额外处理」能省下大量代码:

操作是否天然幂等说明
SET k v是覆盖写,跑几次结果一样
DELETE k是删不存在的键不报错
INCR k否每跑一次加一次,必须去重
INSERT(无唯一约束)否会插入多行,需唯一索引兜底
INSERT(有唯一约束)是冲突即拒绝,天然去重
发送通知 / 邮件否会重复发,需幂等键或去重表

规律很清楚:「赋值型」操作天然幂等,「累加型 / 追加型」操作需要额外保护。最优雅的做法往往不是在应用层写幂等键,而是把幂等性下沉到数据库——给业务表加唯一约束(如 (order_id, action)),重复插入自然被拒绝,比应用层的 SET NX 更可靠,因为它和业务数据在同一个事务里。

8.3.5 死信队列与人工重放

重试耗尽的、参数永远不对的、依赖服务长期不可用的任务,继续留在主队列只会阻塞后面的任务。它们应该被挪进死信队列(DLQ)——一个「失败任务收容所」,等人工介入。

死信不是终点,它需要一条重放(replay)路径:运维看过失败原因、修好问题后,把死信重新投回主队列。重放的关键是重置重试计数,否则任务一进队列就又立刻进死信:

DLQ = "app:q:default:dead"
QUEUE = "app:q:default"

def list_dead(limit: int = 10) -> list[dict]:
    return [json.loads(x) for x in redis.lrange(DLQ, 0, limit - 1)]

def replay(task_id: str) -> bool:
    """把指定死信重新投回主队列,并从死信移除"""
    for raw in redis.lrange(DLQ, 0, -1):
        d = json.loads(raw)
        if d["task"]["id"] == task_id:
            task = d["task"]
            task["attempts"] = 0          # 重置重试计数
            redis.lrem(DLQ, 1, raw)
            redis.lpush(QUEUE, json.dumps(task, ensure_ascii=False))
            return True
    return False

for tid, reason in [("a1", "SMTP 拒收"), ("b2", "下游 503")]:
    redis.lpush(DLQ, json.dumps({"task": {"id": tid, "name": "send_email", "attempts": 3},
                                 "reason": reason}, ensure_ascii=False))

print("死信列表:")
for d in list_dead():
    print("  ", d["task"]["id"], "->", d["reason"])
print("重放 a1:", replay("a1"))
print("重放后 主队列:", redis.llen(QUEUE), " 死信剩余:", redis.llen(DLQ))
死信列表:
   b2 -> 下游 503
   a1 -> SMTP 拒收
重放 a1: True
重放后 主队列: 1  死信剩余: 1

replay("a1") 把任务从死信挪回主队列,attempts 归零,死信从 2 条减到 1 条。生产环境的死信处理通常还有三件事:给死信设上限(超过阈值就丢弃最老的,防止 DLQ 无限膨胀)、按原因分类(SMTP 拒收这种是永久失败,重放也没用;下游 503才是可恢复的)、重放要审计(谁、何时、重放了多少条,否则重复扣款这类事故无从追溯)。

8.3.6 可观测性:让失败被看见

死信队列最大的风险不是「任务进了死信」,而是没人知道任务进了死信。异步任务的失败发生在后台,没有用户投诉来提醒你,所以必须靠指标和告警主动发现。

至少要盯四个指标:

指标含义告警条件
队列深度主队列积压的任务数持续增长或超阈值
处理延迟任务入队到执行的耗时P99 超 SLA
重试率重试次数 / 总执行次数突增说明下游异常
死信增长单位时间进入 DLQ 的数量任何非零都值得看

队列深度持续增长通常意味着消费能力不足(worker 太少或任务太慢);死信突增则意味着下游故障或代码 bug。两类问题的处理方向完全不同,所以指标要分开看。这些指标接到 3.3 节讲过的监控体系里,就能在用户察觉之前发现问题。

延伸阅读:Celery 的 beat 定时、结果后端与死信策略可参考专题 Python Celery 任务队列 。

小结

  • 定时任务分 cron(固定时刻)、interval(固定间隔)、delayed(一次性延时)三类;延时任务用 Redis 的 ZSET 以时间戳为 score 实现。
  • ZSET 延时队列的「取出 + 删除」不是原子操作,worker 崩溃会导致重复执行——「至少一次」是异步系统的常态。
  • 幂等的目标不是「不重复执行」,而是「重复执行与一次执行效果相同」;用 SET NX 抢占幂等键是最小实现。
  • 需要返回一致结果时,把幂等键升级为「状态机 + 结果缓存」,并给 processing 状态设超时兜底。
  • 赋值型操作天然幂等,累加/追加型必须保护;最可靠的做法是把幂等下沉到数据库的唯一约束。
  • 死信队列要有重放路径(重放时重置重试计数)、上限保护、按原因分类和审计;更要盯住队列深度、重试率、死信增长等指标主动告警。

到这里,8.1 到 8.3 把「缓存 + 队列 + 定时 + 幂等 + 死信」这条异步链路串了起来。下一章我们转向另一个绕不开的工程主题——认证与授权:用户是谁(认证)、他能做什么(授权),以及多租户下如何隔离数据。

阅读导航:上一节:任务队列与重试 · 下一节:认证与会话:JWT / OAuth2 。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「python」更多文章

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